diff --git a/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs b/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs index 242fd7f99a..28545d5b60 100644 --- a/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs +++ b/crates/tracedecay-cli/src/sessions_cmd/session_sync.rs @@ -1,6 +1,6 @@ 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 +108,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 +157,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 +411,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-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-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-session-runtime/src/session_retrieval/tests.rs b/crates/tracedecay-session-runtime/src/session_retrieval/tests.rs index 5eafe338a7..25b07e1efc 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, @@ -1402,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`. @@ -1444,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); @@ -1482,6 +1471,78 @@ async fn require_fresh_with_a_current_refresh_worker_is_not_refused_as_worker_mi } } +struct FixedRefreshServing(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.rs b/crates/tracedecay-session-runtime/src/session_sync.rs index ea9c41f0c9..008b78fbb4 100644 --- a/crates/tracedecay-session-runtime/src/session_sync.rs +++ b/crates/tracedecay-session-runtime/src/session_sync.rs @@ -15,16 +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); @@ -116,6 +113,11 @@ pub enum SessionSyncWorkResult { Interrupted(work::SessionSyncInterruption), } +struct ImportHistoryObservation { + coverage: Vec, + failure_codes: Vec, +} + struct SessionSyncTerminalMaterial { termination: OperationTermination, stats: SessionSyncStatsV1, @@ -569,49 +571,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 +617,6 @@ impl DaemonSessionSyncService { .is_ok() && committed && !interrupted - && !projection_current { context.project_refresh.wake(); context.user_refresh.wake(); @@ -637,158 +625,32 @@ impl DaemonSessionSyncService { } } - async fn await_import_history( - &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)); - } - 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)); - } - }; - // 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) - } - } - - async fn await_import_projection( + /// 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, - 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)); - } - }; - 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!( - context.project_refresh.serving_status().state, - SessionProjectionServingState::Current - ) - && 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? + ) -> Result { + if let Some(interruption) = + self.observed_interruption(request.cancellation(), request.deadline()) { - Ok(()) - } else { - Err(None) - } - } - - 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; + return Err(interruption); } + let mut failure_codes = Vec::new(); + 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: transcript_import_requested_coverage(), + failure_codes, + }) } #[hotpath::skip] @@ -1096,6 +958,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 +1208,41 @@ 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 push_historical_admission_failure( + view: &HistoricalAdmissionView, + failure_codes: &mut Vec, +) { + match view { + HistoricalAdmissionView::Scheduled => {} + 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..bd31a3c403 --- /dev/null +++ b/crates/tracedecay-session-runtime/src/session_sync/import_admission_tests.rs @@ -0,0 +1,178 @@ +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 tracedecay_sessions::serving::{ + SessionProjectionServingState, SessionProjectionServingStatusPort, SessionProjectionStaleReason, +}; + +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, + 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"); + 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 control = SessionSyncControlV1::new(scope, admission.idempotency_key); + // 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) => return (receipt, project_state), + SessionSyncOutcomeV1::Accepted(_) | SessionSyncOutcomeV1::Joined(_) => { + 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, 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] +async fn import_after_current_catch_up_defers_until_the_scheduled_pass_runs() { + let (receipt, state) = import_receipt(CatchUp::Current, "import-current").await; + + // 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), 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, state) = import_receipt(CatchUp::Blocked, "import-blocked").await; + assert!(!state.take_historical_dirty()); + + 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..0f76ebe8bb 100644 --- a/crates/tracedecay-session-runtime/src/session_sync/work.rs +++ b/crates/tracedecay-session-runtime/src/session_sync/work.rs @@ -491,23 +491,14 @@ impl SessionSyncProjectContext { request: &SessionSyncRequestV1, project_sessions: RegisteredGlobalDbLeaseV1, ) -> SessionSyncWorkResult { - let history = match service.await_import_history(self, request).await { - Ok(progress) => Some(progress), - Err(Some(interruption)) => { - return SessionSyncWorkResult::Interrupted(interruption); - } - Err(None) => None, + let observation = match service.schedule_import_history(self, request) { + Ok(observation) => observation, + Err(interruption) => return SessionSyncWorkResult::Interrupted(interruption), }; 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 +521,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/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 d40bce8c2b..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 @@ -126,12 +126,10 @@ 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, - 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, @@ -160,12 +158,8 @@ 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), - 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), @@ -356,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) } @@ -404,6 +375,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 { @@ -844,10 +829,39 @@ pub struct SessionTemporalRefreshWake { route: Arc, } -#[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub struct SessionHistoricalRefreshReceipt { - pub stats: tracedecay_sessions::TranscriptIngestStats, - pub committed: bool, +/// Whether a request for historical catch-up could be handed to its worker. +#[derive(Clone, Debug, Eq, PartialEq)] +pub(crate) enum HistoricalAdmissionView { + Scheduled, + Blocked { reason_code: String }, + Unavailable, +} + +fn historical_admission_view(status: &SessionProjectionServingStatus) -> HistoricalAdmissionView { + match &status.state { + SessionProjectionServingState::Current => HistoricalAdmissionView::Scheduled, + SessionProjectionServingState::Stale { reason } => match reason { + SessionProjectionStaleReason::HistoricalConvergence + | SessionProjectionStaleReason::HistoricalRetry { .. } => { + HistoricalAdmissionView::Scheduled + } + SessionProjectionStaleReason::HistoricalBlocked { reason_code } => { + HistoricalAdmissionView::Blocked { + reason_code: reason_code.clone(), + } + } + }, + SessionProjectionServingState::Unavailable { reason } => match reason { + SessionProjectionUnavailableReason::WorkerRecovering => { + HistoricalAdmissionView::Scheduled + } + SessionProjectionUnavailableReason::WorkerMissing + | SessionProjectionUnavailableReason::WorkerStalled + | SessionProjectionUnavailableReason::WorkerStopped => { + HistoricalAdmissionView::Unavailable + } + }, + } } impl SessionTemporalRefreshWake { @@ -932,62 +946,25 @@ impl SessionTemporalRefreshWake { } } - /// Requests a fresh bounded historical-ingest cycle from the retained - /// owner and waits through its existing continuation passes. + /// 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. The pass has not run when this returns, so + /// `Scheduled` is never evidence that sources are admitted. #[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 schedule_historical_admission(&self) -> HistoricalAdmissionView { + let Some(state) = self.target() else { + return HistoricalAdmissionView::Unavailable; + }; + let view = historical_admission_view(&state.serving_status()); + if view == HistoricalAdmissionView::Scheduled { + state.mark_history_pending(); + state.wake_history(); } + view } pub fn status(&self) -> SessionTemporalRefreshWorkerStatus { 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..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, @@ -111,9 +108,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 +120,65 @@ 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 { + // 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 | HistoryContinuation::Backoff => true, + HistoryContinuation::Settled => holding, + }; + if owed { + return WindowFollowUp::PublishBeforeNextHistory; + } + match history { + HistoryContinuation::Immediate => WindowFollowUp::ContinueHistory, + HistoryContinuation::Backoff => WindowFollowUp::BackoffHistory, + HistoryContinuation::Settled if holding => WindowFollowUp::ContinueHistory, + HistoryContinuation::Settled => WindowFollowUp::Settle, + } +} + pub(super) async fn run_session_temporal_refresh_scheduler( database: RegisteredGlobalDbLeaseV1, state: Arc, @@ -143,15 +198,22 @@ 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; } 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), @@ -162,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!( @@ -218,17 +278,20 @@ 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 + // 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 +476,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); } @@ -595,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) @@ -607,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, + }, } } @@ -1263,6 +1329,88 @@ 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" + ); + assert_eq!( + window_follow_up( + HistoryContinuation::Backoff, + projection_still_unpublished(&unpublished), + projection_published_work(&unpublished), + false, + ), + WindowFollowUp::PublishBeforeNextHistory, + "a progressing retryable yield must publish before the backoff opens 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, false, false), + WindowFollowUp::BackoffHistory, + "a backoff whose projection pass moved nothing keeps the retry delay" + ); + } + #[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 397396caac..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; @@ -2115,7 +2128,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(), @@ -2129,7 +2141,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..ce96ee624f 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs @@ -1002,6 +1002,350 @@ 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" + ); + } + + /// 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" + ); + } + + /// 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, + 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(); @@ -2275,6 +2619,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() { diff --git a/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs b/crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs index d26125dcc9..6a2cb076f1 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,17 @@ 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 out-of-scope day + // before the admitted window is projected into search. Once in-scope + // 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 { if remaining == 0 { @@ -274,6 +285,16 @@ impl<'a> ProjectProviderRun<'a> { frontier_committable = false; break; } + 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; + } match codex::try_admit_codex_jsonl_observations_for_project_window( path, self.project_root, @@ -285,6 +306,9 @@ impl<'a> ProjectProviderRun<'a> { .await { Ok(progress) => { + if progress.frames_persisted > 0 && !progress.resumed { + persisted_day = path.parent(); + } deferred |= progress.source_deferred || progress.bytes_consumed > remaining; frontier_committable &= !progress.source_deferred && progress.bytes_consumed <= remaining; 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 6a8c026047..a2adf7aa0f 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( 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 f0db86838a..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()) @@ -804,11 +803,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 +819,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,31 +874,31 @@ 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}"); + 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; }