From 3369aadee25e13ed42f281c5a9c538645442409f Mon Sep 17 00:00:00 2001 From: Zack Jackson <25274700+ScriptedAlchemy@users.noreply.github.com> Date: Mon, 28 Sep 2026 04:08:48 +0000 Subject: [PATCH 01/18] fix(sessions): admit newest Codex rollouts before older days Shared catch-up charged every rollout file as directory work and then shrank that walk to the JSONL parse-slot count, so a large Codex home listed older days for minutes before emitting any session. Co-authored-by: Zack Jackson --- .../src/runtime/hosts/codex.rs | 6 +- .../src/runtime/hosts/codex/tests.rs | 61 +++++++++++++++++++ 2 files changed, 66 insertions(+), 1 deletion(-) diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex.rs index feac9b9cd3..8ba64656ef 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex.rs @@ -2103,7 +2103,6 @@ fn retained_scan_step( )?; match listed.next() { Some(entry) => { - directory_work += 1; let entry = entry.map_err(|source| TranscriptIngestError::ScanIo { operation: "read Codex transcript directory entry", path: dir.clone(), @@ -2117,7 +2116,12 @@ fn retained_scan_step( path: entry.path(), source, })?; + // Rollout files are not structural work. Charging them + // here spends the pass on names that are not child + // directories, so the newest sessions are not emitted + // until every older day has been listed. if file_type.is_dir() && !file_type.is_symlink() { + directory_work += 1; if *depth >= MAX_SCAN_DEPTH { return Err(TranscriptIngestError::ScanIo { operation: "traverse Codex transcript directory depth", diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs index 73954656f5..a7b085cf42 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs @@ -2275,6 +2275,67 @@ mod recent_first_discovery_tests { assert!(pass.report.paths.len() <= bounds.max_files); } + /// A dated tree whose older days hold more rollouts than one structural + /// pass can charge must still surface today's session immediately. Listing + /// those older files is not allowed to postpone the newest rollout. + #[tokio::test] + async fn codex_catch_up_surfaces_the_newest_rollout_before_older_days_are_listed() { + crate::runtime::jsonl_observation_admission::install_test_shared_jsonl_preparation_authority(); + let temp = TempDir::new().unwrap(); + let home = temp.path(); + for index in 0..800 { + write_dated_rollout(home, ("2026", "06", "01"), &format!("older-{index:04}")); + } + let newest = write_dated_rollout(home, ("2026", "09", "28"), "project-newest"); + let hub = CodexDiscoveryHub::default(); + hub.register("project", Some(home)); + let source = CodexSource::with_home(home); + let bounds = TranscriptDiscoveryBounds::default_walk(); + let mut frontier = CodexDiscoveryFrontier::initial(); + let mut surfaced = false; + for _ in 0..2 { + let pass = match hub + .discover("project", &source, bounds, frontier) + .await + .unwrap() + { + CodexDiscoveryDelivery::Ready(pass) => pass, + CodexDiscoveryDelivery::Waiting => { + panic!("a single catch-up consumer must not wait on its own scan") + } + }; + frontier = pass.next_frontier; + hub.acknowledge("project"); + if pass.report.paths.first() == Some(&newest) { + surfaced = true; + break; + } + } + assert!( + surfaced, + "catch-up must surface the newest rollout before it finishes listing older days" + ); + for _ in 0..64 { + if frontier.is_complete() { + break; + } + let pass = match hub + .discover("project", &source, bounds, frontier) + .await + .unwrap() + { + CodexDiscoveryDelivery::Ready(pass) => pass, + CodexDiscoveryDelivery::Waiting => continue, + }; + frontier = pass.next_frontier; + hub.acknowledge("project"); + } + assert!( + frontier.is_complete(), + "catch-up must finish the rollout sweep after the newest session is visible" + ); + } + #[test] #[cfg(unix)] fn codex_same_path_same_size_preserved_mtime_replacement_changes_epoch() { From 4c31655e77e5eea170d654db826a75c484690eb8 Mon Sep 17 00:00:00 2001 From: Zack Jackson <25274700+ScriptedAlchemy@users.noreply.github.com> Date: Mon, 28 Sep 2026 05:06:12 +0000 Subject: [PATCH 02/18] fix(sessions)!: return deferred import while catch-up continues sessions import waited for the historical cycle it had just queued. On a large Codex home that cycle does not finish inside the existing deadline, so the command exited timed out while admission was still moving. BREAKING CHANGE: A transcript import whose catch-up is still pending finishes as partial coverage with remaining work, instead of timed_out. Git sync still requires its own pass to finish. Co-authored-by: Zack Jackson --- .../src/sessions_cmd/session_sync.rs | 65 ++++- .../src/session_sync.rs | 250 +++++++++--------- .../session_sync/import_admission_tests.rs | 170 ++++++++++++ .../src/session_sync/work.rs | 32 +-- .../wake.rs | 111 ++++---- 5 files changed, 433 insertions(+), 195 deletions(-) create mode 100644 crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs diff --git a/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs b/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs index 242fd7f99a..9dd50f2c86 100644 --- a/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs +++ b/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs @@ -1,6 +1,8 @@ use std::path::Path; use tracedecay_contracts::retrieval::{AdminCliSessionSyncV1, AdminCliSurfaceRequestV1}; -use tracedecay_contracts::session_sync::SessionSyncSourceCoverageV1; +use tracedecay_contracts::session_sync::{ + SessionSyncCoverageV1, SessionSyncSourceCoverageV1, +}; use tracedecay_contracts::{IdempotencyKey, OperationTermination, RequestId}; use tracedecay_runtime_core::config::ProfileRoot; @@ -108,6 +110,19 @@ pub(super) fn session_sync_poll_state( ), } })?; + if session_import_deferred_progress( + label, + termination, + &coverage, + &failure_codes, + remaining_work, + ) { + println!( + "{label} scheduled ({}); historical catch-up has remaining work {remaining_work}", + operation_id.as_str() + ); + return Ok(SessionSyncPollState::Completed); + } if termination != OperationTermination::Completed || remaining_work > 0 { let termination = termination_label(termination); let detail = if failure_codes.is_empty() { @@ -144,6 +159,25 @@ fn termination_label(termination: OperationTermination) -> String { } } +fn session_import_deferred_progress( + label: &str, + termination: OperationTermination, + coverage: &[SessionSyncSourceCoverageV1], + failure_codes: &[String], + remaining_work: u64, +) -> bool { + label == "session import" + && termination == OperationTermination::Partial + && failure_codes.is_empty() + && remaining_work > 0 + && coverage.iter().all(|entry| { + matches!( + entry.coverage, + SessionSyncCoverageV1::Complete | SessionSyncCoverageV1::Partial { .. } + ) + }) +} + fn session_sync_remaining_work(coverage: &[SessionSyncSourceCoverageV1]) -> Option { if coverage.is_empty() { return None; @@ -379,6 +413,35 @@ mod tests { } } + #[test] + fn session_import_accepts_deferred_catch_up_without_treating_it_as_failure() { + let outcome = complete( + OperationTermination::Partial, + vec![SessionSyncCoverageV1::Partial { deferred_units: 1 }], + &[], + ); + + assert!(matches!( + session_sync_poll_state("session import", outcome).unwrap(), + SessionSyncPollState::Completed + )); + } + + #[test] + fn session_git_sync_still_rejects_unfinished_coverage() { + let error = session_sync_poll_state( + "session git sync", + complete( + OperationTermination::Partial, + vec![SessionSyncCoverageV1::Partial { deferred_units: 1 }], + &[], + ), + ) + .expect_err("git sync still requires the bounded pass to finish"); + + assert!(error.to_string().contains("remaining work")); + } + #[test] fn session_sync_reports_remaining_coverage_even_if_daemon_mislabels_completion() { let error = session_sync_poll_state( diff --git a/crates/tracedecay-session-runtime/src/session_sync.rs b/crates/tracedecay-session-runtime/src/session_sync.rs index ea9c41f0c9..ded403b8de 100644 --- a/crates/tracedecay-session-runtime/src/session_sync.rs +++ b/crates/tracedecay-session-runtime/src/session_sync.rs @@ -17,6 +17,7 @@ use tracedecay_contracts::{ }; use tracedecay_domain::{BrainId, ProjectId, SessionId, UserProfileId, UtcMicros}; +use crate::session_temporal_refresh_scheduler::wake::HistoricalAdmissionView; use tracedecay_global_db::GlobalDbGitCorrelationStore; use tracedecay_global_db::RegisteredGlobalDbLeaseV1; use tracedecay_runtime_core::background_cpu::ProcessBackgroundCpuV1; @@ -116,6 +117,11 @@ pub enum SessionSyncWorkResult { Interrupted(work::SessionSyncInterruption), } +struct ImportHistoryObservation { + coverage: Vec, + failure_codes: Vec, +} + struct SessionSyncTerminalMaterial { termination: OperationTermination, stats: SessionSyncStatsV1, @@ -569,49 +575,36 @@ impl DaemonSessionSyncService { } } SessionSyncWorkResult::Finished { - mut interruption, + interruption, committed, stats, coverage, source_frontiers, - mut failure_codes, + failure_codes, } => { - let mut interrupted = interruption.is_some(); + let interrupted = interruption.is_some(); let is_import = matches!( request.command(), SessionSyncCommandV1::ImportTranscripts(_) ); - let mut projection_current = false; let coverage_complete = !coverage.is_empty() && coverage.iter().all(|entry| entry.coverage.is_complete()); - if !projection_current - && !interrupted - && coverage_complete - && failure_codes.is_empty() - && is_import - { - match self - .await_import_projection(&context, &project_sessions, &request) - .await - { - Ok(()) => projection_current = true, - Err(Some(reason)) => { - interruption = Some(reason); - interrupted = true; - } - Err(None) => { - failure_codes - .push("session_temporal_projection_not_current".to_owned()); - } - } - } - let termination = completion_termination( + let mut termination = completion_termination( interruption.and_then(work::SessionSyncInterruption::termination), committed, &stats, coverage_complete, failure_codes.is_empty(), ); + if import_reports_deferred_progress( + is_import, + interrupted, + &failure_codes, + coverage_complete, + &coverage, + ) { + termination = OperationTermination::Partial; + } if self .persist_terminal( &context, @@ -628,7 +621,6 @@ impl DaemonSessionSyncService { .is_ok() && committed && !interrupted - && !projection_current { context.project_refresh.wake(); context.user_refresh.wake(); @@ -641,116 +633,82 @@ impl DaemonSessionSyncService { &self, context: &SessionSyncProjectContext, request: &SessionSyncRequestV1, - ) -> Result< - crate::session_temporal_refresh_scheduler::history::SessionHistoricalIngestProgress, - Option, - > { - let remaining_micros = request - .deadline() - .expires_at - .0 - .saturating_sub(now_micros().0); - let Ok(remaining_micros) = u64::try_from(remaining_micros) else { - return Err(Some(work::SessionSyncInterruption::TimedOut)); - }; - if remaining_micros == 0 { - return Err(Some(work::SessionSyncInterruption::TimedOut)); + project_sessions: &RegisteredGlobalDbLeaseV1, + ) -> Result> { + if let Some(interruption) = + self.observed_interruption(request.cancellation(), request.deadline()) + { + return Err(Some(interruption)); } - let timeout = Duration::from_micros(remaining_micros); - let history = async { - tokio::join!( - context - .project_refresh - .wake_history_and_wait_until_idle(timeout), - context - .user_refresh - .wake_history_and_wait_until_idle(timeout), - ) - }; - tokio::pin!(history); - let settled = tokio::select! { - settled = &mut history => settled, - interruption = self.wait_for_interruption(request) => { - return Err(Some(interruption)); + let project = context + .project_refresh + .observe_and_schedule_historical_admission(); + let user = context + .user_refresh + .observe_and_schedule_historical_admission(); + let mut failure_codes = Vec::new(); + push_historical_admission_failure(&project, &mut failure_codes); + push_historical_admission_failure(&user, &mut failure_codes); + let mut coverage = vec![ + historical_admission_coverage("project", &project), + historical_admission_coverage("profile", &user), + ]; + if failure_codes.is_empty() && coverage.iter().all(|entry| entry.coverage.is_complete()) { + match self + .import_projection_is_settled(context, project_sessions, request) + .await + { + Ok(true) => {} + Ok(false) => { + for entry in &mut coverage { + entry.coverage = SessionSyncCoverageV1::Partial { deferred_units: 1 }; + } + } + Err(Some(work::SessionSyncInterruption::TimedOut)) => { + for entry in &mut coverage { + entry.coverage = SessionSyncCoverageV1::Partial { deferred_units: 1 }; + } + } + Err(interruption) => return Err(interruption), } - }; - // Only the historical frontier is decided here. Projection currency is - // `await_import_projection`'s gate, which waits for the projection - // workers and then re-checks these same serving states and stores. This - // gate does not wait for them, so asserting them here reports - // `session_history_not_current` for a history that is current and whose - // projection has simply not drained yet, and skips the stage that would - // have waited for it. - if let (Some(project), Some(user)) = settled { - Ok( - crate::session_temporal_refresh_scheduler::history::SessionHistoricalIngestProgress { - stats: project.stats.merge(user.stats), - committed: project.committed || user.committed, - }, - ) - } else { - Err(None) } + Ok(ImportHistoryObservation { + coverage, + failure_codes, + }) } - async fn await_import_projection( + async fn import_projection_is_settled( &self, context: &SessionSyncProjectContext, project_sessions: &RegisteredGlobalDbLeaseV1, request: &SessionSyncRequestV1, - ) -> Result<(), Option> { - let remaining_micros = request - .deadline() - .expires_at - .0 - .saturating_sub(now_micros().0); - let Ok(remaining_micros) = u64::try_from(remaining_micros) else { - return Err(Some(work::SessionSyncInterruption::TimedOut)); - }; - if remaining_micros == 0 { - return Err(Some(work::SessionSyncInterruption::TimedOut)); - } - let timeout = Duration::from_micros(remaining_micros); - let projection = async { - tokio::join!( - context.project_refresh.wake_and_wait_until_idle(timeout), - context.user_refresh.wake_and_wait_until_idle(timeout), - ) - }; - tokio::pin!(projection); - let settled = tokio::select! { - settled = &mut projection => settled, - interruption = self.wait_for_interruption(request) => { - return Err(Some(interruption)); - } - }; + ) -> Result> { let project = context.project_refresh.status(); let user = context.user_refresh.status(); - if settled.0 - && settled.1 - && project.backlog == 0 - && user.backlog == 0 - && project.unavailable_reason.is_none() - && user.unavailable_reason.is_none() - && matches!( + if project.backlog > 0 + || user.backlog > 0 + || project.unavailable_reason.is_some() + || user.unavailable_reason.is_some() + || !matches!( context.project_refresh.serving_status().state, SessionProjectionServingState::Current ) - && matches!( + || !matches!( context.user_refresh.serving_status().state, SessionProjectionServingState::Current ) - && self - .projection_store_is_current(project_sessions, request) - .await? - && self - .projection_store_is_current(&context.user_sessions, request) - .await? { - Ok(()) - } else { - Err(None) + return Ok(false); + } + if !self + .projection_store_is_current(project_sessions, request) + .await? + { + return Ok(false); } + self.projection_store_is_current(&context.user_sessions, request) + .await } async fn projection_store_is_current( @@ -1096,6 +1054,9 @@ pub mod git_topology; mod project_lifecycle; pub mod work; +#[cfg(test)] +mod import_admission_tests; + pub use project_lifecycle::{SessionSyncProjectContext, SessionSyncTaskV1}; #[cfg(any(test, feature = "test-helpers"))] @@ -1343,6 +1304,59 @@ fn decode_matching_journal( Ok(journal) } +fn import_reports_deferred_progress( + is_import: bool, + interrupted: bool, + failure_codes: &[String], + coverage_complete: bool, + coverage: &[SessionSyncSourceCoverageV1], +) -> bool { + is_import + && !interrupted + && failure_codes.is_empty() + && !coverage_complete + && !coverage.is_empty() + && coverage.iter().all(|entry| { + matches!( + entry.coverage, + SessionSyncCoverageV1::Complete | SessionSyncCoverageV1::Partial { .. } + ) + }) +} + +fn historical_admission_coverage( + store_scope: &str, + view: &HistoricalAdmissionView, +) -> SessionSyncSourceCoverageV1 { + let coverage = match view { + HistoricalAdmissionView::Current => SessionSyncCoverageV1::Complete, + HistoricalAdmissionView::InProgress + | HistoricalAdmissionView::Blocked { .. } + | HistoricalAdmissionView::Unavailable => { + SessionSyncCoverageV1::Partial { deferred_units: 1 } + } + }; + SessionSyncSourceCoverageV1 { + store_scope: store_scope.to_owned(), + coverage, + } +} + +fn push_historical_admission_failure( + view: &HistoricalAdmissionView, + failure_codes: &mut Vec, +) { + match view { + HistoricalAdmissionView::Current | HistoricalAdmissionView::InProgress => {} + HistoricalAdmissionView::Blocked { reason_code } => { + failure_codes.push(reason_code.clone()); + } + HistoricalAdmissionView::Unavailable => { + failure_codes.push("session_history_not_current".to_owned()); + } + } +} + fn completion_termination( requested: Option, committed: bool, diff --git a/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs b/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs new file mode 100644 index 0000000000..ae827a5bc5 --- /dev/null +++ b/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs @@ -0,0 +1,170 @@ +use std::num::NonZeroUsize; +use std::sync::Arc; +use std::time::Duration; + +use tracedecay_contracts::session_sync::{ + SessionSyncCommandV1, SessionSyncControlV1, SessionSyncOutcomeV1, SessionSyncRequestV1, + SessionSyncScopeV1, SessionSyncServicePort, SessionTranscriptImportV1, +}; +use tracedecay_contracts::{ + CancellationSignal, Deadline, IdempotencyKey, OperationTermination, RequestId, now_micros, +}; +use tracedecay_domain::{ProjectId, UtcMicros}; +use tracedecay_global_db::tests::harness::{HostAdmissionScope, HostAdmissionTestRuntimeV1}; +use tracedecay_runtime_core::background_cpu::ProcessBackgroundCpuV1; +use tracedecay_runtime_core::config::ProfileRoot; + +use crate::session_sync::{DaemonSessionSyncConfig, DaemonSessionSyncService}; +use crate::session_temporal_refresh_scheduler::SessionTemporalRefreshWake; +use crate::session_temporal_refresh_scheduler::history::SessionHistoricalIngestOutcome; +use crate::session_temporal_refresh_scheduler::wake::SessionTemporalRefreshWakeState; + +#[derive(Clone, Copy)] +enum CatchUp { + Pending, + Current, + Blocked, +} + +fn catch_up_state(kind: CatchUp) -> Arc { + let state = Arc::new(SessionTemporalRefreshWakeState::default()); + state.mark_running(); + match kind { + CatchUp::Pending => state.mark_history_pending(), + CatchUp::Current => { + state.record_history_outcome(SessionHistoricalIngestOutcome::Complete); + } + CatchUp::Blocked => { + state.record_history_outcome(SessionHistoricalIngestOutcome::Blocked { + reason_code: "invalid_observation_contract", + made_progress: false, + }); + } + } + state +} + +fn bound_wake(state: &Arc) -> SessionTemporalRefreshWake { + let wake = SessionTemporalRefreshWake::unavailable(); + wake.bind(state); + wake +} + +async fn import_receipt( + kind: CatchUp, + label: &str, +) -> tracedecay_contracts::session_sync::SessionSyncCompletionReceiptV1 { + let root = tempfile::tempdir().expect("import fixture directory"); + let project_root = root.path().join("project"); + std::fs::create_dir_all(&project_root).expect("project directory"); + let project_id = ProjectId::new(format!("project.{label}")).expect("project id"); + let runtime = + HostAdmissionTestRuntimeV1::project(root.path(), &project_root, project_id.clone()) + .await + .expect("registered session stores"); + let project_sessions = runtime + .registered_database_lease(HostAdmissionScope::Project) + .expect("project sessions"); + let profile_sessions = runtime + .registered_database_lease(HostAdmissionScope::Profile) + .expect("profile sessions"); + let brain_id = project_sessions.binding().shard_id.brain_id.clone(); + let profile_id = project_sessions.binding().shard_id.profile_id.clone(); + let project_state = catch_up_state(kind); + let user_state = catch_up_state(kind); + let service = DaemonSessionSyncService::default(); + service + .register_project(DaemonSessionSyncConfig { + brain_id, + profile_id: profile_id.clone(), + project_id: project_id.clone(), + profile_root: root.path().to_path_buf(), + project_root: project_root.clone(), + transcript_source_profile: ProfileRoot::new(root.path().to_path_buf()), + project_sessions, + user_sessions: profile_sessions.clone(), + registry: profile_sessions, + background_cpu: Arc::new(ProcessBackgroundCpuV1::new(NonZeroUsize::MIN)), + startup_import: false, + project_refresh: bound_wake(&project_state), + user_refresh: bound_wake(&user_state), + }) + .await + .expect("session sync project"); + let scope = SessionSyncScopeV1::new(project_id, profile_id); + let request = SessionSyncRequestV1::new( + RequestId::new(format!("session-sync.{label}")).expect("operation id"), + IdempotencyKey::new(format!("session-sync.{label}")).expect("idempotency key"), + scope.clone(), + Deadline::new(UtcMicros(now_micros().0.saturating_add(60_000_000))).expect("deadline"), + CancellationSignal::active(format!("session-sync.{label}")).expect("cancellation"), + SessionSyncCommandV1::ImportTranscripts(SessionTranscriptImportV1::all_hosts()), + ); + let accepted = SessionSyncServicePort::execute(&service, request).await; + let SessionSyncOutcomeV1::Accepted(admission) = accepted else { + panic!("import was not admitted: {accepted:?}"); + }; + let started = tokio::time::Instant::now(); + let control = SessionSyncControlV1::new(scope, admission.idempotency_key); + loop { + match SessionSyncServicePort::status(&service, control.clone()).await { + SessionSyncOutcomeV1::Complete(receipt) => { + assert!( + started.elapsed() < Duration::from_secs(1), + "import consumed its observation bound instead of returning the settled catch-up" + ); + return receipt; + } + SessionSyncOutcomeV1::Accepted(_) | SessionSyncOutcomeV1::Joined(_) => { + assert!( + started.elapsed() < Duration::from_secs(1), + "import stayed pending while catch-up state was already known" + ); + tokio::time::sleep(Duration::from_millis(10)).await; + } + other => panic!("import ended without a coverage receipt: {other:?}"), + } + } +} + +fn remaining_work( + receipt: &tracedecay_contracts::session_sync::SessionSyncCompletionReceiptV1, +) -> u64 { + receipt + .coverage + .iter() + .map(|entry| entry.coverage.remaining_work()) + .fold(0, u64::saturating_add) +} + +#[tokio::test] +async fn import_reports_deferred_progress_while_historical_catch_up_is_still_pending() { + let receipt = import_receipt(CatchUp::Pending, "import-pending").await; + + assert_ne!(receipt.termination, OperationTermination::TimedOut); + assert_eq!(receipt.termination, OperationTermination::Partial); + assert!(receipt.failure_codes.is_empty()); + assert!(remaining_work(&receipt) > 0); +} + +#[tokio::test] +async fn import_completes_when_historical_catch_up_is_already_current() { + let receipt = import_receipt(CatchUp::Current, "import-current").await; + + assert_eq!(receipt.termination, OperationTermination::Completed); + assert!(receipt.failure_codes.is_empty()); + assert_eq!(remaining_work(&receipt), 0); +} + +#[tokio::test] +async fn import_keeps_a_blocked_catch_up_as_a_failure() { + let receipt = import_receipt(CatchUp::Blocked, "import-blocked").await; + + assert_eq!(receipt.termination, OperationTermination::Failed); + assert!( + receipt + .failure_codes + .iter() + .any(|code| code == "invalid_observation_contract") + ); +} diff --git a/crates/tracedecay-session-runtime/src/session_sync/work.rs b/crates/tracedecay-session-runtime/src/session_sync/work.rs index 80b21f6e6d..2df4675c50 100644 --- a/crates/tracedecay-session-runtime/src/session_sync/work.rs +++ b/crates/tracedecay-session-runtime/src/session_sync/work.rs @@ -491,23 +491,23 @@ impl SessionSyncProjectContext { request: &SessionSyncRequestV1, project_sessions: RegisteredGlobalDbLeaseV1, ) -> SessionSyncWorkResult { - let history = match service.await_import_history(self, request).await { - Ok(progress) => Some(progress), + let observation = match service + .await_import_history(self, request, &project_sessions) + .await + { + Ok(observation) => observation, Err(Some(interruption)) => { return SessionSyncWorkResult::Interrupted(interruption); } - Err(None) => None, + Err(None) => super::ImportHistoryObservation { + coverage: super::transcript_import_requested_coverage(), + failure_codes: vec!["session_history_not_current".to_owned()], + }, }; let pass = async { - let stats = import_transcript_stats( - history.map_or_else(Default::default, |progress| progress.stats), - ); - let mut coverage = super::transcript_import_requested_coverage(); - if history.is_some() { - for source in &mut coverage { - source.coverage = SessionSyncCoverageV1::Complete; - } - } + let stats = + import_transcript_stats(tracedecay_sessions::TranscriptIngestStats::default()); + let coverage = observation.coverage.clone(); let source_frontiers = hotpath::future!( service.persist_progress( self, @@ -530,15 +530,11 @@ impl SessionSyncProjectContext { } }; let (stats, coverage, source_frontiers) = outcomes; - let mut failure_codes = Vec::new(); - if history.is_none() { - failure_codes.push("session_history_not_current".to_owned()); - } + let mut failure_codes = observation.failure_codes; if source_frontiers.is_err() { failure_codes.push("session_sync_frontier_persist_failed".to_owned()); } - let committed = history.is_some_and(|progress| progress.committed) - || stats != SessionSyncStatsV1::default(); + let committed = stats != SessionSyncStatsV1::default(); SessionSyncWorkResult::Finished { interruption: interrupted, committed, diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs index d40bce8c2b..03781b95a5 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs @@ -844,10 +844,40 @@ pub struct SessionTemporalRefreshWake { route: Arc, } -#[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub struct SessionHistoricalRefreshReceipt { - pub stats: tracedecay_sessions::TranscriptIngestStats, - pub committed: bool, +/// What historical catch-up has already settled, read before this call's wake. +#[derive(Clone, Debug, Eq, PartialEq)] +pub(crate) enum HistoricalAdmissionView { + Current, + InProgress, + Blocked { reason_code: String }, + Unavailable, +} + +fn historical_admission_view(status: &SessionProjectionServingStatus) -> HistoricalAdmissionView { + match &status.state { + SessionProjectionServingState::Current => HistoricalAdmissionView::Current, + SessionProjectionServingState::Stale { reason } => match reason { + SessionProjectionStaleReason::HistoricalConvergence + | SessionProjectionStaleReason::HistoricalRetry { .. } => { + HistoricalAdmissionView::InProgress + } + SessionProjectionStaleReason::HistoricalBlocked { reason_code } => { + HistoricalAdmissionView::Blocked { + reason_code: reason_code.clone(), + } + } + }, + SessionProjectionServingState::Unavailable { reason } => match reason { + SessionProjectionUnavailableReason::WorkerRecovering => { + HistoricalAdmissionView::InProgress + } + SessionProjectionUnavailableReason::WorkerMissing + | SessionProjectionUnavailableReason::WorkerStalled + | SessionProjectionUnavailableReason::WorkerStopped => { + HistoricalAdmissionView::Unavailable + } + }, + } } impl SessionTemporalRefreshWake { @@ -932,62 +962,27 @@ impl SessionTemporalRefreshWake { } } - /// Requests a fresh bounded historical-ingest cycle from the retained - /// owner and waits through its existing continuation passes. + /// Schedules another historical pass and reports the admission already + /// settled by the worker. + /// + /// Callers that must return inside a bound, such as `sessions import`, + /// use this instead of waiting for the pass just queued. Waiting for that + /// pass is historical convergence: a large Codex home does not finish it + /// before the import deadline, so the operation was recorded as timed out + /// while catch-up was still admitting rollouts. #[hotpath::skip] - pub async fn wake_history_and_wait_until_idle( - &self, - timeout: std::time::Duration, - ) -> Option { - let state = self.target()?; - if state.cancelled.load(Ordering::Acquire) { - return None; - } - let requested = state - .history_requested_sequence - .fetch_add(1, Ordering::AcqRel) - .saturating_add(1); - let before_sessions = state.history_sessions_upserted.load(Ordering::Acquire); - let before_messages = state.history_messages_upserted.load(Ordering::Acquire); - let before_commits = state.history_commit_count.load(Ordering::Acquire); - state.mark_history_pending(); - state.wake_history(); - let deadline = tokio::time::Instant::now() + timeout; - loop { - let idle = hotpath::future!( - enabled_idle_notification(&state), - label = "daemon.scheduler.session_temporal.history_idle_wait" - ); - let settled = state.history_completed_sequence.load(Ordering::Acquire) >= requested; - if settled { - if !matches!( - state - .telemetry - .lock() - .unwrap_or_else(PoisonError::into_inner) - .historical_state, - SessionHistoricalServingState::Current - ) { - return None; - } - return Some(SessionHistoricalRefreshReceipt { - stats: tracedecay_sessions::TranscriptIngestStats { - sessions_upserted: state - .history_sessions_upserted - .load(Ordering::Acquire) - .saturating_sub(before_sessions), - messages_upserted: state - .history_messages_upserted - .load(Ordering::Acquire) - .saturating_sub(before_messages), - }, - committed: state.history_commit_count.load(Ordering::Acquire) > before_commits, - }); - } - if tokio::time::timeout_at(deadline, idle).await.is_err() { - return None; - } + pub(crate) fn observe_and_schedule_historical_admission(&self) -> HistoricalAdmissionView { + let Some(state) = self.target() else { + return HistoricalAdmissionView::Unavailable; + }; + let view = historical_admission_view(&state.serving_status()); + if matches!( + view, + HistoricalAdmissionView::Current | HistoricalAdmissionView::InProgress + ) { + state.wake_history(); } + view } pub fn status(&self) -> SessionTemporalRefreshWorkerStatus { From 283b48ef38361c9c7ee8caec6ce4143cedf57d0d Mon Sep 17 00:00:00 2001 From: Zack Jackson <25274700+ScriptedAlchemy@users.noreply.github.com> Date: Mon, 28 Sep 2026 09:30:14 +0000 Subject: [PATCH 03/18] fix(sessions): publish newest Codex day before older days The project Codex pass wrote a coverage cursor for every out-of-scope rollout before it returned, so projection could not activate the newest day until that sweep finished. Yield once an in-scope window is admitted, and finish its projection backlog before the next history window. Co-authored-by: Zack Jackson --- .../wake.rs | 18 ++ .../worker.rs | 174 ++++++++++++++-- .../src/runtime/hosts/codex.rs | 13 ++ .../src/runtime/hosts/codex/tests.rs | 185 ++++++++++++++++++ .../src/runtime/ingest/project_provider.rs | 19 +- 5 files changed, 395 insertions(+), 14 deletions(-) diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs index 03781b95a5..36afaa9bc3 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs @@ -126,6 +126,9 @@ macro_rules! define_wake_state { pub(super) recovery_cycle_pending: std::sync::Mutex>, pub(super) busy: AtomicBool, pub(super) history_retry_pending: AtomicBool, + /// The last history window made progress and its projection backlog + /// is still unpublished. The next history window waits. + pub(super) history_after_projection: AtomicBool, pub(super) pass_count: std::sync::atomic::AtomicUsize, pub(super) history_requested_sequence: std::sync::atomic::AtomicUsize, pub(super) history_completed_sequence: std::sync::atomic::AtomicUsize, @@ -160,6 +163,7 @@ impl Default for SessionTemporalRefreshWakeState { recovery_cycle_pending: std::sync::Mutex::new(VecDeque::new()), busy: AtomicBool::new(false), history_retry_pending: AtomicBool::new(false), + history_after_projection: AtomicBool::new(false), pass_count: std::sync::atomic::AtomicUsize::new(0), history_requested_sequence: std::sync::atomic::AtomicUsize::new(0), history_completed_sequence: std::sync::atomic::AtomicUsize::new(0), @@ -404,6 +408,20 @@ impl SessionTemporalRefreshWakeState { pub fn clear_worker_activity_instrumentation(&self) { self.mark_worker_idle(); self.update_history_retry_state(false); + self.release_history_for_projection(); + } + + pub fn hold_history_for_projection(&self) { + self.history_after_projection.store(true, Ordering::Release); + } + + pub fn release_history_for_projection(&self) { + self.history_after_projection + .store(false, Ordering::Release); + } + + pub fn history_held_for_projection(&self) -> bool { + self.history_after_projection.load(Ordering::Acquire) } pub fn history_retry_pending(&self) -> bool { diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs index 3a58cae896..3de29666cb 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs @@ -111,9 +111,8 @@ enum HistoryContinuation { Settled, } -/// A Codex catch-up yields after one rollout. That yield is progress, so the -/// continuation must not pay the no-progress retry delay or a corpus larger -/// than one window misses the import deadline. +/// A progressing history window continues immediately. The no-progress retry +/// delay is only for a window that admitted nothing. fn history_continuation(outcome: Option) -> HistoryContinuation { match outcome { Some(SessionHistoricalIngestOutcome::Pending { @@ -124,6 +123,59 @@ fn history_continuation(outcome: Option) -> Hist } } +/// What the worker does after one pass while history still owes another window. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum WindowFollowUp { + /// The admitted window is not searchable yet. Run another projection pass + /// before opening the next history window. + PublishBeforeNextHistory, + /// Publication finished or stalled. Open the next history window now. + ContinueHistory, + /// The history window made no progress. Back off. + BackoffHistory, + /// History does not own the next pass. + Settle, +} + +fn projection_published_work(report: &SessionTemporalRefreshPassReport) -> bool { + report.begun > 0 + || report.joined > 0 + || report.projected_batches > 0 + || report.completed > 0 + || report.failed > 0 +} + +fn projection_still_unpublished(report: &SessionTemporalRefreshPassReport) -> bool { + report.backlog.is_some_and(|backlog| backlog > 0) || report.saturated +} + +/// An in-scope Codex window becomes searchable only after its projection +/// backlog reaches an active generation. The next history window, which on a +/// large corpus is the out-of-scope cursor sweep, waits until that publish +/// finishes or a projection pass stops moving. +fn window_follow_up( + history: HistoryContinuation, + publication_pending: bool, + projection_moved: bool, + holding: bool, +) -> WindowFollowUp { + match history { + HistoryContinuation::Backoff => WindowFollowUp::BackoffHistory, + HistoryContinuation::Immediate | HistoryContinuation::Settled => { + let owed = publication_pending + && projection_moved + && (history == HistoryContinuation::Immediate || holding); + if owed { + WindowFollowUp::PublishBeforeNextHistory + } else if history == HistoryContinuation::Immediate || holding { + WindowFollowUp::ContinueHistory + } else { + WindowFollowUp::Settle + } + } + } +} + pub(super) async fn run_session_temporal_refresh_scheduler( database: RegisteredGlobalDbLeaseV1, state: Arc, @@ -227,8 +279,20 @@ pub(super) async fn run_session_temporal_refresh_scheduler( { state.complete_history_sequence(sequence); } + let holding_publication = state.history_held_for_projection(); + // A projection-only iteration that is finishing the admitted window + // must not spend the pass on summary convergence. History did not + // run, so the real outcome is `None`, which would otherwise admit + // the full page. + let convergence_outcome = if holding_publication && history_outcome.is_none() { + Some(SessionHistoricalIngestOutcome::Pending { + made_progress: true, + }) + } else { + history_outcome + }; let convergence_admission = - lcm_convergence_admission(history_outcome, &mut history_priority_passes); + lcm_convergence_admission(convergence_outcome, &mut history_priority_passes); // Derived from the admission so the pass report can never disagree // with what this pass actually ran. let history_needs_another_pass = !matches!( @@ -413,22 +477,35 @@ pub(super) async fn run_session_temporal_refresh_scheduler( } else if history_needs_another_pass { state.mark_running(); retry_attempt = 0; - match history_continuation(history_outcome) { - // A bounded window that admitted work already yielded. The - // next window is continuation of that import, not a failure - // retry: the 250ms backoff below is only for passes that - // made no progress. Sleeping on every successful window - // makes a multi-window corpus miss the import deadline. - HistoryContinuation::Immediate => { + match window_follow_up( + history_continuation(history_outcome), + projection_still_unpublished(&report), + projection_published_work(&report), + holding_publication, + ) { + // The admitted window is not searchable until this backlog + // is published. Opening the next history window first is + // the out-of-scope cursor sweep, and search stays empty + // until that sweep ends. + WindowFollowUp::PublishBeforeNextHistory => { + state.hold_history_for_projection(); + state.update_history_retry_state(false); + state.requeue_projection(); + tokio::task::yield_now().await; + } + WindowFollowUp::ContinueHistory => { + state.release_history_for_projection(); state.update_history_retry_state(false); state.wake_history(); } - HistoryContinuation::Backoff => { + WindowFollowUp::BackoffHistory => { + state.release_history_for_projection(); state.update_history_retry_state(true); } - HistoryContinuation::Settled => {} + WindowFollowUp::Settle => {} } } else { + state.release_history_for_projection(); if history_outcome.is_some() { state.update_history_retry_state(false); } @@ -1263,6 +1340,77 @@ mod tests { ); } + #[test] + fn an_admitted_window_publishes_before_the_next_history_window() { + let unpublished = SessionTemporalRefreshPassReport { + projected_batches: 1, + backlog: Some(84), + ..SessionTemporalRefreshPassReport::default() + }; + assert_eq!( + window_follow_up( + HistoryContinuation::Immediate, + projection_still_unpublished(&unpublished), + projection_published_work(&unpublished), + false, + ), + WindowFollowUp::PublishBeforeNextHistory, + "a progressing history window with a projection backlog must not open the next window" + ); + + let still_moving = SessionTemporalRefreshPassReport { + completed: 16, + backlog: Some(68), + ..SessionTemporalRefreshPassReport::default() + }; + assert_eq!( + window_follow_up( + HistoryContinuation::Settled, + projection_still_unpublished(&still_moving), + projection_published_work(&still_moving), + true, + ), + WindowFollowUp::PublishBeforeNextHistory, + "projection-only passes keep the hold while the admitted window is still unpublished" + ); + + let published = SessionTemporalRefreshPassReport { + completed: 16, + backlog: Some(0), + ..SessionTemporalRefreshPassReport::default() + }; + assert_eq!( + window_follow_up( + HistoryContinuation::Settled, + projection_still_unpublished(&published), + projection_published_work(&published), + true, + ), + WindowFollowUp::ContinueHistory, + "a published window releases the next history window" + ); + + let stalled = SessionTemporalRefreshPassReport { + backlog: Some(68), + ..SessionTemporalRefreshPassReport::default() + }; + assert_eq!( + window_follow_up( + HistoryContinuation::Settled, + projection_still_unpublished(&stalled), + projection_published_work(&stalled), + true, + ), + WindowFollowUp::ContinueHistory, + "a projection pass that moves nothing must not spin ahead of history" + ); + + assert_eq!( + window_follow_up(HistoryContinuation::Backoff, true, true, true), + WindowFollowUp::BackoffHistory + ); + } + #[test] fn pending_history_windows_take_precedence_over_derived_summaries() { assert!(!history_allows_summary_convergence(Some( diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex.rs index 8ba64656ef..83cf41b24a 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex.rs @@ -113,6 +113,19 @@ pub use observation::{ try_admit_codex_jsonl_observations_for_project_with_admission_and_cancellation, }; +/// Project membership from a rollout's leading `session_meta` cwd. +/// +/// `None` means the header could not be read. Callers may leave a rollout for +/// a later pass only on a definitive [`ProjectMembership::NoMatch`]; an +/// `Unknown` git timeout stays in the current pass. +pub(crate) fn codex_rollout_project_membership( + path: &Path, + project_root: &Path, +) -> Option { + let meta = session_meta(path)?; + Some(TranscriptScopeMatcher::project(project_root).membership(Some(&meta.cwd))) +} + const PROVIDER: &str = "codex"; /// `~/.codex/sessions/YYYY/MM/DD/rollout-*.jsonl` → date dirs add depth. const MAX_SCAN_DEPTH: u8 = 6; diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs index a7b085cf42..a92bbbfd83 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs @@ -1002,6 +1002,191 @@ mod goal_event_tests { ); } + /// The newest in-project day is admitted, and its message text is durable, + /// before any older out-of-project day is opened. + /// + /// Search reads the published projection of those admitted messages. A pass + /// that keeps walking older days writes their coverage cursors before that + /// projection can run, so the newest day stays invisible until the older + /// sweep finishes. The follow-up pass must still open the older day, or + /// the yield would retry the same page forever. + #[tokio::test] + async fn newest_day_messages_are_durable_before_older_days_are_opened() { + crate::runtime::jsonl_observation_admission::install_test_shared_jsonl_preparation_authority(); + let temp = tempfile::tempdir().unwrap(); + let home = temp.path().canonicalize().unwrap(); + let project = home.join("project"); + let other = home.join("other"); + std::fs::create_dir_all(&project).unwrap(); + std::fs::create_dir_all(&other).unwrap(); + let marker = "NEWEST_DAY_MARKER cobalt orchard scheduler is ready"; + let newest = [ + ("2026", "08", "29", "newest-a"), + ("2026", "08", "29", "newest-b"), + ]; + let older = [ + ("2026", "08", "28", "older-a"), + ("2026", "08", "28", "older-b"), + ]; + for (year, month, day, session_id) in newest { + write_scoped_rollout(&home, year, month, day, session_id, &project, marker); + } + for (year, month, day, session_id) in older { + write_scoped_rollout( + &home, + year, + month, + day, + session_id, + &other, + "OLDER_DAY_MARKER should stay unopened", + ); + } + + let project_id = ProjectId::new("project-newest-before-older").unwrap(); + let scope = ObservationScopeV1::Project { + project_id: project_id.clone(), + }; + let admission = MemoryHostAdmission::default(); + let cancellation = ObservationCancellation::default(); + let run_pass = || { + with_transcript_source_profile( + tracedecay_runtime_core::config::ProfileRoot::under_home(home.clone()), + ProjectProviderRun { + project_root: &project, + project_id: &project_id, + facade: &admission, + scope: &scope, + candidate: SessionProvider::Codex, + max_new_bytes: u64::MAX, + cancellation: &cancellation, + codex_discovery: None, + } + .run_codex(), + ) + }; + + let first = run_pass().await; + assert!(first.failures.is_empty(), "{:?}", first.failures); + let admitted = session_ids_of(&admission.observations()); + assert_eq!( + admitted, + ["newest-a".to_owned(), "newest-b".to_owned()] + .into_iter() + .collect::>() + ); + assert!( + admission.observations().iter().any(|stored| { + let envelope: CanonicalObservationEnvelopeV1 = + serde_json::from_value(stored.observation().payload().clone()).unwrap(); + envelope.facts().iter().any(|fact| { + matches!( + fact, + CanonicalObservationFactV1::Message { content, .. } + if content.as_str() == Some(marker) + ) + }) + }), + "the newest day's message text must be durable before older days are opened" + ); + for session_id in ["older-a", "older-b"] { + let source = + crate::runtime::hosts::codex::codex_observation_source_v2(session_id).unwrap(); + assert!( + admission + .get_source_cursor(&source, &scope) + .await + .unwrap() + .is_none(), + "{session_id} was opened in the pass that admitted the newest day" + ); + } + assert_eq!( + read_host_provider_coverage(&admission, &scope, "codex") + .await + .unwrap(), + Some(HostProviderCoverage::Partial) + ); + + let mut older_opened = false; + for _ in 0..3 { + let outcome = run_pass().await; + assert!(outcome.failures.is_empty(), "{:?}", outcome.failures); + let source = + crate::runtime::hosts::codex::codex_observation_source_v2("older-a").unwrap(); + if admission + .get_source_cursor(&source, &scope) + .await + .unwrap() + .is_some() + { + older_opened = true; + break; + } + } + assert!( + older_opened, + "deferring the older day must still open it on a later pass" + ); + assert_eq!( + session_ids_of(&admission.observations()), + ["newest-a".to_owned(), "newest-b".to_owned()] + .into_iter() + .collect::>(), + "out-of-project days stay out of the project observation set" + ); + } + + fn write_scoped_rollout( + home: &std::path::Path, + year: &str, + month: &str, + day: &str, + session_id: &str, + cwd: &std::path::Path, + message: &str, + ) { + let directory = home + .join(".codex/sessions") + .join(year) + .join(month) + .join(day); + std::fs::create_dir_all(&directory).unwrap(); + let lines = [ + json!({ + "timestamp": format!("{year}-{month}-{day}T12:00:00.000Z"), + "type": "session_meta", + "payload": {"id": session_id, "cwd": cwd} + }), + json!({ + "timestamp": format!("{year}-{month}-{day}T12:00:01.000Z"), + "type": "event_msg", + "payload": {"type": "user_message", "message": message} + }), + ]; + std::fs::write( + directory.join(format!("rollout-{session_id}.jsonl")), + lines + .iter() + .map(ToString::to_string) + .collect::>() + .join("\n") + + "\n", + ) + .unwrap(); + } + + fn session_ids_of(observations: &[tracedecay_store::StoredObservation]) -> BTreeSet { + observations + .iter() + .map(|stored| { + let envelope: CanonicalObservationEnvelopeV1 = + serde_json::from_value(stored.observation().payload().clone()).unwrap(); + envelope.relations().session_id().as_str().to_owned() + }) + .collect() + } + #[tokio::test] async fn legacy_current_message_migration_records_receipted_duplicate_coverage() { crate::runtime::jsonl_observation_admission::install_test_shared_jsonl_preparation_authority(); diff --git a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs index d26125dcc9..168ebb89cb 100644 --- a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs +++ b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs @@ -8,7 +8,7 @@ use tracedecay_domain::{ObservationScopeV1, ProjectId}; use crate::admission::HostAdmission; use crate::observation::ObservationCancellation; -use crate::runtime::shared::TranscriptIngestStats; +use crate::runtime::shared::{ProjectMembership, TranscriptIngestStats}; use crate::runtime::source::{ HostCoverageReason, HostProviderCoverage, TranscriptDiscoveryBounds, persist_codex_history_frontier, persist_host_provider_coverage, read_codex_history_frontier, @@ -262,6 +262,10 @@ impl<'a> ProjectProviderRun<'a> { let mut remaining = self.max_new_bytes; let mut deferred = discovery.is_truncated(); let mut frontier_committable = true; + // Out-of-scope rollouts consume no byte budget, so a newest-first page + // would otherwise write a cursor for every older day before this pass + // returns and the admitted window can be projected into search. + let mut persisted_in_scope = false; let mut outcome = ProviderRunOutcome::bounded(TranscriptIngestStats::default(), 0, false); for path in &discovery.paths { if remaining == 0 { @@ -274,6 +278,16 @@ impl<'a> ProjectProviderRun<'a> { frontier_committable = false; break; } + if persisted_in_scope + && run_blocking_transcript_section(|| { + codex::codex_rollout_project_membership(path, self.project_root) + == Some(ProjectMembership::NoMatch) + }) + { + deferred = true; + frontier_committable = false; + break; + } match codex::try_admit_codex_jsonl_observations_for_project_window( path, self.project_root, @@ -285,6 +299,9 @@ impl<'a> ProjectProviderRun<'a> { .await { Ok(progress) => { + if progress.frames_persisted > 0 { + persisted_in_scope = true; + } deferred |= progress.source_deferred || progress.bytes_consumed > remaining; frontier_committable &= !progress.source_deferred && progress.bytes_consumed <= remaining; From bc7d42bc85be9450a0e842b721a5de6a76e75a42 Mon Sep 17 00:00:00 2001 From: Zack Jackson <25274700+ScriptedAlchemy@users.noreply.github.com> Date: Mon, 28 Sep 2026 09:48:46 +0000 Subject: [PATCH 04/18] fix(sessions): publish a yielded Codex window before backoff A newest-day yield is retryable backpressure. The worker treated that as a no-progress backoff and opened the older-day sweep after one projection slice, leaving the newest day on a building generation. Co-authored-by: Zack Jackson --- .../worker.rs | 51 ++++++++++++++----- 1 file changed, 38 insertions(+), 13 deletions(-) diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs index 3de29666cb..eb66406d37 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs @@ -158,21 +158,30 @@ fn window_follow_up( publication_pending: bool, projection_moved: bool, holding: bool, + history_made_progress: bool, ) -> WindowFollowUp { + // A deliberate yield is often `Retryable` (`ingest_pass_backpressured`), + // not `Pending`. That continuation is a backoff, and taking it before the + // admitted window is published hands the worker back to the out-of-scope + // cursor sweep with the newest day still on a building generation. + let owed = publication_pending + && projection_moved + && match history { + HistoryContinuation::Immediate => true, + HistoryContinuation::Settled => holding, + // A retryable pass that wrote nothing keeps the backoff. A pass + // that admitted the newest day and yielded is also retryable, and + // that window still has to be published first. + HistoryContinuation::Backoff => history_made_progress, + }; + if owed { + return WindowFollowUp::PublishBeforeNextHistory; + } match history { + HistoryContinuation::Immediate => WindowFollowUp::ContinueHistory, HistoryContinuation::Backoff => WindowFollowUp::BackoffHistory, - HistoryContinuation::Immediate | HistoryContinuation::Settled => { - let owed = publication_pending - && projection_moved - && (history == HistoryContinuation::Immediate || holding); - if owed { - WindowFollowUp::PublishBeforeNextHistory - } else if history == HistoryContinuation::Immediate || holding { - WindowFollowUp::ContinueHistory - } else { - WindowFollowUp::Settle - } - } + HistoryContinuation::Settled if holding => WindowFollowUp::ContinueHistory, + HistoryContinuation::Settled => WindowFollowUp::Settle, } } @@ -482,6 +491,7 @@ pub(super) async fn run_session_temporal_refresh_scheduler( projection_still_unpublished(&report), projection_published_work(&report), holding_publication, + history_outcome.is_some_and(SessionHistoricalIngestOutcome::made_progress), ) { // The admitted window is not searchable until this backlog // is published. Opening the next history window first is @@ -1353,10 +1363,22 @@ mod tests { projection_still_unpublished(&unpublished), projection_published_work(&unpublished), false, + true, ), WindowFollowUp::PublishBeforeNextHistory, "a progressing history window with a projection backlog must not open the next window" ); + assert_eq!( + window_follow_up( + HistoryContinuation::Backoff, + projection_still_unpublished(&unpublished), + projection_published_work(&unpublished), + false, + true, + ), + WindowFollowUp::PublishBeforeNextHistory, + "a progressing retryable yield must publish before the backoff opens the next window" + ); let still_moving = SessionTemporalRefreshPassReport { completed: 16, @@ -1369,6 +1391,7 @@ mod tests { projection_still_unpublished(&still_moving), projection_published_work(&still_moving), true, + false, ), WindowFollowUp::PublishBeforeNextHistory, "projection-only passes keep the hold while the admitted window is still unpublished" @@ -1385,6 +1408,7 @@ mod tests { projection_still_unpublished(&published), projection_published_work(&published), true, + false, ), WindowFollowUp::ContinueHistory, "a published window releases the next history window" @@ -1400,13 +1424,14 @@ mod tests { projection_still_unpublished(&stalled), projection_published_work(&stalled), true, + false, ), WindowFollowUp::ContinueHistory, "a projection pass that moves nothing must not spin ahead of history" ); assert_eq!( - window_follow_up(HistoryContinuation::Backoff, true, true, true), + window_follow_up(HistoryContinuation::Backoff, true, true, true, false), WindowFollowUp::BackoffHistory ); } From e9f7b00eadf34c5ae9249b91701c1c0f247a3db8 Mon Sep 17 00:00:00 2001 From: Zack Jackson <25274700+ScriptedAlchemy@users.noreply.github.com> Date: Mon, 28 Sep 2026 09:54:49 +0000 Subject: [PATCH 05/18] fix(sessions): hold projection when a yield reports no stats Codex catch-up persists observations without ingest counters, so the yield is a no-progress backoff. Publish the projection backlog before that backoff opens the older-day sweep. Co-authored-by: Zack Jackson --- .../worker.rs | 28 +++++++------------ 1 file changed, 10 insertions(+), 18 deletions(-) diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs index eb66406d37..bf49f90d09 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs @@ -158,21 +158,18 @@ fn window_follow_up( publication_pending: bool, projection_moved: bool, holding: bool, - history_made_progress: bool, ) -> WindowFollowUp { - // A deliberate yield is often `Retryable` (`ingest_pass_backpressured`), - // not `Pending`. That continuation is a backoff, and taking it before the - // admitted window is published hands the worker back to the out-of-scope - // cursor sweep with the newest day still on a building generation. + // A newest-day yield is often `Retryable` backpressure whose ingest stats + // stay at zero, so it is a backoff rather than an immediate continuation. + // Taking that backoff before the projection backlog is published hands the + // worker back to the out-of-scope cursor sweep with the newest day still + // on a building generation. The no-progress backoff applies only when this + // pass did not move that backlog. let owed = publication_pending && projection_moved && match history { - HistoryContinuation::Immediate => true, + HistoryContinuation::Immediate | HistoryContinuation::Backoff => true, HistoryContinuation::Settled => holding, - // A retryable pass that wrote nothing keeps the backoff. A pass - // that admitted the newest day and yielded is also retryable, and - // that window still has to be published first. - HistoryContinuation::Backoff => history_made_progress, }; if owed { return WindowFollowUp::PublishBeforeNextHistory; @@ -491,7 +488,6 @@ pub(super) async fn run_session_temporal_refresh_scheduler( projection_still_unpublished(&report), projection_published_work(&report), holding_publication, - history_outcome.is_some_and(SessionHistoricalIngestOutcome::made_progress), ) { // The admitted window is not searchable until this backlog // is published. Opening the next history window first is @@ -1363,7 +1359,6 @@ mod tests { projection_still_unpublished(&unpublished), projection_published_work(&unpublished), false, - true, ), WindowFollowUp::PublishBeforeNextHistory, "a progressing history window with a projection backlog must not open the next window" @@ -1374,7 +1369,6 @@ mod tests { projection_still_unpublished(&unpublished), projection_published_work(&unpublished), false, - true, ), WindowFollowUp::PublishBeforeNextHistory, "a progressing retryable yield must publish before the backoff opens the next window" @@ -1391,7 +1385,6 @@ mod tests { projection_still_unpublished(&still_moving), projection_published_work(&still_moving), true, - false, ), WindowFollowUp::PublishBeforeNextHistory, "projection-only passes keep the hold while the admitted window is still unpublished" @@ -1408,7 +1401,6 @@ mod tests { projection_still_unpublished(&published), projection_published_work(&published), true, - false, ), WindowFollowUp::ContinueHistory, "a published window releases the next history window" @@ -1424,15 +1416,15 @@ mod tests { projection_still_unpublished(&stalled), projection_published_work(&stalled), true, - false, ), WindowFollowUp::ContinueHistory, "a projection pass that moves nothing must not spin ahead of history" ); assert_eq!( - window_follow_up(HistoryContinuation::Backoff, true, true, true, false), - WindowFollowUp::BackoffHistory + window_follow_up(HistoryContinuation::Backoff, true, false, false), + WindowFollowUp::BackoffHistory, + "a backoff whose projection pass moved nothing keeps the retry delay" ); } From fcdf3d1cc658e5938f2ab9d47bb7b1dc339a6b15 Mon Sep 17 00:00:00 2001 From: Zack Jackson <25274700+ScriptedAlchemy@users.noreply.github.com> Date: Mon, 28 Sep 2026 10:00:19 +0000 Subject: [PATCH 06/18] fix(sessions): keep projection ahead of a queued history wake A history wake that arrives during the newest-day window was starting the older-day sweep on the next pass. While that window is unpublished, ignore the wake and keep projecting. Co-authored-by: Zack Jackson --- .../src/session_temporal_refresh_scheduler/worker.rs | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs index bf49f90d09..f4ff102c95 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs @@ -201,7 +201,15 @@ pub(super) async fn run_session_temporal_refresh_scheduler( } loop { let mut projection_requested = state.take_dirty(); - let history_requested = state.take_historical_dirty(); + let mut history_requested = state.take_historical_dirty(); + // A wake that arrives while the admitted window is still + // unpublished must not start the next history window. That window + // is the out-of-scope cursor sweep, and it holds this worker until + // it returns, so the newest day never leaves its building generation. + if state.history_held_for_projection() { + history_requested = false; + projection_requested = true; + } if !projection_requested && !history_requested { break; } From dbc415cf28994124d2122990af375cff2aa183ca Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 28 Sep 2026 15:48:36 +0000 Subject: [PATCH 07/18] fix(sessions): report scheduled imports as deferred, yield per day An import read the worker's serving state before its wake, so a converged profile reported `completed` with no remaining work while rollouts written since the last pass were still unadmitted. Scheduling now marks history pending first and every hand-off is partial with one deferred unit per store; blocked or missing workers stay failures. A project Codex pass that persisted in-scope frames now yields at the next day directory instead of at the next out-of-scope rollout. On an interleaved corpus the old rule ended nearly every pass after a few rollouts. Removes the waiter's leftover progress counters, history sequence, and membership pre-read. --- .../src/session_sync.rs | 148 ++---------------- .../session_sync/import_admission_tests.rs | 35 +++-- .../src/session_sync/work.rs | 13 +- .../history.rs | 39 +---- .../wake.rs | 60 ++----- .../worker.rs | 51 ++---- .../src/runtime/hosts/codex.rs | 13 -- .../src/runtime/ingest/project_provider.rs | 21 ++- 8 files changed, 80 insertions(+), 300 deletions(-) diff --git a/crates/tracedecay-session-runtime/src/session_sync.rs b/crates/tracedecay-session-runtime/src/session_sync.rs index ded403b8de..008b78fbb4 100644 --- a/crates/tracedecay-session-runtime/src/session_sync.rs +++ b/crates/tracedecay-session-runtime/src/session_sync.rs @@ -15,17 +15,13 @@ use tracedecay_contracts::session_sync::{ use tracedecay_contracts::{ CancellationSignal, Deadline, IdempotencyKey, OperationTermination, now_micros, }; -use tracedecay_domain::{BrainId, ProjectId, SessionId, UserProfileId, UtcMicros}; +use tracedecay_domain::{BrainId, ProjectId, UserProfileId, UtcMicros}; use crate::session_temporal_refresh_scheduler::wake::HistoricalAdmissionView; use tracedecay_global_db::GlobalDbGitCorrelationStore; use tracedecay_global_db::RegisteredGlobalDbLeaseV1; use tracedecay_runtime_core::background_cpu::ProcessBackgroundCpuV1; -use tracedecay_session_temporal_store::SessionTemporalAccess; use tracedecay_sessions::admission::{SESSION_INGEST_DISABLED_REASON_V1, session_ingest_disabled}; -use tracedecay_sessions::serving::{ - SessionProjectionServingState, SessionProjectionServingStatusPort, -}; const MAX_SESSION_SYNC_OPERATIONS: usize = 128; const COALESCED_JOURNAL_RECHECK_INTERVAL: Duration = Duration::from_millis(250); @@ -629,126 +625,34 @@ impl DaemonSessionSyncService { } } - async fn await_import_history( + /// Hands historical catch-up to both refresh workers. The import records + /// that hand-off, one deferred unit per store scope, and never waits for + /// the pass: only the workers' own serving state can later say it ran. + fn schedule_import_history( &self, context: &SessionSyncProjectContext, request: &SessionSyncRequestV1, - project_sessions: &RegisteredGlobalDbLeaseV1, - ) -> Result> { + ) -> Result { if let Some(interruption) = self.observed_interruption(request.cancellation(), request.deadline()) { - return Err(Some(interruption)); + return Err(interruption); } - let project = context - .project_refresh - .observe_and_schedule_historical_admission(); - let user = context - .user_refresh - .observe_and_schedule_historical_admission(); let mut failure_codes = Vec::new(); - push_historical_admission_failure(&project, &mut failure_codes); - push_historical_admission_failure(&user, &mut failure_codes); - let mut coverage = vec![ - historical_admission_coverage("project", &project), - historical_admission_coverage("profile", &user), - ]; - if failure_codes.is_empty() && coverage.iter().all(|entry| entry.coverage.is_complete()) { - match self - .import_projection_is_settled(context, project_sessions, request) - .await - { - Ok(true) => {} - Ok(false) => { - for entry in &mut coverage { - entry.coverage = SessionSyncCoverageV1::Partial { deferred_units: 1 }; - } - } - Err(Some(work::SessionSyncInterruption::TimedOut)) => { - for entry in &mut coverage { - entry.coverage = SessionSyncCoverageV1::Partial { deferred_units: 1 }; - } - } - Err(interruption) => return Err(interruption), - } - } + push_historical_admission_failure( + &context.project_refresh.schedule_historical_admission(), + &mut failure_codes, + ); + push_historical_admission_failure( + &context.user_refresh.schedule_historical_admission(), + &mut failure_codes, + ); Ok(ImportHistoryObservation { - coverage, + coverage: transcript_import_requested_coverage(), failure_codes, }) } - async fn import_projection_is_settled( - &self, - context: &SessionSyncProjectContext, - project_sessions: &RegisteredGlobalDbLeaseV1, - request: &SessionSyncRequestV1, - ) -> Result> { - let project = context.project_refresh.status(); - let user = context.user_refresh.status(); - if project.backlog > 0 - || user.backlog > 0 - || project.unavailable_reason.is_some() - || user.unavailable_reason.is_some() - || !matches!( - context.project_refresh.serving_status().state, - SessionProjectionServingState::Current - ) - || !matches!( - context.user_refresh.serving_status().state, - SessionProjectionServingState::Current - ) - { - return Ok(false); - } - if !self - .projection_store_is_current(project_sessions, request) - .await? - { - return Ok(false); - } - self.projection_store_is_current(&context.user_sessions, request) - .await - } - - async fn projection_store_is_current( - &self, - database: &RegisteredGlobalDbLeaseV1, - request: &SessionSyncRequestV1, - ) -> Result> { - let page_limit = - crate::session_temporal_refresh_scheduler::projector::SessionTemporalRefreshPolicy::default() - .max_begin_requests_per_pass; - let active_scan_slots = page_limit / 2; - - let temporal = SessionTemporalAccess::new(&**database); - let mut active_after: Option = None; - loop { - let page = { - let discovery = temporal.pending_session_temporal_refresh_page_result( - page_limit, - active_scan_slots, - active_after.as_ref(), - ); - tokio::pin!(discovery); - tokio::select! { - page = &mut discovery => page.map_err(|_| None)?, - interruption = self.wait_for_interruption(request) => { - return Err(Some(interruption)); - } - } - }; - let (pending, next_active, has_more) = page.into_parts(); - if !pending.is_empty() { - return Ok(false); - } - if !has_more { - return Ok(true); - } - active_after = next_active; - } - } - #[hotpath::skip] async fn transition_running( &self, @@ -1324,30 +1228,12 @@ fn import_reports_deferred_progress( }) } -fn historical_admission_coverage( - store_scope: &str, - view: &HistoricalAdmissionView, -) -> SessionSyncSourceCoverageV1 { - let coverage = match view { - HistoricalAdmissionView::Current => SessionSyncCoverageV1::Complete, - HistoricalAdmissionView::InProgress - | HistoricalAdmissionView::Blocked { .. } - | HistoricalAdmissionView::Unavailable => { - SessionSyncCoverageV1::Partial { deferred_units: 1 } - } - }; - SessionSyncSourceCoverageV1 { - store_scope: store_scope.to_owned(), - coverage, - } -} - fn push_historical_admission_failure( view: &HistoricalAdmissionView, failure_codes: &mut Vec, ) { match view { - HistoricalAdmissionView::Current | HistoricalAdmissionView::InProgress => {} + HistoricalAdmissionView::Scheduled => {} HistoricalAdmissionView::Blocked { reason_code } => { failure_codes.push(reason_code.clone()); } diff --git a/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs b/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs index ae827a5bc5..61a2c7544b 100644 --- a/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs +++ b/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs @@ -13,6 +13,9 @@ use tracedecay_domain::{ProjectId, UtcMicros}; use tracedecay_global_db::tests::harness::{HostAdmissionScope, HostAdmissionTestRuntimeV1}; use tracedecay_runtime_core::background_cpu::ProcessBackgroundCpuV1; use tracedecay_runtime_core::config::ProfileRoot; +use tracedecay_sessions::serving::{ + SessionProjectionServingState, SessionProjectionServingStatusPort, SessionProjectionStaleReason, +}; use crate::session_sync::{DaemonSessionSyncConfig, DaemonSessionSyncService}; use crate::session_temporal_refresh_scheduler::SessionTemporalRefreshWake; @@ -53,7 +56,10 @@ fn bound_wake(state: &Arc) -> SessionTemporalRe async fn import_receipt( kind: CatchUp, label: &str, -) -> tracedecay_contracts::session_sync::SessionSyncCompletionReceiptV1 { +) -> ( + tracedecay_contracts::session_sync::SessionSyncCompletionReceiptV1, + Arc, +) { let root = tempfile::tempdir().expect("import fixture directory"); let project_root = root.path().join("project"); std::fs::create_dir_all(&project_root).expect("project directory"); @@ -113,7 +119,7 @@ async fn import_receipt( started.elapsed() < Duration::from_secs(1), "import consumed its observation bound instead of returning the settled catch-up" ); - return receipt; + return (receipt, project_state); } SessionSyncOutcomeV1::Accepted(_) | SessionSyncOutcomeV1::Joined(_) => { assert!( @@ -139,26 +145,35 @@ fn remaining_work( #[tokio::test] async fn import_reports_deferred_progress_while_historical_catch_up_is_still_pending() { - let receipt = import_receipt(CatchUp::Pending, "import-pending").await; + let (receipt, _) = import_receipt(CatchUp::Pending, "import-pending").await; - assert_ne!(receipt.termination, OperationTermination::TimedOut); assert_eq!(receipt.termination, OperationTermination::Partial); assert!(receipt.failure_codes.is_empty()); - assert!(remaining_work(&receipt) > 0); + assert_eq!(remaining_work(&receipt), 2); } #[tokio::test] -async fn import_completes_when_historical_catch_up_is_already_current() { - let receipt = import_receipt(CatchUp::Current, "import-current").await; +async fn import_after_current_catch_up_defers_until_the_scheduled_pass_runs() { + let (receipt, state) = import_receipt(CatchUp::Current, "import-current").await; - assert_eq!(receipt.termination, OperationTermination::Completed); + // History was current before the request, but the pass it scheduled has + // not run: sources written since the last pass are not admitted yet. + assert_eq!(receipt.termination, OperationTermination::Partial); assert!(receipt.failure_codes.is_empty()); - assert_eq!(remaining_work(&receipt), 0); + assert_eq!(remaining_work(&receipt), 2); + assert_eq!( + bound_wake(&state).serving_status().state, + SessionProjectionServingState::Stale { + reason: SessionProjectionStaleReason::HistoricalConvergence, + } + ); + assert!(state.take_historical_dirty()); } #[tokio::test] async fn import_keeps_a_blocked_catch_up_as_a_failure() { - let receipt = import_receipt(CatchUp::Blocked, "import-blocked").await; + let (receipt, state) = import_receipt(CatchUp::Blocked, "import-blocked").await; + assert!(!state.take_historical_dirty()); assert_eq!(receipt.termination, OperationTermination::Failed); assert!( diff --git a/crates/tracedecay-session-runtime/src/session_sync/work.rs b/crates/tracedecay-session-runtime/src/session_sync/work.rs index 2df4675c50..0f76ebe8bb 100644 --- a/crates/tracedecay-session-runtime/src/session_sync/work.rs +++ b/crates/tracedecay-session-runtime/src/session_sync/work.rs @@ -491,18 +491,9 @@ impl SessionSyncProjectContext { request: &SessionSyncRequestV1, project_sessions: RegisteredGlobalDbLeaseV1, ) -> SessionSyncWorkResult { - let observation = match service - .await_import_history(self, request, &project_sessions) - .await - { + let observation = match service.schedule_import_history(self, request) { Ok(observation) => observation, - Err(Some(interruption)) => { - return SessionSyncWorkResult::Interrupted(interruption); - } - Err(None) => super::ImportHistoryObservation { - coverage: super::transcript_import_requested_coverage(), - failure_codes: vec!["session_history_not_current".to_owned()], - }, + Err(interruption) => return SessionSyncWorkResult::Interrupted(interruption), }; let pass = async { let stats = diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/history.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/history.rs index 68bf9db4bb..28542d9124 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/history.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/history.rs @@ -3,7 +3,6 @@ use std::path::PathBuf; use std::pin::Pin; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::{Mutex, PoisonError}; use tracedecay_runtime_core::config::ProfileRoot; use tracedecay_contracts::ProfileIdentityReadPort; @@ -58,20 +57,10 @@ impl SessionHistoricalIngestOutcome { pub trait SessionHistoricalIngestor: Send + Sync { fn run_pass(&self) -> SessionHistoricalIngestPass<'_>; fn cancel(&self); - - fn take_progress(&self) -> SessionHistoricalIngestProgress { - SessionHistoricalIngestProgress::default() - } } pub type SharedSessionHistoricalIngestor = Arc; -#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] -pub struct SessionHistoricalIngestProgress { - pub stats: tracedecay_sessions::TranscriptIngestStats, - pub committed: bool, -} - pub struct ProjectSessionHistoricalIngestor { database: RegisteredGlobalDbLeaseV1, profile_identity: Arc, @@ -83,7 +72,6 @@ pub struct ProjectSessionHistoricalIngestor { background_cpu: Arc, codex_consumer: String, codex_registered: AtomicBool, - progress: Mutex, } impl ProjectSessionHistoricalIngestor { @@ -118,7 +106,6 @@ impl ProjectSessionHistoricalIngestor { background_cpu, codex_consumer, codex_registered: AtomicBool::new(true), - progress: Mutex::new(SessionHistoricalIngestProgress::default()), } } @@ -154,13 +141,7 @@ impl SessionHistoricalIngestor for ProjectSessionHistoricalIngestor { pass, ) .await; - let progress = SessionHistoricalIngestProgress { - stats: outcome.stats, - committed: outcome.scheduling_state_written || outcome.made_progress(), - }; - let classified = classify_transcript_ingest_outcome(outcome, &self.cancellation); - *self.progress.lock().unwrap_or_else(PoisonError::into_inner) = progress; - classified + classify_transcript_ingest_outcome(outcome, &self.cancellation) }) } @@ -168,10 +149,6 @@ impl SessionHistoricalIngestor for ProjectSessionHistoricalIngestor { self.cancellation.cancel(); self.deregister_codex_once(); } - - fn take_progress(&self) -> SessionHistoricalIngestProgress { - std::mem::take(&mut *self.progress.lock().unwrap_or_else(PoisonError::into_inner)) - } } impl Drop for ProjectSessionHistoricalIngestor { @@ -191,7 +168,6 @@ pub struct ProfileSessionHistoricalIngestor { session_review: SessionReviewPort, codex_consumer: String, codex_registered: AtomicBool, - progress: Mutex, } impl ProfileSessionHistoricalIngestor { @@ -226,7 +202,6 @@ impl ProfileSessionHistoricalIngestor { session_review, codex_consumer, codex_registered: AtomicBool::new(true), - progress: Mutex::new(SessionHistoricalIngestProgress::default()), } } @@ -281,13 +256,7 @@ impl SessionHistoricalIngestor for ProfileSessionHistoricalIngestor { pass, ) .await; - let progress = SessionHistoricalIngestProgress { - stats: outcome.stats, - committed: outcome.scheduling_state_written || outcome.made_progress(), - }; - let classified = classify_transcript_ingest_outcome(outcome, &self.cancellation); - *self.progress.lock().unwrap_or_else(PoisonError::into_inner) = progress; - classified + classify_transcript_ingest_outcome(outcome, &self.cancellation) }) } @@ -295,10 +264,6 @@ impl SessionHistoricalIngestor for ProfileSessionHistoricalIngestor { self.cancellation.cancel(); self.deregister_codex_once(); } - - fn take_progress(&self) -> SessionHistoricalIngestProgress { - std::mem::take(&mut *self.progress.lock().unwrap_or_else(PoisonError::into_inner)) - } } impl Drop for ProfileSessionHistoricalIngestor { diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs index 36afaa9bc3..f5e7703399 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/wake.rs @@ -130,11 +130,6 @@ macro_rules! define_wake_state { /// is still unpublished. The next history window waits. pub(super) history_after_projection: AtomicBool, pub(super) pass_count: std::sync::atomic::AtomicUsize, - pub(super) history_requested_sequence: std::sync::atomic::AtomicUsize, - pub(super) history_completed_sequence: std::sync::atomic::AtomicUsize, - pub(super) history_sessions_upserted: std::sync::atomic::AtomicU64, - pub(super) history_messages_upserted: std::sync::atomic::AtomicU64, - pub(super) history_commit_count: std::sync::atomic::AtomicUsize, pub(super) wake: tokio::sync::Notify, pub(super) idle: tokio::sync::Notify, pub(super) cancelled: AtomicBool, @@ -165,11 +160,6 @@ impl Default for SessionTemporalRefreshWakeState { history_retry_pending: AtomicBool::new(false), history_after_projection: AtomicBool::new(false), pass_count: std::sync::atomic::AtomicUsize::new(0), - history_requested_sequence: std::sync::atomic::AtomicUsize::new(0), - history_completed_sequence: std::sync::atomic::AtomicUsize::new(0), - history_sessions_upserted: std::sync::atomic::AtomicU64::new(0), - history_messages_upserted: std::sync::atomic::AtomicU64::new(0), - history_commit_count: std::sync::atomic::AtomicUsize::new(0), wake: tokio::sync::Notify::new(), idle: tokio::sync::Notify::new(), cancelled: AtomicBool::new(false), @@ -360,29 +350,6 @@ impl SessionTemporalRefreshWakeState { self.wake.notify_one(); } - pub fn history_requested_sequence(&self) -> usize { - self.history_requested_sequence.load(Ordering::Acquire) - } - - pub fn record_history_progress( - &self, - progress: super::history::SessionHistoricalIngestProgress, - ) { - self.history_sessions_upserted - .fetch_add(progress.stats.sessions_upserted, Ordering::AcqRel); - self.history_messages_upserted - .fetch_add(progress.stats.messages_upserted, Ordering::AcqRel); - if progress.committed { - self.history_commit_count.fetch_add(1, Ordering::AcqRel); - } - } - - pub fn complete_history_sequence(&self, sequence: usize) { - self.history_completed_sequence - .fetch_max(sequence, Ordering::AcqRel); - self.idle.notify_waiters(); - } - pub fn has_pending_work(&self) -> bool { self.dirty.load(Ordering::Acquire) || self.historical_dirty.load(Ordering::Acquire) } @@ -862,22 +829,21 @@ pub struct SessionTemporalRefreshWake { route: Arc, } -/// What historical catch-up has already settled, read before this call's wake. +/// Whether a request for historical catch-up could be handed to its worker. #[derive(Clone, Debug, Eq, PartialEq)] pub(crate) enum HistoricalAdmissionView { - Current, - InProgress, + Scheduled, Blocked { reason_code: String }, Unavailable, } fn historical_admission_view(status: &SessionProjectionServingStatus) -> HistoricalAdmissionView { match &status.state { - SessionProjectionServingState::Current => HistoricalAdmissionView::Current, + SessionProjectionServingState::Current => HistoricalAdmissionView::Scheduled, SessionProjectionServingState::Stale { reason } => match reason { SessionProjectionStaleReason::HistoricalConvergence | SessionProjectionStaleReason::HistoricalRetry { .. } => { - HistoricalAdmissionView::InProgress + HistoricalAdmissionView::Scheduled } SessionProjectionStaleReason::HistoricalBlocked { reason_code } => { HistoricalAdmissionView::Blocked { @@ -887,7 +853,7 @@ fn historical_admission_view(status: &SessionProjectionServingStatus) -> Histori }, SessionProjectionServingState::Unavailable { reason } => match reason { SessionProjectionUnavailableReason::WorkerRecovering => { - HistoricalAdmissionView::InProgress + HistoricalAdmissionView::Scheduled } SessionProjectionUnavailableReason::WorkerMissing | SessionProjectionUnavailableReason::WorkerStalled @@ -980,24 +946,22 @@ impl SessionTemporalRefreshWake { } } - /// Schedules another historical pass and reports the admission already - /// settled by the worker. + /// Marks historical catch-up pending and wakes its worker, unless the + /// worker is blocked or gone. /// /// Callers that must return inside a bound, such as `sessions import`, /// use this instead of waiting for the pass just queued. Waiting for that /// pass is historical convergence: a large Codex home does not finish it - /// before the import deadline, so the operation was recorded as timed out - /// while catch-up was still admitting rollouts. + /// before the import deadline. The pass has not run when this returns, so + /// `Scheduled` is never evidence that sources are admitted. #[hotpath::skip] - pub(crate) fn observe_and_schedule_historical_admission(&self) -> HistoricalAdmissionView { + pub(crate) fn schedule_historical_admission(&self) -> HistoricalAdmissionView { let Some(state) = self.target() else { return HistoricalAdmissionView::Unavailable; }; let view = historical_admission_view(&state.serving_status()); - if matches!( - view, - HistoricalAdmissionView::Current | HistoricalAdmissionView::InProgress - ) { + if view == HistoricalAdmissionView::Scheduled { + state.mark_history_pending(); state.wake_history(); } view diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs index f4ff102c95..335dc47c1a 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs @@ -10,10 +10,7 @@ use tracedecay_store::{ SessionRefreshProgressV1, SessionRefreshStore, SessionStoreError, }; -use super::history::{ - SessionHistoricalIngestOutcome, SessionHistoricalIngestProgress, - SharedSessionHistoricalIngestor, -}; +use super::history::{SessionHistoricalIngestOutcome, SharedSessionHistoricalIngestor}; use super::projector::{ SessionTemporalRefreshEffect, SessionTemporalRefreshPolicy, SessionTemporalRefreshProjector, SessionTemporalRefreshProjectorError, SessionTemporalRefreshProjectorErrorClass, @@ -216,8 +213,7 @@ pub(super) async fn run_session_temporal_refresh_scheduler( state.begin_pass(); state.mark_worker_busy(); state.pass_count.fetch_add(1, Ordering::AcqRel); - let history_sequence = history_requested.then(|| state.history_requested_sequence()); - let history_result = if history_requested { + let history_outcome = if history_requested { Some( hotpath::future!( session_history_refresh(&history, &history_admission), @@ -228,10 +224,8 @@ pub(super) async fn run_session_temporal_refresh_scheduler( } else { None }; - let history_outcome = history_result.map(|result| result.0); projection_requested |= state.take_dirty(); - if let Some((outcome, progress)) = history_result { - state.record_history_progress(progress); + if let Some(outcome) = history_outcome { state.record_history_outcome(outcome); } if matches!( @@ -284,15 +278,6 @@ pub(super) async fn run_session_temporal_refresh_scheduler( if state.cancelled.load(Ordering::Acquire) { return; } - if let (Some(sequence), Some(outcome)) = (history_sequence, history_outcome) - && matches!( - outcome, - SessionHistoricalIngestOutcome::Complete - | SessionHistoricalIngestOutcome::Blocked { .. } - ) - { - state.complete_history_sequence(sequence); - } let holding_publication = state.history_held_for_projection(); // A projection-only iteration that is finishing the admitted window // must not spend the pass on summary convergence. History did not @@ -686,10 +671,7 @@ fn observe_retry(class: SessionTemporalRefreshRetryClass, attempt: u32) { async fn session_history_refresh( history: &Arc>>, admission: &tokio::sync::Semaphore, -) -> ( - SessionHistoricalIngestOutcome, - SessionHistoricalIngestProgress, -) { +) -> SessionHistoricalIngestOutcome { let history = history .read() .unwrap_or_else(PoisonError::into_inner) @@ -698,24 +680,17 @@ async fn session_history_refresh( Some(history) => { let Ok(_permit) = admission.try_acquire() else { hotpath::gauge!("session_temporal_refresh_history_admission_deferrals").inc(1.0); - return ( - SessionHistoricalIngestOutcome::Retryable { - reason_code: HISTORY_ADMISSION_SATURATED_REASON, - made_progress: false, - }, - SessionHistoricalIngestProgress::default(), - ); + return SessionHistoricalIngestOutcome::Retryable { + reason_code: HISTORY_ADMISSION_SATURATED_REASON, + made_progress: false, + }; }; - let outcome = history.run_pass().await; - (outcome, history.take_progress()) + history.run_pass().await } - None => ( - SessionHistoricalIngestOutcome::Blocked { - reason_code: "history_ingestor_missing", - made_progress: false, - }, - SessionHistoricalIngestProgress::default(), - ), + None => SessionHistoricalIngestOutcome::Blocked { + reason_code: "history_ingestor_missing", + made_progress: false, + }, } } diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex.rs index 3ff028e871..a4c45d0d54 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex.rs @@ -117,19 +117,6 @@ pub use observation::{ try_admit_codex_jsonl_observations_for_project_with_admission_and_cancellation, }; -/// Project membership from a rollout's leading `session_meta` cwd. -/// -/// `None` means the header could not be read. Callers may leave a rollout for -/// a later pass only on a definitive [`ProjectMembership::NoMatch`]; an -/// `Unknown` git timeout stays in the current pass. -pub(crate) fn codex_rollout_project_membership( - path: &Path, - project_root: &Path, -) -> Option { - let meta = session_meta(path)?; - Some(TranscriptScopeMatcher::project(project_root).membership(Some(&meta.cwd))) -} - const PROVIDER: &str = "codex"; /// `~/.codex/sessions/YYYY/MM/DD/rollout-*.jsonl` → date dirs add depth. const MAX_SCAN_DEPTH: u8 = 6; diff --git a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs index 168ebb89cb..e93ba8e4a8 100644 --- a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs +++ b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs @@ -8,7 +8,7 @@ use tracedecay_domain::{ObservationScopeV1, ProjectId}; use crate::admission::HostAdmission; use crate::observation::ObservationCancellation; -use crate::runtime::shared::{ProjectMembership, TranscriptIngestStats}; +use crate::runtime::shared::TranscriptIngestStats; use crate::runtime::source::{ HostCoverageReason, HostProviderCoverage, TranscriptDiscoveryBounds, persist_codex_history_frontier, persist_host_provider_coverage, read_codex_history_frontier, @@ -262,10 +262,12 @@ impl<'a> ProjectProviderRun<'a> { let mut remaining = self.max_new_bytes; let mut deferred = discovery.is_truncated(); let mut frontier_committable = true; - // Out-of-scope rollouts consume no byte budget, so a newest-first page - // would otherwise write a cursor for every older day before this pass - // returns and the admitted window can be projected into search. - let mut persisted_in_scope = false; + // A Codex day directory is one publication window. Out-of-scope + // rollouts consume no byte budget, so a newest-first page would + // otherwise write a cursor for every older day before the admitted day + // is projected into search. Once a day has persisted in-scope frames, + // the next day belongs to the next pass. + let mut persisted_day: Option<&Path> = None; let mut outcome = ProviderRunOutcome::bounded(TranscriptIngestStats::default(), 0, false); for path in &discovery.paths { if remaining == 0 { @@ -278,12 +280,7 @@ impl<'a> ProjectProviderRun<'a> { frontier_committable = false; break; } - if persisted_in_scope - && run_blocking_transcript_section(|| { - codex::codex_rollout_project_membership(path, self.project_root) - == Some(ProjectMembership::NoMatch) - }) - { + if persisted_day.is_some_and(|day| path.parent() != Some(day)) { deferred = true; frontier_committable = false; break; @@ -300,7 +297,7 @@ impl<'a> ProjectProviderRun<'a> { { Ok(progress) => { if progress.frames_persisted > 0 { - persisted_in_scope = true; + persisted_day = path.parent(); } deferred |= progress.source_deferred || progress.bytes_consumed > remaining; frontier_committable &= From 8f1fd32b7a2d5bd0f7c06f577498131ad1ab1cbc Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 28 Sep 2026 16:44:07 +0000 Subject: [PATCH 08/18] fix(sessions): yield only at a day that opens out of project A per-day yield ended every pass of a project whose history is all in scope, so the discovery frontier never committed and each pass re-read the days before it. The pass now yields only at a day directory that opens with an out-of-project rollout after in-scope frames persisted; mixed days and in-scope days stay in one pass. --- .../src/runtime/hosts/codex.rs | 13 ++++ .../src/runtime/hosts/codex/tests.rs | 60 +++++++++++++++++++ .../src/runtime/ingest/project_provider.rs | 20 ++++--- 3 files changed, 86 insertions(+), 7 deletions(-) diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex.rs index a4c45d0d54..3ff028e871 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex.rs @@ -117,6 +117,19 @@ pub use observation::{ try_admit_codex_jsonl_observations_for_project_with_admission_and_cancellation, }; +/// Project membership from a rollout's leading `session_meta` cwd. +/// +/// `None` means the header could not be read. Callers may leave a rollout for +/// a later pass only on a definitive [`ProjectMembership::NoMatch`]; an +/// `Unknown` git timeout stays in the current pass. +pub(crate) fn codex_rollout_project_membership( + path: &Path, + project_root: &Path, +) -> Option { + let meta = session_meta(path)?; + Some(TranscriptScopeMatcher::project(project_root).membership(Some(&meta.cwd))) +} + const PROVIDER: &str = "codex"; /// `~/.codex/sessions/YYYY/MM/DD/rollout-*.jsonl` → date dirs add depth. const MAX_SCAN_DEPTH: u8 = 6; diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs index a92bbbfd83..d80a8c99d6 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs @@ -1137,6 +1137,66 @@ mod goal_event_tests { ); } + /// An out-of-project rollout inside the day being admitted does not end + /// the pass. On a profile that interleaves projects, ending there admits a + /// few rollouts per pass and rediscovers the same page each time. + #[tokio::test] + async fn a_mixed_day_is_admitted_in_one_pass_before_older_days_are_opened() { + crate::runtime::jsonl_observation_admission::install_test_shared_jsonl_preparation_authority(); + let temp = tempfile::tempdir().unwrap(); + let home = temp.path().canonicalize().unwrap(); + let project = home.join("project"); + let other = home.join("other"); + std::fs::create_dir_all(&project).unwrap(); + std::fs::create_dir_all(&other).unwrap(); + for (session_id, cwd) in [ + ("mixed-a", &project), + ("mixed-m", &other), + ("mixed-z", &project), + ] { + write_scoped_rollout(&home, "2026", "08", "29", session_id, cwd, "mixed day"); + } + write_scoped_rollout(&home, "2026", "08", "28", "older-x", &other, "older day"); + + let project_id = ProjectId::new("project-mixed-day").unwrap(); + let scope = ObservationScopeV1::Project { + project_id: project_id.clone(), + }; + let admission = MemoryHostAdmission::default(); + let outcome = with_transcript_source_profile( + tracedecay_runtime_core::config::ProfileRoot::under_home(home.clone()), + ProjectProviderRun { + project_root: &project, + project_id: &project_id, + facade: &admission, + scope: &scope, + candidate: SessionProvider::Codex, + max_new_bytes: u64::MAX, + cancellation: &ObservationCancellation::default(), + codex_discovery: None, + } + .run_codex(), + ) + .await; + + assert!(outcome.failures.is_empty(), "{:?}", outcome.failures); + assert_eq!( + session_ids_of(&admission.observations()), + ["mixed-a".to_owned(), "mixed-z".to_owned()] + .into_iter() + .collect::>() + ); + let older = crate::runtime::hosts::codex::codex_observation_source_v2("older-x").unwrap(); + assert!( + admission + .get_source_cursor(&older, &scope) + .await + .unwrap() + .is_none(), + "the older out-of-project day belongs to the next pass" + ); + } + fn write_scoped_rollout( home: &std::path::Path, year: &str, diff --git a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs index e93ba8e4a8..b5c383df19 100644 --- a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs +++ b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs @@ -8,7 +8,7 @@ use tracedecay_domain::{ObservationScopeV1, ProjectId}; use crate::admission::HostAdmission; use crate::observation::ObservationCancellation; -use crate::runtime::shared::TranscriptIngestStats; +use crate::runtime::shared::{ProjectMembership, TranscriptIngestStats}; use crate::runtime::source::{ HostCoverageReason, HostProviderCoverage, TranscriptDiscoveryBounds, persist_codex_history_frontier, persist_host_provider_coverage, read_codex_history_frontier, @@ -262,11 +262,12 @@ impl<'a> ProjectProviderRun<'a> { let mut remaining = self.max_new_bytes; let mut deferred = discovery.is_truncated(); let mut frontier_committable = true; - // A Codex day directory is one publication window. Out-of-scope - // rollouts consume no byte budget, so a newest-first page would - // otherwise write a cursor for every older day before the admitted day - // is projected into search. Once a day has persisted in-scope frames, - // the next day belongs to the next pass. + // Out-of-scope rollouts consume no byte budget, so a newest-first page + // would otherwise write a cursor for every older out-of-scope day + // before the admitted window is projected into search. Once in-scope + // frames persisted, a day directory that opens out of scope belongs to + // the next pass. In-scope days and mixed days keep the pass going, so + // a project whose history is all in scope still commits its frontier. let mut persisted_day: Option<&Path> = None; let mut outcome = ProviderRunOutcome::bounded(TranscriptIngestStats::default(), 0, false); for path in &discovery.paths { @@ -280,7 +281,12 @@ impl<'a> ProjectProviderRun<'a> { frontier_committable = false; break; } - if persisted_day.is_some_and(|day| path.parent() != Some(day)) { + if persisted_day.is_some_and(|day| path.parent() != Some(day)) + && run_blocking_transcript_section(|| { + codex::codex_rollout_project_membership(path, self.project_root) + == Some(ProjectMembership::NoMatch) + }) + { deferred = true; frontier_committable = false; break; From 3bc348e3b62651f780e3cbc45e1b521f754811b1 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 28 Sep 2026 19:33:40 +0000 Subject: [PATCH 09/18] test(sessions): assert the scheduled import receipt, then search --- .../tests/core_cli_suite/tracedecay_test.rs | 4 +- .../mcp_handler_test/session_search_test.rs | 43 +++++++------------ 2 files changed, 18 insertions(+), 29 deletions(-) diff --git a/crates/tracedecay-cli/tests/core_cli_suite/tracedecay_test.rs b/crates/tracedecay-cli/tests/core_cli_suite/tracedecay_test.rs index 2fac7acc02..b31fadb3c0 100644 --- a/crates/tracedecay-cli/tests/core_cli_suite/tracedecay_test.rs +++ b/crates/tracedecay-cli/tests/core_cli_suite/tracedecay_test.rs @@ -130,8 +130,8 @@ fn project_path_flags_resolve_a_relative_path_from_inside_the_project() { let stdout = String::from_utf8_lossy(&import.stdout); assert!( import.status.success() - && stdout.starts_with("session import completed (") - && stdout.ends_with(")\n"), + && stdout.starts_with("session import scheduled (") + && stdout.ends_with("); historical catch-up has remaining work 2\n"), "sessions import --project-path . must import into the CLI's project\nstdout:\n{stdout}\nstderr:\n{}", String::from_utf8_lossy(&import.stderr) ); diff --git a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs index a95c57062e..ae0d3c8fa6 100644 --- a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs +++ b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs @@ -804,11 +804,11 @@ async fn production_hook_ingest_reads_only_the_pinned_transcript_home() { #[cfg(feature = "test-transport")] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn completed_session_import_immediately_searches_canonical_message() { +async fn scheduled_session_import_makes_the_final_codex_source_searchable() { let root = test_temp_dir(); let isolation = root.path().join("composition"); - // `sessions_import` is the composition's own pass, so it reads the - // isolated transcript layout rather than the process home. + // `sessions_import` schedules the composition's own workers, so they read + // the isolated transcript layout rather than the process home. let transcripts = composed_transcript_home(&isolation); let project = isolation.join("project"); std::fs::create_dir_all(&project).expect("production composition project"); @@ -820,8 +820,8 @@ async fn completed_session_import_immediately_searches_canonical_message() { .expect("git init"); assert!(init.success(), "git init must succeed"); // More than one bounded transcript pass admits. The searchable message is - // in the final source, so Complete proves the production continuation - // worker consumed every durable Codex frontier before returning. + // in the final source, so finding it proves the continuation worker + // consumed every durable Codex frontier after the import returned. write_production_codex_rollouts(&transcripts, &project, 33); let harness = ProductionProjectCompositionHarnessV1::open_for_session_retrieval( @@ -875,30 +875,19 @@ async fn completed_session_import_immediately_searches_canonical_message() { }) .await .expect("session import completion deadline"); - assert_eq!(completed["termination"], "completed", "{completed}"); - assert!( - completed["stats"]["sessions_imported"] - .as_u64() - .is_some_and(|count| count > 0) - && completed["stats"]["messages_imported"] - .as_u64() - .is_some_and(|count| count > 0), - "{completed}" - ); - assert!( - completed["failure_codes"] - .as_array() - .is_some_and(Vec::is_empty), - "{completed}" - ); - assert!( - completed["coverage"].as_array().is_some_and(|coverage| { - coverage - .iter() - .all(|entry| entry["coverage"]["outcome"] == "complete") - }), + // The import hands catch-up to the workers and admits nothing itself, so + // its receipt is one deferred unit per store, never a completed claim. + assert_eq!(completed["termination"], "partial", "{completed}"); + assert_eq!(completed["failure_codes"], json!([]), "{completed}"); + assert_eq!( + completed["coverage"], + json!([ + {"store_scope": "project", "coverage": {"outcome": "partial", "deferred_units": 1}}, + {"store_scope": "profile", "coverage": {"outcome": "partial", "deferred_units": 1}}, + ]), "{completed}" ); + assert_eq!(completed["stats"]["messages_imported"], 0, "{completed}"); production_codex_message_search(&harness, &project).await; harness.shutdown().await; From f32849691939431e365fb04a8cbb5ccc84962f31 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 28 Sep 2026 20:10:19 +0000 Subject: [PATCH 10/18] test(work): wait for imported transcripts before journey reads --- .../advanced_workflow_journey_test.rs | 35 +++++++++++++++---- 1 file changed, 28 insertions(+), 7 deletions(-) diff --git a/crates/tracedecay/tests/daemon_suite/advanced_workflow_journey_test.rs b/crates/tracedecay/tests/daemon_suite/advanced_workflow_journey_test.rs index be9269c7c9..c21ca4b763 100644 --- a/crates/tracedecay/tests/daemon_suite/advanced_workflow_journey_test.rs +++ b/crates/tracedecay/tests/daemon_suite/advanced_workflow_journey_test.rs @@ -377,13 +377,35 @@ pub(super) fn advance_provider_transcript_participant_generation( .expect("open provider transcript for participant refresh"); writeln!(transcript, "{record}").expect("append provider transcript participant refresh"); drop(transcript); + import_provider_transcript( + home, + project, + "refreshed through the public sessions import authority", + PROVIDER_TRANSCRIPT_REFRESH_MESSAGE_ID, + ); +} + +/// `sessions import` only schedules catch-up, so the journey reads nothing +/// until the imported message is searchable through the same public CLI. +fn import_provider_transcript(home: &Path, project: &Path, query: &str, message_id: &str) { run( common::tracedecay_command_with_home(home) .args(["sessions", "import", "--project-path"]) .arg(project) .current_dir(project), - "tracedecay sessions import participant refresh", + "tracedecay sessions import", ); + wait_until(&format!("{message_id} to become searchable"), || { + let output = common::tracedecay_command_with_home(home) + .args(["sessions", "search", query, "--project-path"]) + .arg(project) + .current_dir(project) + .output() + .expect("tracedecay sessions search"); + String::from_utf8_lossy(&output.stdout) + .contains(message_id) + .then_some(()) + }); } fn initialize_project(home: &Path, project: &Path) -> (String, CommitId) { @@ -1279,12 +1301,11 @@ fn mounted_fan_out_recovers_then_synthesizes_and_hands_off() { .filter(|attempt| attempt.state() == WorkAttemptStateV1::Succeeded) }); write_provider_transcript(&home, &project, completed_synthesis.identity()); - run( - common::tracedecay_command_with_home(&home) - .args(["sessions", "import", "--project-path"]) - .arg(&project) - .current_dir(&project), - "tracedecay sessions import", + import_provider_transcript( + &home, + &project, + "completed through the typed SDK provider session", + PROVIDER_TRANSCRIPT_ASSISTANT_MESSAGE_ID, ); let graph = client .execute::(&WorkGraphReadRequestV1::current( From 53dacc318febd9c3153c199e7b21dc36882832d2 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 28 Sep 2026 20:45:57 +0000 Subject: [PATCH 11/18] style(cli): format session sync imports --- crates/tracedecay-cli/src/sessions_cmd/session_sync.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs b/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs index 9dd50f2c86..28545d5b60 100644 --- a/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs +++ b/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs @@ -1,8 +1,6 @@ use std::path::Path; use tracedecay_contracts::retrieval::{AdminCliSessionSyncV1, AdminCliSurfaceRequestV1}; -use tracedecay_contracts::session_sync::{ - SessionSyncCoverageV1, SessionSyncSourceCoverageV1, -}; +use tracedecay_contracts::session_sync::{SessionSyncCoverageV1, SessionSyncSourceCoverageV1}; use tracedecay_contracts::{IdempotencyKey, OperationTermination, RequestId}; use tracedecay_runtime_core::config::ProfileRoot; From d09b67592c118cdf860938ccf744113e406b93b0 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 28 Sep 2026 19:03:49 -0700 Subject: [PATCH 12/18] fix(ci): install Windows target on pinned toolchain --- .github/workflows/ci.yml | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 4396a9551f..546a94e198 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -956,7 +956,9 @@ jobs: - uses: dtolnay/rust-toolchain@stable with: components: clippy - targets: x86_64-pc-windows-gnu + + - name: Install Windows target for the pinned toolchain + run: rustup target add x86_64-pc-windows-gnu - uses: ./.github/actions/setup-linux-mold From 462471a552a7c0fbd9619b73e22acdb35dd9696c Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 06:33:48 +0000 Subject: [PATCH 13/18] fix(sessions): keep a live Codex tail from ending every catch-up pass The history pass yields after in-scope frames persist so that a newest-first page does not cursor every older out-of-scope day before the admitted window is searchable. Any persisted frames set persisted_day, including appends to a rollout that was already cursored. With an active Codex session writing to today's rollout between passes, every pass ended at the first older out-of-scope day with the frontier uncommittable, so the hub redelivered the same page and older in-scope days were never reached. Only a rollout opened for the first time this pass now arms the yield. A resumed tail already has its earlier window projected, so its appends no longer end the pass. --- .../src/runtime/hosts/codex/tests.rs | 99 +++++++++++++++++++ .../src/runtime/ingest/project_provider.rs | 12 ++- 2 files changed, 107 insertions(+), 4 deletions(-) diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs index d80a8c99d6..ce96ee624f 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs @@ -1197,6 +1197,105 @@ mod goal_event_tests { ); } + /// A session that appends to today's rollout between passes must not keep + /// the pass yielding at yesterday's out-of-project day. The yield exists + /// for a rollout opened this pass; a resumed live tail already has its + /// earlier window searchable, so the pass walks on, admits the older + /// in-project day, and commits its frontier. + #[tokio::test] + async fn a_live_tail_does_not_starve_older_in_project_days() { + crate::runtime::jsonl_observation_admission::install_test_shared_jsonl_preparation_authority(); + let temp = tempfile::tempdir().unwrap(); + let home = temp.path().canonicalize().unwrap(); + let project = home.join("project"); + let other = home.join("other"); + std::fs::create_dir_all(&project).unwrap(); + std::fs::create_dir_all(&other).unwrap(); + write_scoped_rollout(&home, "2026", "08", "29", "live", &project, "live day"); + write_scoped_rollout(&home, "2026", "08", "28", "other-a", &other, "other day"); + write_scoped_rollout(&home, "2026", "08", "27", "oldest", &project, "oldest day"); + let live_rollout = home.join(".codex/sessions/2026/08/29/rollout-live.jsonl"); + + let project_id = ProjectId::new("project-live-tail").unwrap(); + let scope = ObservationScopeV1::Project { + project_id: project_id.clone(), + }; + let admission = MemoryHostAdmission::default(); + let cancellation = ObservationCancellation::default(); + let run_pass = || { + with_transcript_source_profile( + tracedecay_runtime_core::config::ProfileRoot::under_home(home.clone()), + ProjectProviderRun { + project_root: &project, + project_id: &project_id, + facade: &admission, + scope: &scope, + candidate: SessionProvider::Codex, + max_new_bytes: u64::MAX, + cancellation: &cancellation, + codex_discovery: None, + } + .run_codex(), + ) + }; + + let first = run_pass().await; + assert!(first.failures.is_empty(), "{:?}", first.failures); + assert_eq!( + session_ids_of(&admission.observations()), + BTreeSet::from(["live".to_owned()]), + "the newest day yields before the older out-of-project day is opened" + ); + + let appends = 3; + for ordinal in 0..appends { + let mut file = std::fs::OpenOptions::new() + .append(true) + .open(&live_rollout) + .unwrap(); + std::io::Write::write_all( + &mut file, + format!( + "{}\n", + json!({ + "timestamp": format!("2026-08-29T12:00:{:02}.000Z", ordinal + 2), + "type": "event_msg", + "payload": {"type": "user_message", "message": format!("live append {ordinal}")} + }) + ) + .as_bytes(), + ) + .unwrap(); + let outcome = run_pass().await; + assert!(outcome.failures.is_empty(), "{:?}", outcome.failures); + } + + assert_eq!( + session_ids_of(&admission.observations()), + BTreeSet::from(["live".to_owned(), "oldest".to_owned()]), + "a live tail appending every pass starved the older in-project day" + ); + assert_eq!( + admission + .observations() + .iter() + .filter(|stored| { + let envelope: CanonicalObservationEnvelopeV1 = + serde_json::from_value(stored.observation().payload().clone()).unwrap(); + envelope.relations().session_id().as_str() == "live" + }) + .count(), + 2 + appends, + "every live append is admitted exactly once" + ); + assert_eq!( + read_host_provider_coverage(&admission, &scope, "codex") + .await + .unwrap(), + Some(HostProviderCoverage::Complete) + ); + } + fn write_scoped_rollout( home: &std::path::Path, year: &str, diff --git a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs index b5c383df19..6a2cb076f1 100644 --- a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs +++ b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs @@ -265,9 +265,13 @@ impl<'a> ProjectProviderRun<'a> { // Out-of-scope rollouts consume no byte budget, so a newest-first page // would otherwise write a cursor for every older out-of-scope day // before the admitted window is projected into search. Once in-scope - // frames persisted, a day directory that opens out of scope belongs to - // the next pass. In-scope days and mixed days keep the pass going, so - // a project whose history is all in scope still commits its frontier. + // frames persisted from a rollout opened this pass, a day directory + // that opens out of scope belongs to the next pass. In-scope days and + // mixed days keep the pass going, so a project whose history is all in + // scope still commits its frontier. A rollout resumed from its cursor + // is a live tail whose earlier window is already searchable; letting + // its appends end the pass would replay the same page while the + // session stays active and never reach an older in-scope day. let mut persisted_day: Option<&Path> = None; let mut outcome = ProviderRunOutcome::bounded(TranscriptIngestStats::default(), 0, false); for path in &discovery.paths { @@ -302,7 +306,7 @@ impl<'a> ProjectProviderRun<'a> { .await { Ok(progress) => { - if progress.frames_persisted > 0 { + if progress.frames_persisted > 0 && !progress.resumed { persisted_day = path.parent(); } deferred |= progress.source_deferred || progress.bytes_consumed > remaining; From 4179085161c6767b86724134e13ab67eb3a3dec6 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 06:33:51 +0000 Subject: [PATCH 14/18] test(session-runtime): loosen import hand-off bound for loaded runners A 1s wall-clock bound on the import settling is a real-time assertion that a 4-vCPU CI runner under load can miss. 10s still proves the request returned the settled catch-up instead of consuming its 60s deadline. --- .../src/session_sync/import_admission_tests.rs | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs b/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs index 61a2c7544b..f81ca08489 100644 --- a/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs +++ b/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs @@ -112,18 +112,21 @@ async fn import_receipt( }; let started = tokio::time::Instant::now(); let control = SessionSyncControlV1::new(scope, admission.idempotency_key); + // Well under the 60s request deadline the old waiter would have consumed, + // with room for a loaded runner: this bounds hand-off latency, not CPU. + let hand_off_bound = Duration::from_secs(10); loop { match SessionSyncServicePort::status(&service, control.clone()).await { SessionSyncOutcomeV1::Complete(receipt) => { assert!( - started.elapsed() < Duration::from_secs(1), + started.elapsed() < hand_off_bound, "import consumed its observation bound instead of returning the settled catch-up" ); return (receipt, project_state); } SessionSyncOutcomeV1::Accepted(_) | SessionSyncOutcomeV1::Joined(_) => { assert!( - started.elapsed() < Duration::from_secs(1), + started.elapsed() < hand_off_bound, "import stayed pending while catch-up state was already known" ); tokio::time::sleep(Duration::from_millis(10)).await; From a846b169f9710163577bb3c802b929cfd621d5cf Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 07:24:33 +0000 Subject: [PATCH 15/18] fix(sessions): avoid complete zero during import catch-up A deferred import can expose a locally fresh but empty projection while its historical worker still has sources to ingest. Treat that zero as stale until catch-up settles, and keep an end-to-end assertion on the first empty search response. --- .../src/session_retrieval/admitted.rs | 34 ++++++++++++++++--- .../mcp_handler_test/session_search_test.rs | 18 +++++++--- 2 files changed, 44 insertions(+), 8 deletions(-) diff --git a/crates/tracedecay-session-runtime/src/session_retrieval/admitted.rs b/crates/tracedecay-session-runtime/src/session_retrieval/admitted.rs index 52a1b05128..91cd0e34b7 100644 --- a/crates/tracedecay-session-runtime/src/session_retrieval/admitted.rs +++ b/crates/tracedecay-session-runtime/src/session_retrieval/admitted.rs @@ -18,12 +18,13 @@ use tracedecay_session_memory::context::{ CapabilityDigest, ConfigurationDigest, PolicyDigest, RequestBudgets, ResolvedSessionIdentity, }; use tracedecay_session_memory::session::{ - SessionRequestBinding, SessionRetrievalConfiguration, SessionTemporalQuery, - TaskSessionRetrievalOutcomeV1, + SessionDataFreshness, SessionRequestBinding, SessionRetrievalConfiguration, + SessionTemporalQuery, TaskSessionRetrievalOutcomeV1, }; use tracedecay_session_temporal_store::execution::TaskSessionRankSelectorV1; use tracedecay_sessions::serving::{ - RefreshWorkerMissing, SessionProjectionServingStatus, SessionProjectionServingStatusPort, + RefreshWorkerMissing, SessionProjectionServingState, SessionProjectionServingStatus, + SessionProjectionServingStatusPort, }; use tracedecay_store::StoreShardScopeV1; @@ -265,6 +266,10 @@ impl SessionApplicationRetrievalPortV1 for DaemonSessionRetrievalService { Ok(binding) => binding, Err(outcome) => return *outcome, }; + let history_pending = matches!( + self.refresh_status.serving_status().state, + SessionProjectionServingState::Stale { .. } + ); let outcome = self .execute_temporal_query_with_context( context, @@ -273,7 +278,28 @@ impl SessionApplicationRetrievalPortV1 for DaemonSessionRetrievalService { "grant.application.session-retrieval", ) .await; - self.public_outcome(outcome).await + match self.public_outcome(outcome).await { + SessionRetrievalServiceOutcome::CompleteZero { + temporal, + freshness, + } if history_pending + || matches!( + self.refresh_status.serving_status().state, + SessionProjectionServingState::Stale { .. } + ) => + { + SessionRetrievalServiceOutcome::Stale { + temporal, + freshness: match freshness { + SessionDataFreshness::Fresh => { + SessionDataFreshness::Stored { generation_lag: 0 } + } + other => other, + }, + } + } + outcome => outcome, + } }) } diff --git a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs index 2b34ec731c..307c88591e 100644 --- a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs +++ b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs @@ -95,13 +95,12 @@ async fn production_codex_message_search( harness: &ProductionProjectCompositionHarnessV1, project: &Path, ) -> Value { - // A `partial` generation is the store saying "still converging", the same - // not-ready contract as `stale`: re-read it. Every other outcome answers - // now, so an empty `complete_zero` still fails the assertions below. + // An empty partial or stale answer is not yet converged; a complete zero + // must still fail when the final Codex source has not appeared. let payload = tokio::time::timeout(std::time::Duration::from_secs(30), async { loop { let payload = production_codex_message_search_once(harness, project).await; - if payload["outcome"] != "partial" + if !matches!(payload["outcome"].as_str(), Some("partial" | "stale")) || payload["results"] .as_array() .is_some_and(|results| !results.is_empty()) @@ -889,6 +888,17 @@ async fn scheduled_session_import_makes_the_final_codex_source_searchable() { ); assert_eq!(completed["stats"]["messages_imported"], 0, "{completed}"); + let initial = production_codex_message_search_once(&harness, &project).await; + if initial["results"].as_array().is_some_and(Vec::is_empty) { + assert!( + matches!(initial["outcome"].as_str(), Some("partial" | "stale")), + "empty search claimed convergence before the final source arrived: {initial}" + ); + assert_ne!( + initial["temporal"]["freshness"]["state"], "fresh", + "{initial}" + ); + } production_codex_message_search(&harness, &project).await; harness.shutdown().await; } From 39d99e18d4c21188c4bec01a9239c13e76d4add8 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 08:52:36 +0000 Subject: [PATCH 16/18] test(sessions): pin catch-up empty answers and hand-off without clocks --- .../src/session_retrieval/tests.rs | 73 +++++++++++++++++++ .../session_sync/import_admission_tests.rs | 22 ++---- 2 files changed, 79 insertions(+), 16 deletions(-) diff --git a/crates/tracedecay-session-runtime/src/session_retrieval/tests.rs b/crates/tracedecay-session-runtime/src/session_retrieval/tests.rs index 5eafe338a7..d69d0ab7a7 100644 --- a/crates/tracedecay-session-runtime/src/session_retrieval/tests.rs +++ b/crates/tracedecay-session-runtime/src/session_retrieval/tests.rs @@ -24,6 +24,7 @@ use tracedecay_domain::{ }; use tracedecay_lcm::contracts::{LcmDataFreshness, LcmRetrievalOutcome}; use tracedecay_session_temporal_store::SessionTemporalAccess; +use tracedecay_sessions::serving::{SessionProjectionServingState, SessionProjectionStaleReason}; use tracedecay_store::{ AnchoredObservationWrite, ObservationProjectionStore, ObservationStore, ObservationWrite, SessionRecord, SessionTemporalSnapshotRequestV1, build_observation_resolution_authorization_v1, @@ -1482,6 +1483,78 @@ async fn require_fresh_with_a_current_refresh_worker_is_not_refused_as_worker_mi } } +struct FixedRefreshServing(tracedecay_sessions::serving::SessionProjectionServingState); + +impl tracedecay_sessions::serving::SessionProjectionServingStatusPort for FixedRefreshServing { + fn serving_status(&self) -> tracedecay_sessions::serving::SessionProjectionServingStatus { + tracedecay_sessions::serving::SessionProjectionServingStatus { + state: self.0.clone(), + last_progress_at_unix_micros: None, + backlog: 0, + blocker: None, + retry_class: None, + } + } +} + +/// An empty answer is complete only once historical catch-up is current. +/// While catch-up is converging the same empty store must answer stale, or a +/// search right after a scheduled import claims nothing matches (#2512). +#[tokio::test] +async fn empty_answer_is_stale_until_historical_catch_up_is_current() { + let harness = tracedecay_global_db::tests::harness::RegisteredGlobalDbHarness::open( + "session-retrieval-empty-during-catch-up", + ) + .await; + let root = registered_profile_retrieval_root(&harness.registered); + let scope = root + .identity() + .session_request_scope() + .expect("profile session scope"); + let context = admitted_lookup_context(scope); + let answer = |state: SessionProjectionServingState| { + let service = DaemonSessionRetrievalService::new_admitted_profile( + harness.registered.clone(), + root.identity().clone(), + Some(std::sync::Arc::new(FixedRefreshServing(state))), + ) + .expect("registered retrieval service"); + let query = SessionTemporalQuery::new( + SessionId::new("session.empty.catch-up").expect("session identity"), + None, + "", + None, + TemporalModeV1::Current, + tracedecay_domain::RetrievalGrainV1::Occurrence, + 1, + DiversityLimits::unbounded(), + ContextBudget { + max_bytes: APPLICATION_RETRIEVAL_MAX_BYTES, + max_tokens: APPLICATION_RETRIEVAL_MAX_BYTES / 4, + estimator_version: "words-v1".to_owned(), + }, + ) + .expect("temporal query") + .with_execution_limits(admitted_execution_limits(1)); + let context = &context; + async move { service.retrieve_admitted(context, query).await } + }; + + let current = answer(SessionProjectionServingState::Current).await; + assert!( + matches!(current, SessionRetrievalServiceOutcome::CompleteZero { .. }), + "a current empty store answers complete zero: {current:?}" + ); + let converging = answer(SessionProjectionServingState::Stale { + reason: SessionProjectionStaleReason::HistoricalConvergence, + }) + .await; + assert!( + matches!(converging, SessionRetrievalServiceOutcome::Stale { .. }), + "an empty store during catch-up must not claim completeness: {converging:?}" + ); +} + /// yielding a continuation while records remain. #[tokio::test] async fn small_lookup_reads_a_session_larger_than_the_response_budget() { diff --git a/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs b/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs index f81ca08489..bd31a3c403 100644 --- a/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs +++ b/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs @@ -110,25 +110,14 @@ async fn import_receipt( let SessionSyncOutcomeV1::Accepted(admission) = accepted else { panic!("import was not admitted: {accepted:?}"); }; - let started = tokio::time::Instant::now(); let control = SessionSyncControlV1::new(scope, admission.idempotency_key); - // Well under the 60s request deadline the old waiter would have consumed, - // with room for a loaded runner: this bounds hand-off latency, not CPU. - let hand_off_bound = Duration::from_secs(10); + // No worker loop runs behind these states, so an import that waited for + // the pass it scheduled could only end at the request deadline as + // `timed_out`. The callers' literal terminations are the hand-off proof. loop { match SessionSyncServicePort::status(&service, control.clone()).await { - SessionSyncOutcomeV1::Complete(receipt) => { - assert!( - started.elapsed() < hand_off_bound, - "import consumed its observation bound instead of returning the settled catch-up" - ); - return (receipt, project_state); - } + SessionSyncOutcomeV1::Complete(receipt) => return (receipt, project_state), SessionSyncOutcomeV1::Accepted(_) | SessionSyncOutcomeV1::Joined(_) => { - assert!( - started.elapsed() < hand_off_bound, - "import stayed pending while catch-up state was already known" - ); tokio::time::sleep(Duration::from_millis(10)).await; } other => panic!("import ended without a coverage receipt: {other:?}"), @@ -148,11 +137,12 @@ fn remaining_work( #[tokio::test] async fn import_reports_deferred_progress_while_historical_catch_up_is_still_pending() { - let (receipt, _) = import_receipt(CatchUp::Pending, "import-pending").await; + let (receipt, state) = import_receipt(CatchUp::Pending, "import-pending").await; assert_eq!(receipt.termination, OperationTermination::Partial); assert!(receipt.failure_codes.is_empty()); assert_eq!(remaining_work(&receipt), 2); + assert!(state.take_historical_dirty()); } #[tokio::test] From e6f50377ef917431783e797d96d9ec12b73930e1 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 09:02:29 +0000 Subject: [PATCH 17/18] test(sessions): share one fixed serving-status fixture --- .../src/session_retrieval/tests.rs | 20 ++++--------------- 1 file changed, 4 insertions(+), 16 deletions(-) diff --git a/crates/tracedecay-session-runtime/src/session_retrieval/tests.rs b/crates/tracedecay-session-runtime/src/session_retrieval/tests.rs index d69d0ab7a7..25b07e1efc 100644 --- a/crates/tracedecay-session-runtime/src/session_retrieval/tests.rs +++ b/crates/tracedecay-session-runtime/src/session_retrieval/tests.rs @@ -1403,20 +1403,6 @@ async fn describe_without_a_refresh_worker_does_not_pretend_history_is_convergin } } -struct CurrentRefreshServing; - -impl tracedecay_sessions::serving::SessionProjectionServingStatusPort for CurrentRefreshServing { - fn serving_status(&self) -> tracedecay_sessions::serving::SessionProjectionServingStatus { - tracedecay_sessions::serving::SessionProjectionServingStatus { - state: tracedecay_sessions::serving::SessionProjectionServingState::Current, - last_progress_at_unix_micros: None, - backlog: 0, - blocker: None, - retry_class: None, - } - } -} - /// Profile catch-up mounts the real refresh worker's serving-status port. When /// that port reports current, `RequireFresh` must not be refused as /// `RefreshWorkerMissing`. @@ -1445,7 +1431,9 @@ async fn require_fresh_with_a_current_refresh_worker_is_not_refused_as_worker_mi let service = DaemonSessionRetrievalService::new_admitted_profile( harness.registered.clone(), root.identity().clone(), - Some(std::sync::Arc::new(CurrentRefreshServing)), + Some(std::sync::Arc::new(FixedRefreshServing( + SessionProjectionServingState::Current, + ))), ) .expect("registered retrieval service"); let context = admitted_lookup_context(scope); @@ -1483,7 +1471,7 @@ async fn require_fresh_with_a_current_refresh_worker_is_not_refused_as_worker_mi } } -struct FixedRefreshServing(tracedecay_sessions::serving::SessionProjectionServingState); +struct FixedRefreshServing(SessionProjectionServingState); impl tracedecay_sessions::serving::SessionProjectionServingStatusPort for FixedRefreshServing { fn serving_status(&self) -> tracedecay_sessions::serving::SessionProjectionServingStatus { From adf071ce710784f2fbcf5b1f88bf926681df8e3f Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 10:10:30 +0000 Subject: [PATCH 18/18] test(cli): wait for code-index readiness before graph reads --- .../tool_surface_transport_test.rs | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/crates/tracedecay-cli/tests/core_cli_suite/tool_surface_transport_test.rs b/crates/tracedecay-cli/tests/core_cli_suite/tool_surface_transport_test.rs index 0bdac44c7e..c5f8335f42 100644 --- a/crates/tracedecay-cli/tests/core_cli_suite/tool_surface_transport_test.rs +++ b/crates/tracedecay-cli/tests/core_cli_suite/tool_surface_transport_test.rs @@ -303,6 +303,25 @@ fn application_surface_primitive_tools_resolve_the_working_directory_project() { "storage_status", r#"{"format":"json"}"#, ); + // `init` returns before the sealed generation serves, and until then the + // graph read answers a typed `application.code-graph.unavailable`. Wait + // for the published readiness point, then require the resolved route. + let ready = run_surface_tool_from( + &home_path, + &project_path, + "status", + &format!( + r#"{{"wait_for":{{"state":"ready","timeout_ms":{}}},"format":"json"}}"#, + SURFACE_TIMEOUT.as_millis() + ), + ); + assert_eq!( + ready.payload()["wait"]["outcome"], + "reached", + "code index never became ready\nstdout:\n{}\nstderr:\n{}", + ready.stdout, + ready.stderr + ); assert_surface_resolves_project( &home_path, &project_path,