From 56adcc9199d7f0b05d7586d5d54bff1ab8bed88e Mon Sep 17 00:00:00 2001 From: "zackary.l.jackson" Date: Sun, 4 Oct 2026 16:10:03 +0000 Subject: [PATCH 1/2] perf(sessions): skip unchanged host bookkeeping writes on idle passes The daemon reruns the history pass every idle minute. Each pass rewrote every provider's host-coverage row and the Kimi and Pi discovery frontiers with a bumped mtime even when nothing changed: nine commits per idle project pass, plus the same on the user store, each adding WAL frames that readers and checkpoints reread. revise_host_record now writes such a record only when its byte_offset or file_id changes; mtime stays the revision that lets a changed record move backwards past the monotonic cursor guard. The Codex callers' own unchanged-coverage guard is gone, since the shared write owns it. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../tests/session_store_read_cost.rs | 70 +++++++++++++++++++ .../src/runtime/hosts/kimi.rs | 38 +++++----- .../src/runtime/hosts/pi.rs | 31 ++++---- .../src/runtime/ingest/project_provider.rs | 19 +++-- .../src/runtime/ingest/user_provider.rs | 19 +++-- .../tracedecay-sessions/src/runtime/source.rs | 41 +++++++++-- 6 files changed, 155 insertions(+), 63 deletions(-) diff --git a/crates/tracedecay-session-runtime/tests/session_store_read_cost.rs b/crates/tracedecay-session-runtime/tests/session_store_read_cost.rs index 6d4e257963..660a3e1cff 100644 --- a/crates/tracedecay-session-runtime/tests/session_store_read_cost.rs +++ b/crates/tracedecay-session-runtime/tests/session_store_read_cost.rs @@ -35,11 +35,13 @@ use tracedecay_domain::{ use tracedecay_global_db::observation::retention::ObservationRetentionConfig; use tracedecay_global_db::tests::harness::{HostAdmissionTestRuntimeV1, writer_telemetry}; use tracedecay_global_db::{RegisteredGlobalDb, RegisteredGlobalDbLeaseV1}; +use tracedecay_host_admission::session_ingest_authority::GlobalDbSessionIngestAuthority; use tracedecay_host_admission::{HostAdmissionAuthorities, HostAdmissionFacade}; use tracedecay_lcm::LcmRetentionConfig; use tracedecay_maintenance::retention::registered_store::run_registered_store_retention; use tracedecay_privacy::{ObservationRecordParseErrorV1, parse_normalized_observation_record_v1}; use tracedecay_runtime_core::background_cpu::ProcessBackgroundCpuV1; +use tracedecay_runtime_core::config::ProfileRoot; use tracedecay_session_runtime::session_sync::test_harness::{ SessionTemporalRefreshWakeState, run_session_temporal_refresh_pass, }; @@ -51,6 +53,9 @@ use tracedecay_sessions::observation::{ CaptureObservationOutcome, CaptureObservationRequest, ObservationCancellation, }; use tracedecay_sessions::repository_provenance::RepositoryProvenanceAdmissionContext; +use tracedecay_sessions::runtime::{ + TranscriptIngestOutcome, ingest_project_sources_for_provider, with_transcript_source_profile, +}; use tracedecay_store::WAL_SOFT_LIMIT_BYTES; const PROVIDER: &str = "codex"; @@ -975,6 +980,71 @@ async fn streamed_message_commits_once_per_durability_boundary() { ); } +/// One project history pass over every host provider, as the session +/// temporal refresh runs it, reading transcripts under `home`. +async fn project_history_pass( + fixture: &DrainFixture, + database: &RegisteredGlobalDbLeaseV1, + home: &Path, +) -> TranscriptIngestOutcome { + let authority = GlobalDbSessionIngestAuthority::new(database.clone()).with_background_cpu( + Arc::new(ProcessBackgroundCpuV1::new(NonZeroUsize::new(4).unwrap())), + ); + let shard = &database.binding().shard_id; + with_transcript_source_profile( + ProfileRoot::under_home(home), + ingest_project_sources_for_provider( + &shard.brain_id, + &shard.profile_id, + &authority, + &fixture.project, + Some(fixture.project_id.clone()), + None, + true, + ), + ) + .await +} + +/// The daemon reruns the history pass every idle minute whether or not a +/// host wrote anything. A pass that finds nothing new must leave the store +/// as it was: every commit it makes is WAL every reader and checkpoint then +/// rereads, for as long as the daemon idles. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn idle_history_pass_commits_nothing() { + let _measured = MEASURED.lock().await; + let fixture = DrainFixture::open().await; + let database = fixture + .runtime + .registered_database_lease(HostAdmissionScope::Project) + .unwrap(); + let store = session_store_path(&fixture); + let home = fixture._tmp.path().join("home"); + std::fs::create_dir_all(&home).unwrap(); + + let start = wal_mark(&store); + let first = project_history_pass(&fixture, &database, &home).await; + let settled = wal_mark(&store); + let idle = project_history_pass(&fixture, &database, &home).await; + let idled = wal_mark(&store); + + assert!( + first.failures.is_empty() && idle.failures.is_empty(), + "history passes must not fail: first={:?} idle={:?}", + first.failures, + idle.failures, + ); + assert!( + wal_commits(&store, start, settled).is_some_and(|commits| commits > 0), + "the first pass records each provider's coverage" + ); + assert_eq!( + wal_commits(&store, settled, idled), + Some(0), + "a history pass that finds nothing new must commit nothing" + ); +} + /// Levels of every B-tree in the fixture's project session store: the pages /// one point lookup in that tree reads on a cold reader. fn session_store_tree_depths(fixture: &DrainFixture) -> BTreeMap { diff --git a/crates/tracedecay-sessions/src/runtime/hosts/kimi.rs b/crates/tracedecay-sessions/src/runtime/hosts/kimi.rs index 4decad93a8..67468c5cff 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/kimi.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/kimi.rs @@ -29,7 +29,7 @@ use crate::runtime::snapshot_observation::{ use crate::runtime::source::{ FileDiscoveryLimit, HostProviderCoverage, TranscriptDiscoveryBounds, TranscriptIngestError, TranscriptIngestResult, bound_path_list, canonical_framed_sha256, jsonl_file_identity, - persist_host_provider_coverage, run_blocking_transcript_section, + persist_host_provider_coverage, revise_host_record, run_blocking_transcript_section, }; mod discovery; @@ -837,29 +837,27 @@ pub async fn capture_kimi_observations( && !cancellation.is_cancelled() { let next_frontier = if discovery.reached_end { - Some(ParseOffset { - byte_offset: 0, - mtime: discovery_frontier.mtime.saturating_add(1), - file_id: 0, - }) + Some((0, 0)) } else { - last_discovered_entry.map(|entry| ParseOffset { - byte_offset: entry.sequence, - mtime: discovery_frontier.mtime.saturating_add(1), - file_id: entry.sequence, - }) + last_discovered_entry.map(|entry| (entry.sequence, entry.sequence)) }; - if let Some(next_frontier) = next_frontier + if let Some((byte_offset, file_id)) = next_frontier && !cancellation.is_cancelled() { - facade - .advance_parse_offset(&scope, KIMI_DISCOVERY_FRONTIER_KEY, next_frontier) - .await - .map_err(|outcome| { - crate::runtime::snapshot_observation::host_admission_error( - PROVIDER, outcome, - ) - })?; + revise_host_record( + facade, + &scope, + KIMI_DISCOVERY_FRONTIER_KEY, + discovery_frontier, + byte_offset, + file_id, + ) + .await + .map_err(|outcome| { + crate::runtime::snapshot_observation::host_admission_error( + PROVIDER, outcome, + ) + })?; } } let deferred_units = outcome diff --git a/crates/tracedecay-sessions/src/runtime/hosts/pi.rs b/crates/tracedecay-sessions/src/runtime/hosts/pi.rs index 0a2415f6ab..fb3df87f1c 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/pi.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/pi.rs @@ -38,7 +38,8 @@ use crate::runtime::snapshot_observation::MAX_SNAPSHOT_METADATA_BYTES; use crate::runtime::source::{ FileDiscoveryLimit, FileDiscoveryReport, HostProviderCoverage, TranscriptDiscoveryBounds, TranscriptIngestError, TranscriptIngestResult, bound_path_list, canonical_framed_sha256, - jsonl_file_identity, persist_host_provider_coverage, run_blocking_transcript_section, + jsonl_file_identity, persist_host_provider_coverage, revise_host_record, + run_blocking_transcript_section, }; /// Environment override Pi reads for its agent directory. @@ -657,23 +658,21 @@ pub async fn capture_pi_observations( && !cancellation.is_cancelled() { let next_frontier = if discovery.reached_end { - Some(ParseOffset { - byte_offset: 0, - mtime: discovery_frontier.mtime.saturating_add(1), - file_id: 0, - }) + Some((0, 0)) } else { - last_discovered_entry.map(|entry| ParseOffset { - byte_offset: entry.sequence, - mtime: discovery_frontier.mtime.saturating_add(1), - file_id: entry.sequence, - }) + last_discovered_entry.map(|entry| (entry.sequence, entry.sequence)) }; - if let Some(next_frontier) = next_frontier { - facade - .advance_parse_offset(&scope, PI_DISCOVERY_FRONTIER_KEY, next_frontier) - .await - .map_err(admission_error)?; + if let Some((byte_offset, file_id)) = next_frontier { + revise_host_record( + facade, + &scope, + PI_DISCOVERY_FRONTIER_KEY, + discovery_frontier, + byte_offset, + file_id, + ) + .await + .map_err(admission_error)?; } } persist_coverage(facade, &scope, &outcome).await?; diff --git a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs index 944a4e96c2..2533dc8752 100644 --- a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs +++ b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs @@ -392,16 +392,15 @@ impl<'a> ProjectProviderRun<'a> { } else { HostProviderCoverage::Complete }; - if stored_coverage != Some(coverage) - && let Err(error) = persist_host_provider_coverage( - self.facade, - self.scope, - "codex", - coverage, - u64::from(coverage != HostProviderCoverage::Complete), - None, - ) - .await + if let Err(error) = persist_host_provider_coverage( + self.facade, + self.scope, + "codex", + coverage, + u64::from(coverage != HostProviderCoverage::Complete), + None, + ) + .await { outcome.add_failure(warn_transcript_catch_up_failure( "codex", diff --git a/crates/tracedecay-sessions/src/runtime/ingest/user_provider.rs b/crates/tracedecay-sessions/src/runtime/ingest/user_provider.rs index 75e5fb905d..bb9756c7cd 100644 --- a/crates/tracedecay-sessions/src/runtime/ingest/user_provider.rs +++ b/crates/tracedecay-sessions/src/runtime/ingest/user_provider.rs @@ -210,16 +210,15 @@ impl UserProviderUnit<'_, S> { } else { HostProviderCoverage::Complete }; - if stored_coverage != Some(coverage) - && let Err(coverage_error) = persist_host_provider_coverage( - self.facade, - &ObservationScopeV1::Profile, - "codex", - coverage, - u64::from(coverage != HostProviderCoverage::Complete), - None, - ) - .await + if let Err(coverage_error) = persist_host_provider_coverage( + self.facade, + &ObservationScopeV1::Profile, + "codex", + coverage, + u64::from(coverage != HostProviderCoverage::Complete), + None, + ) + .await { run.add_failure(warn_transcript_catch_up_failure( "codex", diff --git a/crates/tracedecay-sessions/src/runtime/source.rs b/crates/tracedecay-sessions/src/runtime/source.rs index 1e01b09df4..7ec1488cae 100644 --- a/crates/tracedecay-sessions/src/runtime/source.rs +++ b/crates/tracedecay-sessions/src/runtime/source.rs @@ -30,7 +30,7 @@ use tracedecay_store::{ParseOffset, TranscriptStoreError}; pub use tracedecay_domain::canonical_text::canonical_framed_sha256; pub use super::host_coverage::HostCoverageReason; -use crate::admission::HostAdmission; +use crate::admission::{HostAdmission, HostAdmissionOutcome}; pub use crate::runtime::shared::{NewRows, StoredCursor, TranscriptIngestStats}; use tracedecay_framing::{WireReadOutcome, read_bounded_to_string}; @@ -168,20 +168,47 @@ pub(super) async fn persist_host_provider_coverage( crate::runtime::snapshot_observation::host_admission_error(provider, outcome) })? .unwrap_or_default(); + revise_host_record( + admission, + scope, + &key, + current, + deferred_units, + coverage.file_id(reason), + ) + .await + .map_err(|outcome| { + crate::runtime::snapshot_observation::host_admission_error(provider, outcome) + }) +} + +/// Writes a host bookkeeping record whose readers consult only its +/// `byte_offset` and `file_id`. Its `mtime` is a revision that lets a changed +/// record move `byte_offset` backwards past the monotonic cursor guard. An +/// unchanged record stays as stored, so a sweep that finds nothing new +/// commits nothing. +pub(super) async fn revise_host_record( + admission: &(impl HostAdmission + ?Sized), + scope: &ObservationScopeV1, + key: &str, + current: ParseOffset, + byte_offset: u64, + file_id: u64, +) -> Result<(), HostAdmissionOutcome> { + if (current.byte_offset, current.file_id) == (byte_offset, file_id) { + return Ok(()); + } admission .advance_parse_offset( scope, - &key, + key, ParseOffset { - byte_offset: deferred_units, + byte_offset, mtime: current.mtime.saturating_add(1).max(1), - file_id: coverage.file_id(reason), + file_id, }, ) .await - .map_err(|outcome| { - crate::runtime::snapshot_observation::host_admission_error(provider, outcome) - }) } #[derive(Debug, Error)] From f3a6445bfa9da5f0afdeead22428527be77e409f Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 5 Oct 2026 03:27:53 -0700 Subject: [PATCH 2/2] fix(ci): repair all-target feature checks --- .../src/agents/context_scout/owner.rs | 30 ++++++++----------- crates/tracedecay/Cargo.toml | 6 +--- crates/tracedecay/benches/coverage/admin.rs | 5 ++-- crates/tracedecay/benches/coverage/mod.rs | 3 +- 4 files changed, 18 insertions(+), 26 deletions(-) diff --git a/crates/tracedecay-agent-hosts/src/agents/context_scout/owner.rs b/crates/tracedecay-agent-hosts/src/agents/context_scout/owner.rs index e06de7ab96..7cd16aea08 100644 --- a/crates/tracedecay-agent-hosts/src/agents/context_scout/owner.rs +++ b/crates/tracedecay-agent-hosts/src/agents/context_scout/owner.rs @@ -1012,26 +1012,20 @@ impl ProjectContextScoutOwnerV1 { let retired = match (&mutation, &receipt.result) { ( ContextScoutPublicMutationV1::Cancel { work }, - ContextScoutMutationResultV1::Cancel(outcome), - ) if matches!( - outcome, - ContextScoutDurableStoreOutcomeV1::Stored - | ContextScoutDurableStoreOutcomeV1::Duplicate - ) => - { - Some((*work, None)) - } + ContextScoutMutationResultV1::Cancel( + ContextScoutDurableStoreOutcomeV1::Stored + | ContextScoutDurableStoreOutcomeV1::Duplicate, + ), + ) => Some((*work, None)), ( ContextScoutPublicMutationV1::Delivery { work, .. }, - ContextScoutMutationResultV1::Delivery { outcome, receipt }, - ) if matches!( - outcome, - ContextScoutDurableStoreOutcomeV1::Stored - | ContextScoutDurableStoreOutcomeV1::Duplicate - ) => - { - Some((*work, Some(receipt.outcome))) - } + ContextScoutMutationResultV1::Delivery { + outcome: + ContextScoutDurableStoreOutcomeV1::Stored + | ContextScoutDurableStoreOutcomeV1::Duplicate, + receipt, + }, + ) => Some((*work, Some(receipt.outcome))), _ => None, }; if let Some((work, delivery)) = retired diff --git a/crates/tracedecay/Cargo.toml b/crates/tracedecay/Cargo.toml index ed0d8d3ece..2ced06ab52 100644 --- a/crates/tracedecay/Cargo.toml +++ b/crates/tracedecay/Cargo.toml @@ -1,5 +1,6 @@ [package] name = "tracedecay" +autobenches = false version.workspace = true publish = false edition.workspace = true @@ -373,11 +374,6 @@ path = "benches/large_repos.rs" harness = false required-features = ["test-transport"] -[[bench]] -name = "queries" -path = "benches/queries.rs" -required-features = ["test-transport"] - [[bench]] name = "session_temporal" path = "benches/session_temporal.rs" diff --git a/crates/tracedecay/benches/coverage/admin.rs b/crates/tracedecay/benches/coverage/admin.rs index 4548b204d3..d1713195ed 100644 --- a/crates/tracedecay/benches/coverage/admin.rs +++ b/crates/tracedecay/benches/coverage/admin.rs @@ -1046,7 +1046,7 @@ pub(crate) fn groups(ctx: &QueryContext, out: &mut Vec) { args["include_storage_health"] = json!(true); } if tool == "tracedecay_configuration_observed_state" { - Query::prepared_read(label, tool, args, super::configuration_read_prime) + Query::prepared_read(label, tool, args, i, super::configuration_read_prime) } else { rq(tool, label, args) } @@ -1084,11 +1084,12 @@ pub(crate) fn groups(ctx: &QueryContext, out: &mut Vec) { } out.push(ToolGroup { tool: "tracedecay_configuration_get", - queries: five(|_i| { + queries: five(|i| { Query::prepared_read( "config_get", "tracedecay_configuration_get", json!({"key": ctx.seeds.config_key.clone().unwrap_or_else(|| "diagnostics.prewarm.v1".into())}), + i, super::configuration_read_prime, ) }), diff --git a/crates/tracedecay/benches/coverage/mod.rs b/crates/tracedecay/benches/coverage/mod.rs index 014afd058f..a15602f7a9 100644 --- a/crates/tracedecay/benches/coverage/mod.rs +++ b/crates/tracedecay/benches/coverage/mod.rs @@ -13,7 +13,8 @@ mod work; pub(crate) use admin_fixture::verify_admin_fixture; pub(crate) use admin::{ - verify_context_scout_fixture, verify_github_stack_signal_fixture, verify_native_fixture, + finish_context_scout_fixture, verify_context_scout_fixture, verify_github_stack_signal_fixture, + verify_native_fixture, }; #[cfg(unix)] pub(crate) use code::prepare_source_reconciliation;