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 @@ -2364,10 +2364,14 @@ impl LatestCodeTextGenerationV1 {
Err(VerifiedSealedLexicalCursorRestoreErrorV1::IncompatiblePosition) => {
drop(builder);
store.discard_incompatible_staging(&staging_path, control)?;
builder = CodeLexicalArtifactBuilderV1::create_with_memory_budget(
// Keep the same V14 admission writer as the cold create path.
// Defaulting to V16 here would admit clone fingerprints on the
// rebuild lane and diverge from the successor-backed cutover.
builder = CodeLexicalArtifactBuilderV1::create_with_memory_budget_and_format_revision(
&staging_path,
metadata,
builder_budget,
CodeLexicalArtifactWriterRevisionV1::V14,
)
.map_err(map_text_artifact_error)?;
progress = builder.progress().map_err(map_text_artifact_error)?;
Expand Down Expand Up @@ -2801,16 +2805,18 @@ impl LatestCodeTextGenerationV1 {
.map_err(map_text_artifact_error)?;
let needs_clone_successor = !reader.has_clone_fingerprints();
let prior = reader.verified_artifact().clone();
// Publish Ready before installing owners so status cannot observe
// query-ready owners while the last finalization wake still says
// Verification.
self.publish_text_progress_phase(CodeIndexBuildPhaseV1::Ready, 0, 0);
self.install_artifact_owners(reader, reader_reservation)?;
Comment on lines +2811 to 2812

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 Publish Ready only after owner installation succeeds

If install_artifact_owners returns an error—for example because the reader exceeds its admitted ceiling or the reservation cannot shrink—the progress snapshot has already been permanently advanced to Ready, while no query owners were installed. The caller latches the projection failure but does not roll this phase back, so dashboard/MCP progress reports a ready generation that cannot serve queries; coordinate these transitions atomically or publish Ready only after successful installation.

AGENTS.md reference: AGENTS.md:L9-L12

Useful? React with 👍 / 👎.

if needs_clone_successor {
let source = store.open_sealed_source(&sealed_identity, control)?;
let build =
self.begin_clone_successor(descriptor, prior, sealed_identity, source, control)?;
drop(publish_claim.install(TextHeadOpenBuildV1::CloneSuccessor(build)));
self.publish_text_progress_phase(CodeIndexBuildPhaseV1::Ready, 0, 0);
return Ok(false);
}
self.publish_text_progress_phase(CodeIndexBuildPhaseV1::Ready, 0, 0);
Ok(true)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8547,28 +8547,41 @@ async fn graph_off_changed_source_advances_text_authority_without_full_decode()
let scheduler = scheduler
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let active = scheduler
// Seal releases the decoded active generation so text projection does
// not keep a whole-generation owner. Inspect the durable pointer — do
// not call load_active_shared here or the probe itself would decode.
assert_eq!(
scheduler.sealed_decode_count(),
0,
"graph-off A-to-C publication must not decode a sealed generation"
);
let pointer = scheduler
.publication
.load_active_shared()
.expect("load generation C from the in-memory publication authority")
.expect("generation C is active");
.read_publication_pointer()
.expect("read generation C pointer")
.expect("generation C is durable");
let current = scheduler
.capture_authoritative_snapshot_without_active_generation_reuse(None)
.expect("capture current generation C revision");
assert_ne!(
active.manifest().generation_id,
unpublished_b_generation,
pointer.generation_id,
unpublished_b_generation.as_str(),
"generation B must never become active after generation C is observed"
);
assert_eq!(
active.snapshot().source_revision,
current.snapshot.source_revision,
pointer.generation_id,
generation_c.as_str(),
"the durable pointer must name generation C"
);
assert_eq!(
pointer.snapshot_content_identity,
current.snapshot.content_identity.as_str(),
"the successor must record the allow-empty generation C revision"
);
assert_eq!(
scheduler.sealed_decode_count(),
0,
"graph-off A-to-C publication must not decode a sealed generation"
"pointer and snapshot inspection must not decode a sealed generation"
);
}
registry.shutdown().await;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1697,13 +1697,18 @@ fn page_aligned_final_source_page_converges_the_text_projection() {
}
assert!(latest.query_owners_are_ready());
let progress = build_progress_snapshot(&scheduler);
// Every committed page must be chunk-full: an early commit from the page
// byte bound or an import record would leave the final page partial, which
// is exactly the shape that does not trip this invariant.
// Clone-body lanes mint their own pages after chunk pages. Chunk packing
// must still land on the page-record bound; clone-body pages may add
// additional page ordinals beyond the chunk-only count.
let page_chunks = super::super::TEXT_ARTIFACT_PAGE_CHUNKS_V1 as u64;
assert_eq!(
progress.committed_chunks,
progress.committed_pages * super::super::TEXT_ARTIFACT_PAGE_CHUNKS_V1 as u64,
"the fixture must keep every page chunk-full so the final page ends on the last record"
progress.committed_chunks % page_chunks,
0,
"chunk pages must remain record-full so the final chunk page ends on a chunk boundary"
);
assert!(
progress.committed_pages * page_chunks >= progress.committed_chunks,
"clone-body pages may follow chunk pages but must not shrink chunk packing"
);
assert_eq!(
progress.committed_imports, 0,
Expand Down Expand Up @@ -1843,8 +1848,12 @@ fn invalid_partial_text_artifact_cursor_is_discarded_and_rebuilt() {
cursor[3].as_u64().is_some_and(|ordinal| ordinal > 0),
"first page must advance within the file's chunks"
);
// Rewind the chunk ordinal while leaving a non-zero import ordinal so
// restore_cursor_classified refuses the authenticated-but-impossible
// mid-file position (chunk < count && import != 0).
// Persisted layout: [3]=next_chunk_ordinal, [8]=next_import_ordinal.
cursor[3] = serde_json::Value::from(0_u64);
cursor[6] = serde_json::Value::from(1_u64);
cursor[8] = serde_json::Value::from(1_u64);

let text = |index: usize| {
cursor[index]
Expand All @@ -1861,6 +1870,10 @@ fn invalid_partial_text_artifact_cursor_is_discarded_and_rebuilt() {
.to_le_bytes(),
);
hasher.update(text(0));
// Integrity order matches VerifiedSealedLexicalCursorV1::integrity_digest:
// file_ord, file_off, chunk_ord, import_ord, clone_ord, page_ord,
// emitted_chunks, emitted_payload, emitted_imports, emitted_import_bytes,
// emitted_clones, emitted_clone_bytes.
for index in [1, 2, 3, 8, 4, 5, 6, 7, 9, 10, 11, 12] {
hasher.update(number(index).to_le_bytes());
}
Expand Down Expand Up @@ -5983,13 +5996,23 @@ async fn graph_off_overflow_preserves_text_owner_progress_without_full_decode()
.filter_map(Result::ok)
.map(|entry| entry.file_name().to_string_lossy().into_owned())
.collect::<Vec<_>>();
assert_eq!(
artifact_names
.iter()
.filter(|name| name.starts_with("text-artifact-") && name.ends_with(".bin"))
.count(),
1,
"one durable artifact owns ready text serving"
let pointer: serde_json::Value = serde_json::from_slice(
&std::fs::read(scoped_store.join("active-code-generation-v1.json"))
.expect("read active publication pointer"),
)
.expect("decode active publication pointer");
let active_artifact = pointer["generation_index"]
.as_array()
.expect("generation index")
.iter()
.find(|entry| entry["generation_id"] == pointer["generation_id"])
.and_then(|entry| entry.get("text_artifact"))
.filter(|value| !value.is_null())
.and_then(|artifact| artifact["artifact_file"].as_str())
.expect("active generation owns one durable text artifact");
assert!(
artifact_names.iter().any(|name| name == active_artifact),
"the attached text artifact must exist on disk"
);
assert_eq!(
artifact_names
Expand Down
37 changes: 30 additions & 7 deletions crates/tracedecay-code-index/src/chunks/artifacts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -379,17 +379,40 @@ impl CodeFileIndexArtifactsV1 {
.collect::<std::collections::BTreeSet<_>>();
if self.clone_bodies.windows(2).any(|pair| {
pair[0].occurrence.symbol_occurrence_id >= pair[1].occurrence.symbol_occurrence_id
}) || self.clone_bodies.iter().any(|body| {
body.occurrence.path.is_empty()
|| body.occurrence.body_span.is_empty()
|| body.occurrence.payload_digest != body.payload.payload_digest
|| (validate_payloads && body.payload.validate().is_err())
|| !occurrences.contains(&body.occurrence.symbol_occurrence_id)
}) {
return Err(ChunkingFailureV1::NonCanonicalIdentity(
"clone body evidence is not canonically bound to file symbols".to_owned(),
"clone body evidence is not in strict symbol-occurrence order".to_owned(),
));
}
for body in &self.clone_bodies {
if body.occurrence.path.is_empty() {
return Err(ChunkingFailureV1::NonCanonicalIdentity(
"clone body evidence has an empty path".to_owned(),
));
}
if body.occurrence.body_span.is_empty() {
return Err(ChunkingFailureV1::NonCanonicalIdentity(
"clone body evidence has an empty body span".to_owned(),
));
}
if body.occurrence.payload_digest != body.payload.payload_digest {
return Err(ChunkingFailureV1::NonCanonicalIdentity(
"clone body evidence payload digest does not match its payload".to_owned(),
));
}
if validate_payloads {
if let Err(detail) = body.payload.validate() {
return Err(ChunkingFailureV1::NonCanonicalIdentity(format!(
"clone body evidence payload is not canonical: {detail}"
)));
}
}
if !occurrences.contains(&body.occurrence.symbol_occurrence_id) {
return Err(ChunkingFailureV1::NonCanonicalIdentity(
"clone body evidence is not bound to a file symbol".to_owned(),
));
}
}
Ok(())
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1412,6 +1412,11 @@ fn decode_verified_file_segment(
file.artifacts.edges.sort_by(|left, right| {
crate::chunks::canonical_edge_key(left).cmp(&crate::chunks::canonical_edge_key(right))
});
file.artifacts.clone_bodies.sort_by(|left, right| {
left.occurrence
.symbol_occurrence_id
.cmp(&right.occurrence.symbol_occurrence_id)
});
file.artifacts.unresolved_references.sort();
});
Ok(file)
Expand Down
46 changes: 37 additions & 9 deletions crates/tracedecay/src/daemon/tests/rmcp_route.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,13 +29,37 @@ struct RmcpRouteFixture {
}

async fn rmcp_route_fixture(label: &str) -> RmcpRouteFixture {
rmcp_route_fixture_with_projects(label, &[]).await
}

/// Every project a test will mount is initialized here, before the engine
/// exists. Initialization takes the profile's exclusive maintenance lease and
/// writes the profile database; production never overlaps that with a live
/// daemon, and an in-process engine writing the same file turns the fixture's
/// `BEGIN IMMEDIATE` into a typed `database is locked` refusal under load.
async fn rmcp_route_fixture_with_projects(
label: &str,
extra_projects: &[(&str, &str, &str)],
) -> RmcpRouteFixture {
let temp = TempDir::new().expect("route fixture");
let project = temp.path().join("project");
let profile_root = temp.path().join("profile");
std::fs::create_dir_all(project.join("src")).expect("fixture source directory");
std::fs::write(project.join("src/main.rs"), "fn main() {}\n").expect("fixture source");
let client_identity = test_client_identity_for(profile_root.clone());
initialize_test_project(&project, &client_identity).await;
for (directory, source_path, source) in extra_projects {
let extra = temp.path().join(directory);
let source_file = extra.join(source_path);
std::fs::create_dir_all(
source_file
.parent()
.expect("extra project source directory"),
)
.expect("extra project source directory");
std::fs::write(&source_file, source).expect("extra project source");
initialize_test_project(&extra, &client_identity).await;
}
let _database_scope = enter_test_daemon_database_scope(&profile_root, label);
let handshake = DaemonHandshake {
project_path: Some(project),
Expand Down Expand Up @@ -622,18 +646,20 @@ async fn wait_for_count(counter: &AtomicUsize, expected: usize, message: &str) {
);
}

#[cfg(unix)]
const RMCP_TARGET_PROJECT: (&str, &str, &str) = (
"target-project",
"src/target.rs",
"pub const RMCP_SELECTED_TARGET_MARKER: &str = \"target-beta\";\n",
);

/// Mounts the target project the fixture already initialized through
/// [`RMCP_TARGET_PROJECT`].
#[cfg(unix)]
async fn mount_rmcp_target(
fixture: &RmcpRouteFixture,
) -> (DaemonHandshake, String, ProjectServerKey, ProjectRouteKey) {
let project = fixture._temp.path().join("target-project");
std::fs::create_dir_all(project.join("src")).expect("target source directory");
std::fs::write(
project.join("src/target.rs"),
"pub const RMCP_SELECTED_TARGET_MARKER: &str = \"target-beta\";\n",
)
.expect("target source marker");
initialize_test_project(&project, &fixture.handshake.client_identity).await;
let project = fixture._temp.path().join(RMCP_TARGET_PROJECT.0);
let handshake = DaemonHandshake {
project_path: Some(project.clone()),
client_identity: fixture.handshake.client_identity.clone(),
Expand Down Expand Up @@ -670,7 +696,9 @@ fn response_text(response: &Value) -> &str {
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn selected_target_rmcp_flushes_response_and_disconnect_cancels_selector_owner() {
let fixture = rmcp_route_fixture("rmcp-selected-target-disconnect").await;
let fixture =
rmcp_route_fixture_with_projects("rmcp-selected-target-disconnect", &[RMCP_TARGET_PROJECT])
.await;
let (_target_handshake, target_project_id, target_key, _target_route) =
mount_rmcp_target(&fixture).await;
let target_server = {
Expand Down
Loading