Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -283,6 +283,29 @@ fn assert_surface_resolves_project(
payload
}

/// Holds until the cold daemon serves a published generation's code graph.
/// Graph-backed surfaces answer the retryable `application.code-graph.unavailable`
/// before that, which is the truthful state of a daemon still activating.
fn await_graph_ready(home: &Path, project: &Path) {
let outcome = run_surface_tool_from(
home,
project,
"status",
r#"{"format":"json","wait_for":{"state":"graph_ready","timeout_ms":50000}}"#,
);
assert!(
outcome.success,
"status wait failed\nstdout:\n{}\nstderr:\n{}",
outcome.stdout, outcome.stderr
);
assert_eq!(
outcome.payload()["wait"],
serde_json::json!({ "outcome": "reached" }),
"the daemon never served a published code graph: {}",
outcome.stdout
);
}

fn surface_fixture() -> (TempDir, TempDir, PathBuf, PathBuf) {
let home = TempDir::new().unwrap();
let project = TempDir::new().unwrap();
Expand All @@ -303,6 +326,7 @@ fn application_surface_primitive_tools_resolve_the_working_directory_project() {
"storage_status",
r#"{"format":"json"}"#,
);
await_graph_ready(&home_path, &project_path);
assert_surface_resolves_project(
&home_path,
&project_path,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -353,7 +353,7 @@ impl CodeIndexSchedulerRegistryV1 {
// longer held, so no later retry can serve it either.
.ok_or(CallableCodeCursorError::Stale)
} else if is_unpinned_latest(requested) {
self.latest_complete_fresh_for_scope(request.scope())
self.latest_complete_fresh_for_scope_awaiting_seat(request.scope())
.await
.ok_or(CallableCodeCursorError::Unavailable)
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,10 +18,11 @@ use super::super::{
use super::graph_cursor_retention::GraphCursorRetentionV1;
use super::scope_identity::{latest_matches_scope_identity, text_matches_scope_identity};
use super::{
CodeIndexMountedScopeV1, CodeIndexSchedulerRegistryV1, CodeIndexServingScopeV1,
MountedCodeIndexWorktreeV1, PendingWakeClaimV1, ReadyProbeServingPartsV1,
dashboard_code_graph_serving, dashboard_freshness_identity, dashboard_terminal_status,
dashboard_text_freshness_identity, project_graph_publication_phase, unique_mounted_for_scope,
CodeIndexMountedScopeV1, CodeIndexOwnerSignalsV1, CodeIndexSchedulerRegistryV1,
CodeIndexServingScopeV1, MountedCodeIndexWorktreeV1, PendingWakeClaimV1,
ReadyProbeServingPartsV1, dashboard_code_graph_serving, dashboard_freshness_identity,
dashboard_terminal_status, dashboard_text_freshness_identity, project_graph_publication_phase,
unique_mounted_for_scope,
};
use tracedecay_runtime_core::path_safety::canonical_existing_identity;

Expand Down Expand Up @@ -838,6 +839,61 @@ impl CodeIndexSchedulerRegistryV1 {
latest_matches_scope_identity(&latest, scope).then_some(latest)
}

/// [`Self::latest_complete_fresh_for_scope`] for a read that needs the
/// whole decoded generation. A publication seats only its text owner and
/// defers the decode until a reader needs it, so this read's own demand
/// is what starts that decode: it waits for the seat rather than answering
/// the demanding request unavailable. It stops waiting once no decode is
/// pending (seated, memory-refused, shutting down, or unmounted); the
/// caller's resolution deadline bounds the rest.
pub(crate) async fn latest_complete_fresh_for_scope_awaiting_seat(
&self,
scope: &tracedecay_contracts::ResolvedScope,
) -> Option<LatestCompleteCodeIndexV1> {
let root = {
let mounted = self.mounted.lock().await;
unique_mounted_for_scope(&mounted, scope)
.unique()?
.0
.clone()
};
let mut signals = CodeIndexOwnerSignalsV1::subscribe(self, &root).await;
loop {
if let Some(latest) = self.latest_complete_fresh_for_scope(scope).await {
return Some(latest);
}
if !self.complete_seat_pending(&root).await {
// The seat can land between the miss above and this check.
return self.latest_complete_fresh_for_scope(scope).await;
}
signals.changed().await.ok()?;
}
}

/// Whether demand has asked the worker to seat the complete generation
/// of a published text owner and nothing yet stops it from doing so.
async fn complete_seat_pending(&self, project_root: &Path) -> bool {
let mounted = self.mounted.lock().await;
let Some(worktree) = mounted.get(project_root) else {
return false;
};
worktree
.complete_generation_requested
.load(Ordering::Acquire)
&& !worktree.memory_retry.waiting()
&& !worktree.shutting_down.load(Ordering::Acquire)
Comment on lines +880 to +884

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Stop waiting after terminal convergence parks

When a retained text owner has no decoded seat and its decode ends in one of the worker's terminal convergence-park paths (for example, reproducible input failure or unrecoverable publication corruption), complete_generation_requested remains permanently true while memory_retry is false and shutdown has not begun. This predicate therefore continues to report work as pending after the worker's final notification; the loop consumes the entire generation-resolution deadline—up to 30 seconds—before returning the already-known unavailable result. Include the convergence park or actual worker/decode activity in this check so terminal background repair cannot block an ordinary code-facet read.

AGENTS.md reference: AGENTS.md:L222-L224

Useful? React with 👍 / 👎.

&& worktree
.serving_generation
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_none()
&& worktree
.text_generation
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_some()
}

/// Resolve one exact scope and admit only an already-current generation.
pub async fn latest_complete_ready_for_scope(
&self,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -231,6 +231,14 @@ impl DaemonGitAuthoritySlot {
Ok(())
}

fn installed(&self) -> Result<bool, GitIndexTransactionPortError> {
Ok(self
.source
.read()
.map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)?
.is_some())
}

fn clear(&self) -> Result<(), GitIndexTransactionPortError> {
self.source
.write()
Expand Down Expand Up @@ -646,6 +654,9 @@ impl DaemonGitIndexTransactionServiceRegistry {

/// Resolve only an owner already mounted by project-open admission.
/// Missing and ambiguous roots deliberately share the same outcome.
/// Project open publishes the service before it installs the service's
/// authority; until then the owner is still mounting, and resolving it
/// would answer every read as a policy denial.
pub async fn for_repository_root(
&self,
repository_root: &std::path::Path,
Expand All @@ -662,7 +673,7 @@ impl DaemonGitIndexTransactionServiceRegistry {
let Some(entry) = matches.next() else {
return Ok(None);
};
if matches.next().is_some() {
if matches.next().is_some() || !entry.authority.installed()? {
return Ok(None);
}
Ok(Some(DaemonGitInvocationOwner {
Expand Down
61 changes: 56 additions & 5 deletions crates/tracedecay-code-index-runtime/src/git_transactions/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,16 +6,20 @@ use std::time::Duration;
use std::collections::BTreeSet;
use std::sync::Mutex;

use tracedecay_application::ProjectSourceAccessSnapshot;
use tracedecay_contracts::catalog_composition::CatalogCompositionError;
use tracedecay_contracts::{
CancellationStage, GitIndexApplyRequestV1, GitIndexTransactionPort,
GitIndexTransactionPortError, OperationTermination,
};
use tracedecay_domain::configuration::{
AuthorityRef, ConfigurationRevisionId, ScopeSourceBinding, SourceBindingId, SourceKindV1,
};
use tracedecay_domain::{
GitIndexIdempotencyKey, GitIndexJournalPhaseV1, GitIndexPreviewId, GitIndexPreviewInputV1,
GitIndexPreviewV1, GitIndexReceiptOutcomeV1, GitIndexTransactionId,
GitIndexTransactionJournalV1, GitIndexTransactionReceiptV1, GitOperationStateV1, ProjectId,
RepositoryId, RepositoryIndexStateV1, UtcMicros,
ActorId, GitIndexIdempotencyKey, GitIndexJournalPhaseV1, GitIndexPreviewId,
GitIndexPreviewInputV1, GitIndexPreviewV1, GitIndexReceiptOutcomeV1, GitIndexTransactionId,
GitIndexTransactionJournalV1, GitIndexTransactionReceiptV1, GitOperationStateV1, LocatorDigest,
ManifestDigest, ProjectId, RepositoryId, RepositoryIndexStateV1, UtcMicros, WorktreeId,
};
use tracedecay_policy::{GitConflictRiskV1, GitEffectClassifierV1};
use tracedecay_store::{
Expand Down Expand Up @@ -960,7 +964,7 @@ async fn daemon_owner_isolates_worktrees_sharing_a_project_database() {
.ensure(
database.registered.clone(),
alternate.path().to_path_buf(),
project_id,
project_id.clone(),
fixture_time(21),
)
.await
Expand All @@ -977,6 +981,25 @@ async fn daemon_owner_isolates_worktrees_sharing_a_project_database() {
),
"all Git mutation services must serialize through the daemon registry queue"
);
assert!(
registry
.for_repository_root(directory.path())
.await
.expect("mounting owner lookup")
.is_none(),
"an owner whose authority project open has not installed yet is still mounting"
);
for root in [directory.path(), alternate.path()] {
registry
.install_authority(
root,
source_access(&project_id),
database.registered.clone(),
tokio::runtime::Handle::current(),
)
.await
.expect("project open installs each worktree's authority");
}
let primary_owner = registry
.for_repository_root(directory.path())
.await
Expand All @@ -991,6 +1014,34 @@ async fn daemon_owner_isolates_worktrees_sharing_a_project_database() {
assert!(Arc::ptr_eq(&linked_owner.service, &linked));
}

fn source_access(project_id: &ProjectId) -> ProjectSourceAccessSnapshot {
ProjectSourceAccessSnapshot {
scope: tracedecay_contracts::ResolvedScope::new(
project_id.clone(),
RepositoryId::new("repository.singleton.fixture").expect("fixture id"),
WorktreeId::new("worktree.singleton.fixture").expect("fixture id"),
None,
)
.expect("scope"),
requester: ActorId::new("actor.singleton.fixture").expect("fixture id"),
binding: ScopeSourceBinding::new(
SourceBindingId::new("binding.singleton.fixture").expect("fixture id"),
SourceKindV1::Cursor,
LocatorDigest::new(format!("sha256:{}", "a".repeat(64))).expect("locator digest"),
AuthorityRef::Project(project_id.clone()),
)
.expect("binding"),
configuration_revision: ConfigurationRevisionId::new("revision.singleton.fixture")
.expect("fixture id"),
configuration_digest: ManifestDigest::new(format!("sha256:{}", "b".repeat(64)))
.expect("configuration digest"),
configuration_provenance_digest: ManifestDigest::new(format!("sha256:{}", "c".repeat(64)))
.expect("configuration provenance"),
effective_capabilities: BTreeSet::new(),
grant_expires_at: fixture_time(100),
}
}

/// Symlink alias of an already-mounted root must resolve to the same owner
/// entry. Distinct filesystem identities remain rejected (see rebinding test).
#[cfg(unix)]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,14 @@ async fn git_owner_uses_explicit_canonical_catalog_and_rechecks_authorization()
)
.await
.unwrap();
assert!(
registry
.for_repository_root(&project_root)
.await
.unwrap()
.is_none(),
"an owner whose authority project open has not installed yet is still mounting"
);
registry
.install_authority(
&project_root,
Expand Down
7 changes: 6 additions & 1 deletion crates/tracedecay/tests/common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -216,11 +216,16 @@ impl IsolatedHome {

/// Reuses the installed toolchain in a child whose `HOME` is this
/// isolated home: rustup and cargo would otherwise resolve their state
/// under the empty throwaway home.
/// under the empty throwaway home. The child's hermetic `PATH` gains only
/// the toolchain's own `$CARGO_HOME/bin`, where rustup installs `cargo`
/// and the `rust-analyzer` proxy the daemon resolves through `PATH`.
pub fn apply_toolchain_env(&self, command: &mut Command) {
for (key, value) in &self.toolchain_environment {
match value {
Some(value) => {
if *key == "CARGO_HOME" {
command.env("PATH", hermetic_path(&[Path::new(value).join("bin")]));
}
command.env(key, value);
}
None => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -277,14 +277,22 @@ fn assert_hook_completed_analytics(provider: &str, state: &str, project: &Path,
);
}

let project_path = project.to_string_lossy();
// The payload's session is attribution, not content: the row names it so
// analytics can join a hook to its session, and no other field repeats it.
let session_id = format!("{provider}-host-fixture");
assert_eq!(row["session_id"], session_id.as_str(), "{label}");
let mut unattributed = row.clone();
unattributed
.as_object_mut()
.expect("hook row is an object")
.remove("session_id");
let project_path = project.to_string_lossy();
for private in [
project_path.as_ref(),
session_id.as_str(),
"Inspect the fixture.",
] {
assert_json_strings_omit(row, private, &label);
assert_json_strings_omit(&unattributed, private, &label);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -327,6 +327,9 @@ async fn git_runtime_fixture() -> RuntimeFixture {
)
.expect("write staged Git change");
git(&project, &["add", "src/main.rs"]);
// Git reads answer a typed retryable `mounting` problem until project open
// installs the owner's authority; the journeys start from a served project.
await_published_code_index(environment.home(), &project);

let handshake = tracedecay::daemon::handshake_for_current_client(
environment.profile(),
Expand Down
Loading