From e1d55899dbc2045a7a4193f72894bd1884d8784f Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Wed, 30 Sep 2026 03:18:31 +0000 Subject: [PATCH] refactor(sessions)!: delete the Codex v1 source-cursor replay Codex admission read the retired per-session source cursor on every pass and, when the v2 cursor lagged it, replayed current user messages into the v2 source. That carried data from an old identity into the new one. Only binaries before the unified observation identity wrote that cursor, and a store holding observation rows without the unified-identity marker is already refused with the scoped `tracedecay wipe --stale --yes` reset. Admission now reads only the v2 source, and the replay-only JSONL admission plumbing (required start cursor, end-offset bound, existence precheck) and `HostAdmission::observation_receipt` are gone. BREAKING CHANGE: Codex admission ignores cursors of the retired per-session source identity; stores that still carry them beside observation rows must be reset with `tracedecay wipe --stale --yes`. --- crates/tracedecay-host-admission/src/lib.rs | 31 +- .../tracedecay-sessions/src/admission/mod.rs | 50 +- .../src/runtime/hosts/codex/observation.rs | 405 +-------- .../codex/observation/retired_source_tests.rs | 137 +++ .../src/runtime/hosts/codex/tests.rs | 161 +--- .../jsonl_observation_admission.rs | 336 +++----- .../host_event_fixture_test.rs | 785 ++---------------- 7 files changed, 311 insertions(+), 1594 deletions(-) create mode 100644 crates/tracedecay-sessions/src/runtime/hosts/codex/observation/retired_source_tests.rs diff --git a/crates/tracedecay-host-admission/src/lib.rs b/crates/tracedecay-host-admission/src/lib.rs index 16e56f4c2f..f3833e5c29 100644 --- a/crates/tracedecay-host-admission/src/lib.rs +++ b/crates/tracedecay-host-admission/src/lib.rs @@ -8,8 +8,8 @@ use std::path::PathBuf; use std::sync::{Arc, Mutex}; use tracedecay_domain::{ - CanonicalObservationIdV1, FactOwnerV1, ObservationScopeV1, ObservationSourceCursorV1, - ObservationSourceIdentityV1, RetrievalAnchorId, SanitizationReceiptV1, + FactOwnerV1, ObservationScopeV1, ObservationSourceCursorV1, ObservationSourceIdentityV1, + RetrievalAnchorId, }; use tracedecay_store::observation::{CursorAdvanceOutcome, ObservationCursorAdvance}; use tracedecay_store::{ @@ -32,8 +32,8 @@ use tracedecay_sessions::admission::{ }; use tracedecay_sessions::observation::{ AdvanceNonDurableSourceCursorRequest, CaptureObservationOutcome, CaptureObservationRequest, - ExternalSourceProjectionRetryHandleV1, ExternalSourceProjectionStateV1, GetObservationRequest, - ObservationApplication, ObservationApplicationError, ObservationCancellation, + ExternalSourceProjectionRetryHandleV1, ExternalSourceProjectionStateV1, ObservationApplication, + ObservationApplicationError, ObservationCancellation, }; use tracedecay_sessions::repository_provenance::RepositoryProvenanceAdmissionContext; use tracedecay_sessions::runtime::git_correlation::{ @@ -311,29 +311,6 @@ impl tracedecay_sessions::admission::HostAdmission for HostAdmissionFacade<'_> { Box::pin(HostAdmissionFacade::get_source_cursor(self, source, scope)) } - fn observation_receipt<'a>( - &'a self, - provider: &'a str, - scope: &'a ObservationScopeV1, - observation_id: &'a CanonicalObservationIdV1, - cancellation: &'a ObservationCancellation, - ) -> tracedecay_sessions::admission::AdmissionFuture<'a, Option> { - Box::pin(async move { - let application = self.application(provider, scope)?; - application - .get_observation(GetObservationRequest::new( - observation_id.clone(), - cancellation.clone(), - )) - .await - .map(|read| { - read.observation() - .map(|stored| stored.observation().receipt().clone()) - }) - .map_err(|error| classify_error(&error)) - }) - } - fn drain_projection_queue<'a>( &'a self, provider: &'a str, diff --git a/crates/tracedecay-sessions/src/admission/mod.rs b/crates/tracedecay-sessions/src/admission/mod.rs index e9c199387e..87112dfbc9 100644 --- a/crates/tracedecay-sessions/src/admission/mod.rs +++ b/crates/tracedecay-sessions/src/admission/mod.rs @@ -19,8 +19,7 @@ use std::pin::Pin; use serde::Serialize; use tracedecay_domain::{ - CanonicalObservationIdV1, ObservationScopeV1, ObservationSourceCursorV1, - ObservationSourceIdentityV1, SanitizationReceiptV1, + ObservationScopeV1, ObservationSourceCursorV1, ObservationSourceIdentityV1, }; use tracedecay_store::ParseOffset; use tracedecay_store::observation::{CursorAdvanceOutcome, ObservationCursorAdvance}; @@ -339,15 +338,6 @@ impl HostAdmissionOutcome { Some("parse_offset_conflict"), ) } - - #[hotpath::skip] - pub const fn observation_point_read_unavailable() -> Self { - Self::new( - HostAdmissionStatus::Unavailable, - true, - Some("observation_point_read_unavailable"), - ) - } } pub(crate) fn is_admission_cancellation( @@ -408,19 +398,6 @@ pub trait HostAdmission: Send + Sync { scope: &'a ObservationScopeV1, ) -> AdmissionFuture<'a, Option>; - /// Content-free point existence check for idempotent recovery. The - /// default is a typed unavailable state; production composition overrides - /// it with the canonical observation authority. - fn observation_receipt<'a>( - &'a self, - _provider: &'a str, - _scope: &'a ObservationScopeV1, - _observation_id: &'a CanonicalObservationIdV1, - _cancellation: &'a ObservationCancellation, - ) -> AdmissionFuture<'a, Option> { - Box::pin(async { Err(HostAdmissionOutcome::observation_point_read_unavailable()) }) - } - /// Drains up to `max` queued projections for one provider. fn drain_projection_queue<'a>( &'a self, @@ -607,8 +584,7 @@ pub(crate) mod test_support { type SessionBackfillPagePause = (Arc, Arc); use crate::observation::{ - AdvanceNonDurableSourceCursorRequest, GetObservationRequest, ObservationApplication, - ObservationApplicationError, + AdvanceNonDurableSourceCursorRequest, ObservationApplication, ObservationApplicationError, }; use super::*; @@ -1136,28 +1112,6 @@ pub(crate) mod test_support { }) } - fn observation_receipt<'a>( - &'a self, - _provider: &'a str, - _scope: &'a ObservationScopeV1, - observation_id: &'a CanonicalObservationIdV1, - cancellation: &'a ObservationCancellation, - ) -> AdmissionFuture<'a, Option> { - Box::pin(async move { - self.application()? - .get_observation(GetObservationRequest::new( - observation_id.clone(), - cancellation.clone(), - )) - .await - .map(|read| { - read.observation() - .map(|stored| stored.observation().receipt().clone()) - }) - .map_err(Self::application_error) - }) - } - fn drain_projection_queue<'a>( &'a self, provider: &'a str, diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/observation.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/observation.rs index b02eecaab6..657bf4f1eb 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/observation.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/observation.rs @@ -11,14 +11,13 @@ pub(super) use tracedecay_capture::codex::codex_native_record_id; #[cfg(test)] pub use tracedecay_capture::codex::normalize_codex_observation; use tracedecay_capture::codex::{ - CodexObservationLocation, codex_current_user_message, codex_observation_record_supported, + CodexObservationLocation, codex_observation_record_supported, normalize_codex_observation_with_location, }; use tracedecay_domain::canonical_text::encode_lowercase_hex; use tracedecay_domain::{ - CanonicalObservationIdV1, ObservationIdentityMaterialV1, ObservationOrderingDomainV1, - ObservationScopeV1, ObservationSourceGenerationV1, ObservationSourceIdentityV1, ProjectId, - ProviderId, RetentionClass, SessionId, + ObservationScopeV1, ObservationSourceIdentityV1, ProjectId, ProviderId, RetentionClass, + SessionId, }; use tracedecay_store::observation::ObservationCoverageReason; @@ -33,16 +32,15 @@ use crate::runtime::jsonl_observation_admission::{ SharedJsonlFileIdentity, admit_jsonl_observations, reserve_shared_jsonl_page, shared_jsonl_background_cpu, shared_jsonl_file_identity, shared_jsonl_preparation_capacity, }; -use crate::runtime::shared::{StoredCursor, TranscriptScopeMatcher}; -use crate::runtime::source::{ - JsonlResumeState, MAX_JSONL_RECORD_BYTES, TranscriptIngestError, TranscriptIngestResult, - try_stream_new_jsonl_raw_strict_with_resume, -}; +use crate::runtime::shared::TranscriptScopeMatcher; +use crate::runtime::source::{TranscriptIngestError, TranscriptIngestResult}; use tracedecay_privacy::{ObservationRecordParseErrorV1, normalize_prepared_observation_record_v1}; use tracedecay_runtime_core::resident_memory::ProcessSharedMemoryReservationV1; #[cfg(test)] mod meta_cache_tests; +#[cfg(test)] +mod retired_source_tests; const CODEX_OBSERVATION_RETENTION: &str = "retention.provider-observation"; pub const CODEX_HOOK_MAX_NEW_BYTES: u64 = crate::runtime::source::MAX_JSONL_RECORD_BYTES as u64; @@ -448,41 +446,6 @@ impl CodexObservationAdmission<'_> { struct CodexAdmissionState { context: CodexContextState, scope_verdict: Option, - replay_through: Option, -} - -#[derive(Clone)] -enum CodexAdmissionMode { - Ordinary, - CurrentUserMessageReplay { - expected_start_cursor: Option, - generation: u64, - through: u64, - }, -} - -impl CodexAdmissionMode { - fn replay_through(&self, scan_generation: u64) -> Option { - match *self { - Self::CurrentUserMessageReplay { - generation, - through, - .. - } if generation == scan_generation => Some(through), - Self::Ordinary | Self::CurrentUserMessageReplay { .. } => None, - } - } - - fn replay_window(&self) -> Option<(Option, u64)> { - match self { - Self::CurrentUserMessageReplay { - expected_start_cursor, - through, - .. - } => Some((expected_start_cursor.clone(), *through)), - Self::Ordinary => None, - } - } } #[derive(Clone, Copy)] @@ -518,96 +481,6 @@ pub fn codex_observation_source_v2( )?) } -async fn replay_advanced_current_user_messages( - context: CodexAdmissionContext<'_>, - ordinary_source: &ObservationSourceIdentityV1, - canonical_source: &ObservationSourceIdentityV1, - target: &tracedecay_domain::ObservationSourceCursorV1, - max_new_bytes: Option, -) -> TranscriptIngestResult> { - let CodexAdmissionContext { - scope: admission_scope, - admission, - .. - } = context; - let scope = admission_scope.scope(); - let replay_cursor = admission - .get_source_cursor(canonical_source, &scope) - .await - .map_err(|outcome| { - crate::runtime::snapshot_observation::host_admission_error(PROVIDER, outcome) - })?; - let replay_position = match replay_cursor.as_ref() { - Some(cursor) if cursor.generation() == target.generation() => cursor.position(), - Some(_) => return Ok(None), - None => 0, - }; - let remaining = target.position().saturating_sub(replay_position); - if remaining == 0 { - return Ok(None); - } - let replay_limit = max_new_bytes - .unwrap_or(CODEX_HOOK_MAX_NEW_BYTES) - .min(CODEX_HOOK_MAX_NEW_BYTES) - .min(remaining); - let mut progress = admit_codex_jsonl_page( - context, - canonical_source.clone(), - Some((ordinary_source, target.generation())), - Some(replay_limit), - None, - CodexAdmissionMode::CurrentUserMessageReplay { - expected_start_cursor: replay_cursor, - generation: target.generation().generation_id(), - through: target.position(), - }, - ) - .await?; - // This invocation was reserved for historical catch-up. Ordinary - // admission resumes on the next bounded scheduler pass. - progress.source_deferred = true; - Ok(Some(progress)) -} - -async fn legacy_cursor_matches_current_file( - path: &Path, - target: &tracedecay_domain::ObservationSourceCursorV1, -) -> TranscriptIngestResult { - if target.position() == 0 { - return Ok(true); - } - let (Some(file_identity), Some(fingerprint)) = - (target.file_identity(), target.resume_fingerprint()) - else { - return Ok(false); - }; - let generation = target.generation().generation_id(); - let position = target.position(); - let path = path.to_path_buf(); - let scan = tokio::task::spawn_blocking(move || { - try_stream_new_jsonl_raw_strict_with_resume( - &path, - StoredCursor { - position, - mtime: 0, - file_id: generation, - }, - Some(0), - MAX_JSONL_RECORD_BYTES, - Some(JsonlResumeState { - generation, - file_identity, - fingerprint, - }), - ) - }) - .await - .map_err(|_| TranscriptIngestError::BlockingScanTaskFailed { provider: PROVIDER })??; - Ok(scan.start_offset == position - && scan.new_cursor.file_id == generation - && !scan.replacement_generation) -} - async fn shared_session_meta_with_provenance( path: &Path, cancellation: &ObservationCancellation, @@ -716,10 +589,6 @@ async fn try_admit_codex_jsonl_observations( if !admission_scope.accepts_session(&meta.session_id) { return Ok(CodexJsonlAdmissionProgress::default()); } - let ordinary_source = ObservationSourceIdentityV1::for_provider( - ProviderId::new(PROVIDER)?, - SessionId::new(meta.session_id.clone())?, - )?; let canonical_source = codex_observation_source_v2(&meta.session_id)?; let context = CodexAdmissionContext { path, @@ -735,54 +604,14 @@ async fn try_admit_codex_jsonl_observations( // this pass is about to commit. let gate = codex_admission_gate(&scope, path); let _admitting = gate.lock().await; - if let Some(target) = admission - .get_source_cursor(&ordinary_source, &scope) - .await - .map_err(|outcome| { - crate::runtime::snapshot_observation::host_admission_error(PROVIDER, outcome) - })? - { - let canonical_cursor = admission - .get_source_cursor(&canonical_source, &scope) - .await - .map_err(|outcome| { - crate::runtime::snapshot_observation::host_admission_error(PROVIDER, outcome) - })?; - let replay_is_pending = canonical_cursor.as_ref().is_none_or(|cursor| { - cursor.generation() == target.generation() && cursor.position() < target.position() - }); - if replay_is_pending - && legacy_cursor_matches_current_file(path, &target).await? - && let Some(progress) = replay_advanced_current_user_messages( - context, - &ordinary_source, - &canonical_source, - &target, - max_new_bytes, - ) - .await? - { - return Ok(progress); - } - } - admit_codex_jsonl_page( - context, - canonical_source, - None, - max_new_bytes, - max_frames, - CodexAdmissionMode::Ordinary, - ) - .await + admit_codex_jsonl_page(context, canonical_source, max_new_bytes, max_frames).await } async fn admit_codex_jsonl_page( context: CodexAdmissionContext<'_>, source: ObservationSourceIdentityV1, - ordinary_identity: Option<(&ObservationSourceIdentityV1, ObservationSourceGenerationV1)>, max_new_bytes: Option, max_frames: Option, - mode: CodexAdmissionMode, ) -> TranscriptIngestResult { let CodexAdmissionContext { path, @@ -814,11 +643,6 @@ async fn admit_codex_jsonl_page( if let Some(max_frames) = max_frames { request = request.with_max_frames(max_frames); } - if let Some((expected_start_cursor, through)) = mode.replay_window() { - request = request - .with_required_start_cursor(expected_start_cursor) - .with_max_end_offset(through); - } let progress = admit_jsonl_observations( request, |scan| { @@ -834,12 +658,10 @@ async fn admit_codex_jsonl_page( CodexAdmissionState { context, scope_verdict: None, - replay_through: mode.replay_through(scan.generation), } }, |state, _bytes, range, _, prepared, hints| { let mut stable_record_id = None; - let mut ordinary_observation_id = None; let mut non_durable_reason = None; // Scope is consulted before the record is decoded, not after. A // rollout that belongs to another project answers the same verdict @@ -874,35 +696,8 @@ async fn admit_codex_jsonl_page( non_durable_reason = Some(ObservationCoverageReason::UnsupportedFact); return Err(ObservationRecordParseErrorV1::NormalizationFailed); } - let replaying_frame = state - .replay_through - .is_some_and(|through| range.end() <= through); - if replaying_frame { - let payload = native.get("payload").unwrap_or(native); - if codex_current_user_message(payload).is_none() { - non_durable_reason = Some(ObservationCoverageReason::UnsupportedFact); - return Err(ObservationRecordParseErrorV1::NormalizationFailed); - } - } let record_id = codex_native_record_id(&meta.session_id, native) .map_err(|_| ObservationRecordParseErrorV1::NormalizationFailed)?; - if replaying_frame { - let (ordinary_source, ordinary_generation) = ordinary_identity - .ok_or(ObservationRecordParseErrorV1::NormalizationFailed)?; - let identity = ObservationIdentityMaterialV1::for_native_record( - ordinary_source.clone(), - scope.clone(), - ordinary_generation, - range, - ObservationOrderingDomainV1::FileBytes, - record_id.clone(), - ) - .map_err(|_| ObservationRecordParseErrorV1::NormalizationFailed)?; - ordinary_observation_id = Some( - CanonicalObservationIdV1::derive(&identity) - .map_err(|_| ObservationRecordParseErrorV1::NormalizationFailed)?, - ); - } let envelope = normalize_codex_observation_with_location( native, &meta.session_id, @@ -918,15 +713,7 @@ async fn admit_codex_jsonl_page( Ok(parsed) => { let record_id = stable_record_id .ok_or(TranscriptIngestError::InvalidFrameState { provider: PROVIDER })?; - if let Some(observation_id) = ordinary_observation_id { - Ok(JsonlFrameAdmission::durable_unless_observation_exists( - parsed, - record_id, - observation_id, - )) - } else { - Ok(JsonlFrameAdmission::durable(parsed, record_id)) - } + Ok(JsonlFrameAdmission::durable(parsed, record_id)) } Err(_) => Ok(JsonlFrameAdmission::non_durable( non_durable_reason.unwrap_or(ObservationCoverageReason::MalformedFrame), @@ -947,177 +734,3 @@ async fn admit_codex_jsonl_page( resumed: progress.resumed, }) } - -#[cfg(test)] -#[allow(clippy::unwrap_used)] -mod replay_boundary_tests { - use std::io::Write as _; - - use serde_json::json; - use tempfile::TempDir; - use tracedecay_domain::{CanonicalObservationEnvelopeV1, CanonicalObservationFactV1}; - - use super::*; - use crate::admission::test_support::MemoryHostAdmission; - - #[tokio::test] - async fn stale_replay_stops_when_peer_and_legacy_writer_advance() { - crate::runtime::jsonl_observation_admission::install_test_shared_jsonl_preparation_authority(); - let tmp = TempDir::new().unwrap(); - let project = tmp.path().join("project"); - std::fs::create_dir_all(&project).unwrap(); - let path = tmp.path().join("rollout.jsonl"); - let session_id = "peer-winner-session"; - let lines = [ - json!({ - "timestamp": "2026-09-04T12:00:00.000Z", - "type": "session_meta", - "payload": {"id": session_id, "cwd": project} - }), - json!({ - "timestamp": "2026-09-04T12:00:01.000Z", - "type": "event_msg", - "payload": { - "type": "item_completed", - "item": { - "type": "UserMessage", - "id": "peer-winner-item", - "content": [{"type": "text", "text": "Admit before peer race."}] - } - } - }), - ]; - std::fs::write( - &path, - lines - .iter() - .map(ToString::to_string) - .collect::>() - .join("\n") - + "\n", - ) - .unwrap(); - let admission = MemoryHostAdmission::default(); - let project_id = ProjectId::new("project.peer-winner").unwrap(); - let canonical_source = codex_observation_source_v2(session_id).unwrap(); - let scope = ObservationScopeV1::Project { - project_id: project_id.clone(), - }; - let stale_start_cursor = admission - .get_source_cursor(&canonical_source, &scope) - .await - .unwrap(); - assert!(stale_start_cursor.is_none()); - try_admit_codex_jsonl_observations_for_project_with_admission( - &path, - &project, - project_id.clone(), - &admission, - None, - ) - .await - .unwrap(); - let target = admission - .get_source_cursor(&canonical_source, &scope) - .await - .unwrap() - .expect("peer winner cursor"); - let usage = json!({ - "timestamp": "2026-09-04T12:00:02.000Z", - "type": "event_msg", - "payload": { - "type": "token_count", - "info": {"last_token_usage": { - "input_tokens": 17, - "output_tokens": 3, - "cached_input_tokens": 2, - "reasoning_output_tokens": 1, - "total_tokens": 20 - }} - } - }); - let usage_line = usage.to_string() + "\n"; - let mut file = std::fs::OpenOptions::new() - .append(true) - .open(&path) - .unwrap(); - file.write_all(usage_line.as_bytes()).unwrap(); - drop(file); - - let cancellation = ObservationCancellation::default(); - let parsed_meta = shared_session_meta_with_provenance(&path, &cancellation) - .await - .unwrap(); - let admission_scope = CodexObservationAdmission::Project { - root: &project, - project_id: project_id.clone(), - }; - let ordinary_source = ObservationSourceIdentityV1::for_provider( - ProviderId::new(PROVIDER).unwrap(), - SessionId::new(session_id).unwrap(), - ) - .unwrap(); - let context = CodexAdmissionContext { - path: &path, - scope: &admission_scope, - admission: &admission, - meta: &parsed_meta.meta, - native_thread_id: parsed_meta.native_thread_id.as_deref(), - cancellation: &cancellation, - }; - let old_writer = admit_codex_jsonl_page( - context, - ordinary_source.clone(), - None, - None, - None, - CodexAdmissionMode::Ordinary, - ) - .await - .unwrap(); - assert_eq!(old_writer.frames_persisted, 3); - let usage_count = || { - admission - .observations() - .iter() - .filter(|stored| { - serde_json::from_value::( - stored.observation().payload().clone(), - ) - .is_ok_and(|envelope| { - envelope.facts().iter().any(|fact| { - matches!(fact, CanonicalObservationFactV1::ProviderUsage { .. }) - }) - }) - }) - .count() - }; - assert_eq!(usage_count(), 1); - - let stale_replay = admit_codex_jsonl_page( - context, - canonical_source.clone(), - Some((&ordinary_source, target.generation())), - Some(usage_line.len() as u64), - None, - CodexAdmissionMode::CurrentUserMessageReplay { - expected_start_cursor: stale_start_cursor, - generation: target.generation().generation_id(), - through: target.position(), - }, - ) - .await - .unwrap(); - assert!(stale_replay.source_deferred); - assert_eq!(stale_replay.bytes_consumed, 0); - assert_eq!(stale_replay.frames_persisted, 0); - - let refreshed = try_admit_codex_jsonl_observations_for_project_with_admission( - &path, &project, project_id, &admission, None, - ) - .await - .unwrap(); - assert!(refreshed.source_deferred); - assert_eq!(usage_count(), 1); - } -} diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/observation/retired_source_tests.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/observation/retired_source_tests.rs new file mode 100644 index 0000000000..b21a8c94f8 --- /dev/null +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/observation/retired_source_tests.rs @@ -0,0 +1,137 @@ +//! The retired per-session Codex source identity has no admission authority. +//! +//! Stores that hold observation rows written under it predate the unified +//! observation identity and are refused with the scoped session-store reset +//! before admission runs. A leftover cursor for it in an admitted store must +//! therefore neither gate nor narrow what the canonical source admits. + +use serde_json::json; +use tempfile::TempDir; +use tracedecay_domain::ObservationSourceRangeV1; +use tracedecay_store::observation::ObservationCursorAdvance; + +use super::*; +use crate::admission::test_support::MemoryHostAdmission; +use crate::runtime::hosts::codex::try_admit_codex_jsonl_observations_for_project_with_admission; +use crate::runtime::observation::jsonl_observation_admission::install_test_shared_jsonl_preparation_authority; + +const SESSION_ID: &str = "retired-source-session"; + +fn write_rollout(dir: &Path, project: &Path) -> PathBuf { + let path = dir.join("rollout.jsonl"); + let lines = [ + json!({ + "timestamp": "2026-09-03T21:08:01.000Z", + "type": "session_meta", + "payload": {"id": SESSION_ID, "cwd": project} + }), + json!({ + "timestamp": "2026-09-03T21:08:01.250Z", + "type": "event_msg", + "payload": {"type": "token_count", "info": {"last_token_usage": { + "input_tokens": 10, + "output_tokens": 2, + "cached_input_tokens": 3, + "reasoning_output_tokens": 1, + "total_tokens": 12 + }}} + }), + json!({ + "timestamp": "2026-09-03T21:08:01.300Z", + "type": "event_msg", + "payload": {"type": "item_completed", "item": { + "type": "UserMessage", + "id": "retired-source-user-item", + "content": [{"type": "text", "text": "Admit every frame once."}] + }} + }), + ]; + std::fs::write( + &path, + lines + .iter() + .map(ToString::to_string) + .collect::>() + .join("\n") + + "\n", + ) + .unwrap(); + path +} + +#[tokio::test] +async fn retired_source_cursor_neither_gates_nor_narrows_canonical_admission() { + install_test_shared_jsonl_preparation_authority(); + let tmp = TempDir::new().unwrap(); + let project = tmp.path().join("project"); + std::fs::create_dir_all(&project).unwrap(); + let path = write_rollout(tmp.path(), &project); + let project_id = ProjectId::new("project.retired-codex-source").unwrap(); + let scope = ObservationScopeV1::Project { + project_id: project_id.clone(), + }; + + // The file's generation checkpoint, as any admission of it records it. + let scratch = MemoryHostAdmission::default(); + try_admit_codex_jsonl_observations_for_project_with_admission( + &path, + &project, + project_id.clone(), + &scratch, + None, + ) + .await + .unwrap(); + let checkpoint = scratch + .get_source_cursor(&codex_observation_source_v2(SESSION_ID).unwrap(), &scope) + .await + .unwrap() + .unwrap(); + + let admission = MemoryHostAdmission::default(); + let retired_source = ObservationSourceIdentityV1::for_provider( + ProviderId::new(PROVIDER).unwrap(), + SessionId::new(SESSION_ID).unwrap(), + ) + .unwrap(); + admission + .advance_non_durable_source_cursor( + ObservationCursorAdvance::new( + retired_source, + scope.clone(), + checkpoint.generation(), + None, + ObservationSourceRangeV1::new(0, checkpoint.position()).unwrap(), + ObservationCoverageReason::UnsupportedFact, + ) + .unwrap() + .with_resume_checkpoint( + checkpoint.file_identity().unwrap(), + checkpoint.resume_fingerprint().unwrap(), + ), + ObservationCancellation::default(), + ) + .await + .unwrap(); + + let first = try_admit_codex_jsonl_observations_for_project_with_admission( + &path, + &project, + project_id.clone(), + &admission, + None, + ) + .await + .unwrap(); + assert!(!first.source_deferred); + assert_eq!(first.frames_persisted, 3); + assert_eq!(admission.observations().len(), 3); + + let second = try_admit_codex_jsonl_observations_for_project_with_admission( + &path, &project, project_id, &admission, None, + ) + .await + .unwrap(); + assert_eq!(second.frames_persisted, 0); + assert_eq!(admission.observations().len(), 3); +} diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs index ce96ee624f..1703afd5f2 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs @@ -23,17 +23,11 @@ mod goal_event_tests { use super::*; use serde_json::json; - use tracedecay_domain::{ - CanonicalObservationEnvelopeV1, ObservationIdentityMaterialV1, ObservationOrderingDomainV1, - ObservationScopeV1, ObservationSourceGenerationV1, ObservationSourceIdentityV1, - ObservationSourceRangeV1, ProviderId, RetentionClass, SessionId, - }; - use tracedecay_privacy::parse_normalized_observation_record_v1; - use tracedecay_store::observation::{ObservationCoverageReason, ObservationCursorAdvance}; + use tracedecay_domain::{CanonicalObservationEnvelopeV1, ObservationScopeV1}; use crate::admission::HostAdmission; use crate::admission::test_support::MemoryHostAdmission; - use crate::observation::{CaptureObservationRequest, ObservationCancellation}; + use crate::observation::ObservationCancellation; use crate::runtime::hosts::codex::{ try_admit_codex_jsonl_observations_for_project_window, try_admit_codex_jsonl_observations_for_project_with_admission, @@ -1345,157 +1339,6 @@ mod goal_event_tests { }) .collect() } - - #[tokio::test] - async fn legacy_current_message_migration_records_receipted_duplicate_coverage() { - crate::runtime::jsonl_observation_admission::install_test_shared_jsonl_preparation_authority(); - let temp = tempfile::tempdir().unwrap(); - let project = temp.path().join("project"); - std::fs::create_dir_all(&project).unwrap(); - let transcript = temp.path().join("rollout.jsonl"); - let session_id = "session-legacy-current"; - let session_meta = json!({ - "timestamp": "2026-09-04T12:00:00.000Z", - "type": "session_meta", - "payload": {"id": session_id, "cwd": project} - }) - .to_string(); - let current = json!({ - "timestamp": "2026-09-04T12:00:01.004Z", - "type": "event_msg", - "payload": { - "type": "item_completed", - "thread_id": session_id, - "turn_id": "turn-1", - "item": { - "type": "UserMessage", - "id": "user-item-legacy", - "content": [{"type": "text", "text": "recover this prompt"}] - } - } - }); - let current_line = current.to_string(); - std::fs::write(&transcript, format!("{session_meta}\n{current_line}\n")).unwrap(); - - let source = ObservationSourceIdentityV1::for_provider( - ProviderId::new("codex").unwrap(), - SessionId::new(session_id).unwrap(), - ) - .unwrap(); - let project_id = ProjectId::new("project-legacy-current").unwrap(); - let scope = ObservationScopeV1::Project { - project_id: project_id.clone(), - }; - let scanned = crate::runtime::source::try_stream_new_jsonl_raw_strict_with_resume( - &transcript, - StoredCursor::default(), - None, - crate::runtime::source::MAX_JSONL_RECORD_BYTES, - None, - ) - .unwrap(); - assert_eq!(scanned.frames.len(), 2); - let generation = ObservationSourceGenerationV1::new(scanned.new_cursor.file_id).unwrap(); - let meta_end = u64::try_from(session_meta.len() + 1).unwrap(); - let current_end = u64::try_from(session_meta.len() + 1 + current_line.len() + 1).unwrap(); - let meta_range = ObservationSourceRangeV1::new(0, meta_end).unwrap(); - let current_range = ObservationSourceRangeV1::new(meta_end, current_end).unwrap(); - let admission = MemoryHostAdmission::default(); - let cancellation = ObservationCancellation::default(); - admission - .advance_non_durable_source_cursor( - ObservationCursorAdvance::new( - source.clone(), - scope.clone(), - generation, - None, - meta_range, - ObservationCoverageReason::UnsupportedFact, - ) - .unwrap() - .with_resume_checkpoint( - scanned.file_identity, - scanned.frames[0].resume_fingerprint, - ), - cancellation.clone(), - ) - .await - .unwrap(); - - let native_record_id = codex_native_record_id(session_id, ¤t).unwrap(); - let envelope = normalize_codex_observation( - ¤t, - session_id, - Some(session_id), - native_record_id.clone(), - current_range, - ) - .unwrap(); - let parsed = parse_normalized_observation_record_v1( - format!("{current_line}\n").as_bytes(), - current_range, - ObservationOrderingDomainV1::FileBytes, - |_| Ok(envelope), - ) - .unwrap(); - let identity = ObservationIdentityMaterialV1::for_native_record( - source, - scope.clone(), - generation, - current_range, - ObservationOrderingDomainV1::FileBytes, - native_record_id, - ) - .unwrap(); - let expected = admission - .get_source_cursor(identity.source(), &scope) - .await - .unwrap(); - admission - .capture_observation( - CaptureObservationRequest::new( - parsed, - identity, - expected, - RetentionClass::new("retention.provider-observation").unwrap(), - cancellation, - ) - .unwrap() - .with_resume_checkpoint( - scanned.file_identity, - scanned.frames[1].resume_fingerprint, - ), - ) - .await - .unwrap(); - let original = admission.observations(); - assert_eq!(original.len(), 1); - let original_receipt = original[0].observation().receipt().clone(); - - let progress = try_admit_codex_jsonl_observations_for_project_with_admission( - &transcript, - &project, - project_id, - &admission, - None, - ) - .await - .unwrap(); - - assert_eq!(progress.frames_persisted, 0); - assert_eq!(admission.observations().len(), 1); - let duplicate_advances = admission - .non_durable_advances() - .into_iter() - .filter(|advance| advance.reason() == ObservationCoverageReason::DuplicateObservation) - .collect::>(); - assert_eq!(duplicate_advances.len(), 1); - assert_eq!( - duplicate_advances[0].sanitization_receipt(), - Some(&original_receipt) - ); - assert_eq!(duplicate_advances[0].covered(), current_range); - } } #[cfg(test)] diff --git a/crates/tracedecay-sessions/src/runtime/observation/jsonl_observation_admission.rs b/crates/tracedecay-sessions/src/runtime/observation/jsonl_observation_admission.rs index bcab861baa..bd3e913a5d 100644 --- a/crates/tracedecay-sessions/src/runtime/observation/jsonl_observation_admission.rs +++ b/crates/tracedecay-sessions/src/runtime/observation/jsonl_observation_admission.rs @@ -9,10 +9,9 @@ use std::sync::{Arc, Mutex, OnceLock, PoisonError}; use tokio::sync::Notify; use tracedecay_domain::{ - CanonicalObservationIdV1, ObservationId, ObservationIdentityMaterialV1, - ObservationOrderingDomainV1, ObservationScopeV1, ObservationSourceCursorV1, - ObservationSourceGenerationV1, ObservationSourceIdentityV1, RetentionClass, - SanitizationReceiptV1, + ObservationId, ObservationIdentityMaterialV1, ObservationOrderingDomainV1, ObservationScopeV1, + ObservationSourceCursorV1, ObservationSourceGenerationV1, ObservationSourceIdentityV1, + RetentionClass, SanitizationReceiptV1, }; use tracedecay_store::observation::{ ObservationCoverageReason, ObservationCursorAdvance, ObservationIdentityCollisionDispositionV1, @@ -84,8 +83,6 @@ pub(in crate::runtime) struct JsonlObservationAdmissionRequest<'request> { retention_class: RetentionClass, max_new_bytes: Option, max_frames: Option, - required_start_cursor: Option>, - max_end_offset: Option, persisted_cursor_update: PersistedCursorUpdate, cancellation: ObservationCancellation, shared_frame_preparation: SharedJsonlFramePreparation, @@ -109,8 +106,6 @@ impl<'request> JsonlObservationAdmissionRequest<'request> { retention_class, max_new_bytes: None, max_frames: None, - required_start_cursor: None, - max_end_offset: None, persisted_cursor_update: PersistedCursorUpdate::Monotonic, cancellation: ObservationCancellation::default(), shared_frame_preparation: SharedJsonlFramePreparation::None, @@ -127,23 +122,6 @@ impl<'request> JsonlObservationAdmissionRequest<'request> { self } - /// Require this admission to start from the cursor observed by its caller. - /// A concurrent winner changes the cursor into a no-op so the caller can - /// refresh any external watermark before deciding what the next bytes mean. - pub(in crate::runtime) fn with_required_start_cursor( - mut self, - required_start_cursor: Option, - ) -> Self { - self.required_start_cursor = Some(required_start_cursor); - self - } - - /// Bound this admission to an absolute source offset. - pub(in crate::runtime) fn with_max_end_offset(mut self, max_end_offset: u64) -> Self { - self.max_end_offset = Some(max_end_offset); - self - } - pub(in crate::runtime) fn with_persisted_cursor_update( mut self, persisted_cursor_update: PersistedCursorUpdate, @@ -171,7 +149,6 @@ pub(in crate::runtime) enum JsonlFrameAdmission { parsed_record: ParsedObservationRecordV1, native_record_id: ObservationId, identity_collision_retry: bool, - skip_if_observation_exists: Option, }, NonDurable { reason: ObservationCoverageReason, @@ -211,20 +188,6 @@ impl JsonlFrameAdmission { parsed_record, native_record_id, identity_collision_retry: true, - skip_if_observation_exists: None, - } - } - - pub(in crate::runtime) fn durable_unless_observation_exists( - parsed_record: ParsedObservationRecordV1, - native_record_id: ObservationId, - observation_id: CanonicalObservationIdV1, - ) -> Self { - Self::Durable { - parsed_record, - native_record_id, - identity_collision_retry: false, - skip_if_observation_exists: Some(observation_id), } } @@ -1849,7 +1812,6 @@ struct DurableJsonlFrame { range: tracedecay_domain::ObservationSourceRangeV1, parsed_record: ParsedObservationRecordV1, native_record_id: ObservationId, - skip_if_observation_exists: Option, bytes: Arc<[u8]>, fallback_prepared: Option, fallback_hints: JsonlFrameHints, @@ -2175,29 +2137,6 @@ impl ActiveAdmission<'_> { persisted_cursor_update: PersistedCursorUpdate, ) -> TranscriptIngestResult { let checkpoint = frame.checkpoint; - let existing_receipt = match frame.skip_if_observation_exists.as_ref() { - Some(observation_id) => self - .admission - .observation_receipt( - self.provider, - &self.scope, - observation_id, - &self.cancellation, - ) - .await - .map_err(|outcome| host_admission_error(self.provider, outcome))?, - None => None, - }; - if let Some(receipt) = existing_receipt { - self.advance_coverage( - expected_cursor, - checkpoint, - ObservationCoverageReason::DuplicateObservation, - Some(receipt), - ) - .await?; - return Ok(DurableFrameDisposition::AlreadyDurable); - } let result = self .capture_attempt(expected_cursor.clone(), frame, retention_class) .await?; @@ -2232,34 +2171,6 @@ impl ActiveAdmission<'_> { if frames.is_empty() { return Ok(()); } - if frames - .iter() - .any(|frame| frame.skip_if_observation_exists.is_some()) - { - for frame in frames { - match self - .capture( - expected_cursor, - frame, - retention_class, - persisted_cursor_update, - ) - .await? - { - DurableFrameDisposition::Persisted => { - progress.frames_accepted = progress.frames_accepted.saturating_add(1); - progress.frames_persisted = progress.frames_persisted.saturating_add(1); - } - DurableFrameDisposition::Refused => { - progress.frames_refused = progress.frames_refused.saturating_add(1); - } - DurableFrameDisposition::AlreadyDurable => { - progress.frames_skipped = progress.frames_skipped.saturating_add(1); - } - } - } - return Ok(()); - } crate::runtime::pipeline_metrics::record_capture_window(frames.len()); let batch_bytes = frames.iter().try_fold(0_u64, |total, frame| { let bytes = u64::try_from(frame.bytes.len()) @@ -2383,10 +2294,8 @@ pub(in crate::runtime) async fn admit_jsonl_observations( source, scope, retention_class, - mut max_new_bytes, + max_new_bytes, max_frames, - required_start_cursor, - max_end_offset, persisted_cursor_update, cancellation, shared_frame_preparation, @@ -2408,15 +2317,6 @@ pub(in crate::runtime) async fn admit_jsonl_observations( if cancellation.is_cancelled() { return Err(TranscriptIngestError::Cancelled { provider }); } - if required_start_cursor - .as_ref() - .is_some_and(|required| required != &expected_cursor) - { - return Ok(JsonlObservationAdmissionProgress { - source_deferred: true, - ..JsonlObservationAdmissionProgress::default() - }); - } let previous = expected_cursor .as_ref() .map_or(StoredCursor::default(), |cursor| StoredCursor { @@ -2424,16 +2324,6 @@ pub(in crate::runtime) async fn admit_jsonl_observations( mtime: 0, file_id: cursor.generation().generation_id(), }); - if let Some(max_end_offset) = max_end_offset { - let remaining = max_end_offset.saturating_sub(previous.position); - if remaining == 0 { - return Ok(JsonlObservationAdmissionProgress { - source_deferred: true, - ..JsonlObservationAdmissionProgress::default() - }); - } - max_new_bytes = Some(max_new_bytes.map_or(remaining, |limit| limit.min(remaining))); - } let resume_state = expected_cursor.as_ref().and_then(|cursor| { Some(JsonlResumeState { generation: cursor.generation().generation_id(), @@ -2595,14 +2485,12 @@ pub(in crate::runtime) async fn admit_jsonl_observations( parsed_record, native_record_id, identity_collision_retry, - skip_if_observation_exists, } => { let frame = DurableJsonlFrame { checkpoint, range, parsed_record, native_record_id, - skip_if_observation_exists, bytes: Arc::clone(&bytes), fallback_prepared: None, fallback_hints: hints, @@ -2612,64 +2500,52 @@ pub(in crate::runtime) async fn admit_jsonl_observations( ObservationIdentityCollisionDispositionV1::SettleTerminal }, }; - let disposition = if frame.skip_if_observation_exists.is_some() { - active - .capture( - expected_cursor, - frame, - policy.retention_class, - policy.persisted_cursor_update, - ) - .await? - } else { - let first = active - .capture_attempt( - expected_cursor.clone(), - frame, - policy.retention_class, - ) - .await?; - match first { - Err(outcome) - if identity_collision_retry - && matches!( - &outcome, - HostAdmissionOutcome { - reason_code: Some( - "observation_identity_collision" - ), - retryable: false, - .. - } - ) => - { - // The provider explicitly proved that the - // primary identity had no native record - // id. Re-normalize from the frame's input - // state so provider carry is applied once, - // then try exactly one positional identity. - let mut retry_state = frame_start_state; - let mut retry_hints = hints; - retry_hints.identity_collision_retry = true; - let retry = normalize( - &mut retry_state, - bytes.as_ref(), - range, - checkpoint.offset, - prepared, - retry_hints, - )?; - let JsonlFrameAdmission::Durable { - parsed_record, - native_record_id, - .. - } = retry - else { - return Err(TranscriptIngestError::InvalidFrameState { - provider: active.provider, - }); - }; - let disposition = active + let first = active + .capture_attempt( + expected_cursor.clone(), + frame, + policy.retention_class, + ) + .await?; + let disposition = match first { + Err(outcome) + if identity_collision_retry + && matches!( + &outcome, + HostAdmissionOutcome { + reason_code: Some("observation_identity_collision"), + retryable: false, + .. + } + ) => + { + // The provider explicitly proved that the + // primary identity had no native record + // id. Re-normalize from the frame's input + // state so provider carry is applied once, + // then try exactly one positional identity. + let mut retry_state = frame_start_state; + let mut retry_hints = hints; + retry_hints.identity_collision_retry = true; + let retry = normalize( + &mut retry_state, + bytes.as_ref(), + range, + checkpoint.offset, + prepared, + retry_hints, + )?; + let JsonlFrameAdmission::Durable { + parsed_record, + native_record_id, + .. + } = retry + else { + return Err(TranscriptIngestError::InvalidFrameState { + provider: active.provider, + }); + }; + let disposition = active .capture( expected_cursor, DurableJsonlFrame { @@ -2677,7 +2553,6 @@ pub(in crate::runtime) async fn admit_jsonl_observations( range, parsed_record, native_record_id, - skip_if_observation_exists: None, bytes, fallback_prepared: None, fallback_hints: retry_hints, @@ -2688,19 +2563,18 @@ pub(in crate::runtime) async fn admit_jsonl_observations( policy.persisted_cursor_update, ) .await?; - fallback_state = retry_state; - disposition - } - result => { - active - .apply_capture_result( - expected_cursor, - checkpoint, - result, - policy.persisted_cursor_update, - ) - .await? - } + fallback_state = retry_state; + disposition + } + result => { + active + .apply_capture_result( + expected_cursor, + checkpoint, + result, + policy.persisted_cursor_update, + ) + .await? } }; match disposition { @@ -2857,54 +2731,47 @@ pub(in crate::runtime) async fn admit_jsonl_observations( progress.frames_decoded = progress.frames_decoded.saturating_add(1); } state = frame_state; - let (parsed_record, native_record_id, identity_collision_retry, skip_if_observation_exists) = - match admission { - JsonlFrameAdmission::Durable { - parsed_record, - native_record_id, - identity_collision_retry, - skip_if_observation_exists, - } => ( - parsed_record, - native_record_id, - identity_collision_retry, - skip_if_observation_exists, - ), - JsonlFrameAdmission::NonDurable { - reason, - before_decode, - } => { - flush_pending( - &active, - PendingAdmissionWindow { - expected_cursor: &mut expected_cursor, - frames: &mut pending, - bytes: &mut pending_bytes, - start_state: &mut pending_start_state, - progress: &mut progress, - }, - FlushPolicy { - retention_class: &retention_class, - persisted_cursor_update, - }, - &mut normalize, - ) + let (parsed_record, native_record_id, identity_collision_retry) = match admission { + JsonlFrameAdmission::Durable { + parsed_record, + native_record_id, + identity_collision_retry, + } => (parsed_record, native_record_id, identity_collision_retry), + JsonlFrameAdmission::NonDurable { + reason, + before_decode, + } => { + flush_pending( + &active, + PendingAdmissionWindow { + expected_cursor: &mut expected_cursor, + frames: &mut pending, + bytes: &mut pending_bytes, + start_state: &mut pending_start_state, + progress: &mut progress, + }, + FlushPolicy { + retention_class: &retention_class, + persisted_cursor_update, + }, + &mut normalize, + ) + .await?; + active + .advance_coverage(&mut expected_cursor, checkpoint, reason, None) .await?; - active - .advance_coverage(&mut expected_cursor, checkpoint, reason, None) - .await?; - crate::runtime::pipeline_metrics::record_frame_skipped(reason); - progress.frames_skipped = progress.frames_skipped.saturating_add(1); - if before_decode { - progress.frames_rejected_before_decode = - progress.frames_rejected_before_decode.saturating_add(1); - } - continue; - } - JsonlFrameAdmission::NeedsPreparation => { - return Err(TranscriptIngestError::InvalidFrameState { provider }); + crate::runtime::pipeline_metrics::record_frame_skipped(reason); + progress.frames_skipped = progress.frames_skipped.saturating_add(1); + if before_decode { + progress.frames_rejected_before_decode = + progress.frames_rejected_before_decode.saturating_add(1); } - }; + continue; + } + JsonlFrameAdmission::NeedsPreparation => { + return Err(TranscriptIngestError::InvalidFrameState { provider }); + } + }; let frame_bytes = u64::try_from(frame.bytes.len()) .map_err(|_| TranscriptIngestError::InvalidFrameState { provider })?; if !pending.is_empty() @@ -2932,7 +2799,6 @@ pub(in crate::runtime) async fn admit_jsonl_observations( range, parsed_record, native_record_id, - skip_if_observation_exists, bytes: Arc::clone(&frame.bytes), fallback_prepared: frame .prepared diff --git a/crates/tracedecay/tests/runtime_acceptance_suite/host_event_fixture_test.rs b/crates/tracedecay/tests/runtime_acceptance_suite/host_event_fixture_test.rs index a9ae7e632d..66d47fa8f6 100644 --- a/crates/tracedecay/tests/runtime_acceptance_suite/host_event_fixture_test.rs +++ b/crates/tracedecay/tests/runtime_acceptance_suite/host_event_fixture_test.rs @@ -26,7 +26,6 @@ use tracedecay_sessions::admission::{ }; use tracedecay_sessions::observation::{CaptureObservationRequest, ObservationCancellation}; use tracedecay_sessions::runtime::source::TranscriptSource; -use tracedecay_sessions::runtime::source::try_stream_new_jsonl_raw_strict_with_resume; use tracedecay_sessions::runtime::{hosts::claude, hosts::codex, hosts::cursor, hosts::hermes}; use tracedecay_store::ObservationReplayRequest; use tracedecay_store::observation::{ObservationCoverageReason, ObservationCursorAdvance}; @@ -1088,45 +1087,27 @@ async fn current_codex_goal_pair_projects_once_with_native_identity_and_goal_sem assert_eq!(metadata["codex_goal"]["tokens_remaining"], 11000); } +/// Binaries before the unified observation identity are the only writers of +/// the per-session Codex source cursor, so a store holding that cursor beside +/// admitted rows is a store without the unified-identity marker. Its session +/// authority is refused with the scoped reset, and its rows are neither +/// re-admitted nor rewritten. #[tokio::test] -async fn legacy_codex_cursor_migrates_to_v2_without_usage_duplication() { +async fn retired_codex_source_cursor_store_refuses_sessions_with_scoped_reset() { let tmp = TempDir::new().unwrap(); + let profile = tmp.path().join("profile"); let project = tmp.path().join("project"); std::fs::create_dir_all(&project).unwrap(); - let project_id = ProjectId::new("project.codex-current-user-replay").unwrap(); - assert!( - Command::new(git_program()) - .arg("init") - .arg(&project) - .status() - .unwrap() - .success() - ); - assert!( - tracedecay_runtime_core::storage::write_repository_identity_marker( - &project, - project_id.as_str() - ) - .unwrap() - ); - let runtime = HostAdmissionTestRuntimeV1::project( - tmp.path().join("profile"), - &project, - project_id.clone(), - ) - .await - .unwrap(); - let facade = runtime.facade(); - let transcript = tmp.path().join("codex-current-user-replay.jsonl"); - let session_id = "session-current-user-replay"; - let initial_lines = [ + let transcript = tmp.path().join("retired-source.jsonl"); + let session_id = "session-retired-codex-source"; + let lines = [ json!({"timestamp":"2026-09-03T21:08:01.000Z","type":"session_meta","payload":{"id":session_id,"cwd":project}}), json!({"timestamp":"2026-09-03T21:08:01.250Z","type":"event_msg","payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":10,"output_tokens":2,"cached_input_tokens":3,"reasoning_output_tokens":1,"total_tokens":12}}}}), - json!({"timestamp":"2026-09-03T21:08:01.300Z","type":"event_msg","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"ordinarily-admitted-user-item","content":[{"type":"text","text":"Recover replay-marker ordinary prompt."}]}}}), + json!({"timestamp":"2026-09-03T21:08:01.300Z","type":"event_msg","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"retired-source-user-item","content":[{"type":"text","text":"Admitted before the unified identity."}]}}}), ]; std::fs::write( &transcript, - initial_lines + lines .iter() .map(ToString::to_string) .collect::>() @@ -1134,80 +1115,10 @@ async fn legacy_codex_cursor_migrates_to_v2_without_usage_duplication() { + "\n", ) .unwrap(); - let ordinary_source = ObservationSourceIdentityV1::for_provider( - ProviderId::new("codex").unwrap(), - SessionId::new(session_id).unwrap(), - ) - .unwrap(); - let project_scope = ObservationScopeV1::Project { - project_id: project_id.clone(), - }; - codex::try_admit_codex_jsonl_observations_for_project_with_admission( - &transcript, - &project, - project_id.clone(), - &facade, - None, - ) - .await - .unwrap(); - let initial_length = std::fs::metadata(&transcript).unwrap().len(); - let project_observations = runtime - .replay_observations( - HostAdmissionScope::Project, - ObservationReplayRequest::new(0, 32).unwrap(), - ) - .await - .unwrap(); - let canonical_source = project_observations - .first() - .expect("ordinary v2 admission must persist project observations") - .observation() - .source() - .clone(); - assert!(canonical_source.explicit_source_key().is_some()); - let canonical_cursor = facade - .get_source_cursor(&canonical_source, &project_scope) - .await - .unwrap() - .expect("ordinary v2 admission must publish its exact cursor"); - assert_eq!(canonical_cursor.position(), initial_length); - let scope = ObservationScopeV1::Profile; - facade - .advance_non_durable_source_cursor( - ObservationCursorAdvance::new( - ordinary_source.clone(), - scope.clone(), - canonical_cursor.generation(), - None, - ObservationSourceRangeV1::new(0, initial_length).unwrap(), - ObservationCoverageReason::UnsupportedFact, - ) - .unwrap() - .with_resume_checkpoint( - canonical_cursor - .file_identity() - .expect("project cursor must retain file identity"), - canonical_cursor - .resume_fingerprint() - .expect("project cursor must retain its prefix fingerprint"), - ), - ObservationCancellation::default(), - ) - .await - .unwrap(); - drop(facade); - drop(runtime); - let runtime = HostAdmissionTestRuntimeV1::project( - tmp.path().join("profile"), - &project, - project_id.clone(), - ) - .await - .unwrap(); + let runtime = HostAdmissionTestRuntimeV1::profile(&profile).await.unwrap(); let facade = runtime.facade(); - let recovery = codex::try_admit_codex_jsonl_observations_for_profile_with_admission( + let admitted = codex::try_admit_codex_jsonl_observations_for_profile_with_admission( &transcript, Some(session_id), &[], @@ -1216,659 +1127,75 @@ async fn legacy_codex_cursor_migrates_to_v2_without_usage_duplication() { ) .await .unwrap(); - assert!( - recovery.source_deferred, - "the restart pass is reserved for bounded replay reconciliation" - ); - let recovered_observations = runtime - .replay_observations( - HostAdmissionScope::Profile, - ObservationReplayRequest::new(0, 32).unwrap(), + assert_eq!(admitted.frames_persisted, 3); + let scope = ObservationScopeV1::Profile; + let checkpoint = facade + .get_source_cursor( + &codex::codex_observation_source_v2(session_id).unwrap(), + &scope, ) .await - .unwrap(); - assert_eq!( - recovered_observations - .iter() - .filter(|stored| { - serde_json::from_value::( - stored.observation().payload().clone(), - ) - .is_ok_and(|envelope| { - envelope - .facts() - .iter() - .any(|fact| matches!(fact, CanonicalObservationFactV1::Message { .. })) - }) - }) - .count(), - 1 - ); - assert_eq!( - recovered_observations - .iter() - .filter(|stored| { - serde_json::from_value::( - stored.observation().payload().clone(), - ) - .is_ok_and(|envelope| { - envelope.facts().iter().any(|fact| { - matches!(fact, CanonicalObservationFactV1::ProviderUsage { .. }) - }) - }) - }) - .count(), - 0, - "migration recovery must not re-admit historical provider usage" - ); - let current = json!({"timestamp":"2026-09-03T21:08:01.541Z","type":"event_msg","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"replayed-user-item","content":[{"type":"text","text":"Recover replay-marker current prompt."}]}}}); - let mut file = std::fs::OpenOptions::new() - .append(true) - .open(&transcript) - .unwrap(); - writeln!(file, "{current}").unwrap(); - drop(file); - let first_advanced_length = std::fs::metadata(&transcript).unwrap().len(); - codex::try_admit_codex_jsonl_observations_for_project_with_admission( - &transcript, - &project, - project_id.clone(), - &facade, - None, - ) - .await - .unwrap(); - let project_cursor = facade - .get_source_cursor(&canonical_source, &project_scope) - .await .unwrap() - .expect("project v2 cursor must cover the appended current message"); - assert_eq!(project_cursor.position(), first_advanced_length); - let ordinary_cursor = facade - .get_source_cursor(&ordinary_source, &scope) - .await - .unwrap() - .expect("legacy cursor must survive the admission restart"); - facade - .advance_non_durable_source_cursor( - ObservationCursorAdvance::new( - ordinary_source.clone(), - scope.clone(), - ordinary_cursor.generation(), - Some(ordinary_cursor), - ObservationSourceRangeV1::new(initial_length, first_advanced_length).unwrap(), - ObservationCoverageReason::UnsupportedFact, - ) - .unwrap() - .with_resume_checkpoint( - project_cursor - .file_identity() - .expect("project cursor must retain file identity"), - project_cursor - .resume_fingerprint() - .expect("project cursor must retain its prefix fingerprint"), - ), - ObservationCancellation::default(), - ) - .await .unwrap(); - - let progress = codex::try_admit_codex_jsonl_observations_for_profile_with_admission( - &transcript, - Some(session_id), - &[], - &facade, - None, - ) - .await - .unwrap(); - assert!( - progress.source_deferred, - "replay owns one bounded admission pass" - ); - facade - .drain_projection_queue("codex", &scope, &ObservationCancellation::default(), 16) - .await - .unwrap(); - let messages = runtime - .search_session_messages_for_test( - HostAdmissionScope::Profile, - "codex", - None, - "replay-marker", - 8, - ) - .await - .unwrap(); - assert_eq!(messages.len(), 2); - assert!( - messages - .iter() - .any(|message| message.message.message_id == "replayed-user-item") - ); - let later_usage = json!({"timestamp":"2026-09-03T21:08:02.500Z","type":"event_msg","payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":20,"output_tokens":4,"cached_input_tokens":6,"reasoning_output_tokens":2,"total_tokens":24}}}}); - let later = json!({"timestamp":"2026-09-03T21:08:02.541Z","type":"event_msg","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"replayed-user-item-later","content":[{"type":"text","text":"Recover replay-marker later current prompt."}]}}}); - let mut file = std::fs::OpenOptions::new() - .append(true) - .open(&transcript) - .unwrap(); - writeln!(file, "{later_usage}").unwrap(); - writeln!(file, "{later}").unwrap(); - drop(file); - drop(messages); - drop(facade); - drop(runtime); - - let reopened = HostAdmissionTestRuntimeV1::project( - tmp.path().join("profile"), - &project, - project_id.clone(), - ) - .await - .unwrap(); - let reopened_facade = reopened.facade(); - codex::try_admit_codex_jsonl_observations_for_profile_with_admission( - &transcript, - Some(session_id), - &[], - &reopened_facade, - None, - ) - .await - .unwrap(); - reopened_facade - .drain_projection_queue("codex", &scope, &ObservationCancellation::default(), 16) - .await - .unwrap(); - let replayed = reopened - .search_session_messages_for_test( - HostAdmissionScope::Profile, - "codex", - None, - "replay-marker", - 8, - ) - .await - .unwrap(); - assert_eq!(replayed.len(), 3); - let replayed_ids = replayed - .iter() - .map(|message| message.message.message_id.as_str()) - .collect::>(); - assert_eq!( - replayed_ids, - std::collections::BTreeSet::from([ - "ordinarily-admitted-user-item", - "replayed-user-item", - "replayed-user-item-later" - ]) - ); - let replayed_observations = reopened - .replay_observations( - HostAdmissionScope::Profile, - ObservationReplayRequest::new(0, 32).unwrap(), - ) - .await - .unwrap(); - let replayed_usage_count = replayed_observations - .iter() - .filter(|stored| { - serde_json::from_value::( - stored.observation().payload().clone(), - ) - .is_ok_and(|envelope| { - envelope - .facts() - .iter() - .any(|fact| matches!(fact, CanonicalObservationFactV1::ProviderUsage { .. })) - }) - }) - .count(); - assert_eq!( - replayed_usage_count, 1, - "historical current-message recovery must not re-admit provider usage" - ); - let replayed_message_count = replayed_observations - .iter() - .filter(|stored| { - serde_json::from_value::( - stored.observation().payload().clone(), - ) - .is_ok_and(|envelope| { - envelope - .facts() - .iter() - .any(|fact| matches!(fact, CanonicalObservationFactV1::Message { .. })) - }) - }) - .count(); - assert_eq!( - replayed_message_count, 3, - "ordinary and recovery sources must not duplicate current-message observations" - ); -} - -#[tokio::test] -async fn rewritten_codex_rollout_abandons_legacy_migration_and_admits_replacement() { - let tmp = TempDir::new().unwrap(); - let project = tmp.path().join("project"); - std::fs::create_dir_all(&project).unwrap(); - let project_id = ProjectId::new("project.codex-rewritten-migration").unwrap(); - assert!( - Command::new(git_program()) - .arg("init") - .arg(&project) - .status() - .unwrap() - .success() - ); - assert!( - tracedecay_runtime_core::storage::write_repository_identity_marker( - &project, - project_id.as_str() - ) - .unwrap() - ); - let runtime = HostAdmissionTestRuntimeV1::project( - tmp.path().join("profile"), - &project, - project_id.clone(), - ) - .await - .unwrap(); - let facade = runtime.facade(); - let transcript = tmp.path().join("codex-rewritten-migration.jsonl"); - let session_id = "session-rewritten-migration"; - let session_meta = json!({ - "timestamp":"2026-09-03T21:08:01.000Z", - "type":"session_meta", - "payload":{"id":session_id,"cwd":project} - }); - let original_lines = [ - session_meta.clone(), - json!({"timestamp":"2026-09-03T21:08:01.250Z","type":"event_msg","payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":100,"output_tokens":20,"cached_input_tokens":30,"reasoning_output_tokens":10,"total_tokens":120}}}}), - json!({"timestamp":"2026-09-03T21:08:01.300Z","type":"event_msg","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"rewritten-user-item","content":[{"type":"text","text":"Original prompt whose long suffix ensures the replacement is shorter than the stale migration watermark."}]}}}), - ]; - std::fs::write( - &transcript, - original_lines - .iter() - .map(ToString::to_string) - .collect::>() - .join("\n") - + "\n", - ) - .unwrap(); - codex::try_admit_codex_jsonl_observations_for_project_with_admission( - &transcript, - &project, - project_id.clone(), - &facade, - None, - ) - .await - .unwrap(); - let project_scope = ObservationScopeV1::Project { - project_id: project_id.clone(), - }; - let project_observations = runtime - .replay_observations( - HostAdmissionScope::Project, - ObservationReplayRequest::new(0, 16).unwrap(), - ) - .await - .unwrap(); - let canonical_source = project_observations[0].observation().source().clone(); - let project_cursor = facade - .get_source_cursor(&canonical_source, &project_scope) - .await - .unwrap() - .expect("project admission cursor"); - let original_length = project_cursor.position(); - let legacy_source = ObservationSourceIdentityV1::for_provider( - ProviderId::new("codex").unwrap(), - SessionId::new(session_id).unwrap(), - ) - .unwrap(); - let profile_scope = ObservationScopeV1::Profile; facade .advance_non_durable_source_cursor( ObservationCursorAdvance::new( - legacy_source, - profile_scope.clone(), - project_cursor.generation(), + ObservationSourceIdentityV1::for_provider( + ProviderId::new("codex").unwrap(), + SessionId::new(session_id).unwrap(), + ) + .unwrap(), + scope, + checkpoint.generation(), None, - ObservationSourceRangeV1::new(0, original_length).unwrap(), + ObservationSourceRangeV1::new(0, checkpoint.position()).unwrap(), ObservationCoverageReason::UnsupportedFact, ) .unwrap() .with_resume_checkpoint( - project_cursor.file_identity().unwrap(), - project_cursor.resume_fingerprint().unwrap(), + checkpoint.file_identity().unwrap(), + checkpoint.resume_fingerprint().unwrap(), ), ObservationCancellation::default(), ) .await .unwrap(); - - let migrated = codex::try_admit_codex_jsonl_observations_for_profile_with_admission( - &transcript, - Some(session_id), - &[], - &facade, - None, - ) - .await - .unwrap(); - assert!(migrated.source_deferred); - let pre_rewrite_cursor = facade - .get_source_cursor(&canonical_source, &profile_scope) - .await - .unwrap() - .expect("v2 migration cursor"); - assert_eq!(pre_rewrite_cursor.position(), original_length); - - let replacement_lines = [ - session_meta, - json!({"timestamp":"2026-09-03T21:08:02.250Z","type":"event_msg","payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":7,"output_tokens":2,"cached_input_tokens":1,"reasoning_output_tokens":0,"total_tokens":9}}}}), - json!({"timestamp":"2026-09-03T21:08:02.300Z","type":"event_msg","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"rewritten-user-item","content":[{"type":"text","text":"Replacement prompt."}]}}}), - ]; - std::fs::write( - &transcript, - replacement_lines - .iter() - .map(ToString::to_string) - .collect::>() - .join("\n") - + "\n", - ) - .unwrap(); - let replacement_length = std::fs::metadata(&transcript).unwrap().len(); - assert!(replacement_length < original_length); - - let replacement = codex::try_admit_codex_jsonl_observations_for_profile_with_admission( - &transcript, - Some(session_id), - &[], - &facade, - None, - ) - .await - .unwrap(); - assert!( - !replacement.source_deferred, - "a stale legacy watermark must not retain migration ownership" - ); - assert_eq!(replacement.frames_persisted, 3); - let replacement_cursor = facade - .get_source_cursor(&canonical_source, &profile_scope) - .await + let db_path = runtime + .database_path(HostAdmissionScope::Profile) .unwrap() - .expect("replacement v2 cursor"); - assert_eq!(replacement_cursor.position(), replacement_length); - assert_ne!( - replacement_cursor.generation(), - pre_rewrite_cursor.generation() - ); - facade - .drain_projection_queue( - "codex", - &profile_scope, - &ObservationCancellation::default(), - 16, - ) - .await - .unwrap(); - let messages = runtime - .search_session_messages_for_test( - HostAdmissionScope::Profile, - "codex", - None, - "Replacement prompt", - 8, - ) - .await - .unwrap(); - assert_eq!(messages.len(), 1); - assert_eq!(messages[0].message.message_id, "rewritten-user-item"); - let observations = runtime - .replay_observations( - HostAdmissionScope::Profile, - ObservationReplayRequest::new(0, 32).unwrap(), - ) - .await - .unwrap(); - assert_eq!( - observations - .iter() - .filter(|stored| { - serde_json::from_value::( - stored.observation().payload().clone(), - ) - .is_ok_and(|envelope| { - envelope.facts().iter().any(|fact| { - matches!(fact, CanonicalObservationFactV1::ProviderUsage { .. }) - }) - }) - }) - .count(), - 1, - "replacement usage must be admitted instead of skipped below a stale watermark" - ); -} + .to_path_buf(); + drop(facade); + drop(runtime); -#[tokio::test] -async fn replacement_generation_supersedes_a_valid_legacy_prefix_watermark() { - let tmp = TempDir::new().unwrap(); - let project = tmp.path().join("project"); - std::fs::create_dir_all(&project).unwrap(); - let project_id = ProjectId::new("project.codex-replacement-generation").unwrap(); - assert!( - Command::new(git_program()) - .arg("init") - .arg(&project) - .status() + let observation_count = || -> i64 { + rusqlite::Connection::open(&db_path) + .unwrap() + .query_row("SELECT COUNT(*) FROM observations", [], |row| row.get(0)) .unwrap() - .success() - ); - assert!( - tracedecay_runtime_core::storage::write_repository_identity_marker( - &project, - project_id.as_str() - ) - .unwrap() - ); - let runtime = HostAdmissionTestRuntimeV1::project( - tmp.path().join("profile"), - &project, - project_id.clone(), - ) - .await - .unwrap(); - let facade = runtime.facade(); - let transcript = tmp.path().join("codex-replacement-generation.jsonl"); - let session_id = "session-replacement-generation"; - let session_meta = json!({ - "timestamp":"2026-09-03T21:08:01.000Z", - "type":"session_meta", - "payload":{"id":session_id,"cwd":project} - }); - let session_meta_line = session_meta.to_string() + "\n"; - let original_suffix = [ - json!({"timestamp":"2026-09-03T21:08:01.250Z","type":"event_msg","payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":100,"output_tokens":20,"cached_input_tokens":30,"reasoning_output_tokens":10,"total_tokens":120}}}}), - json!({"timestamp":"2026-09-03T21:08:01.300Z","type":"event_msg","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"old-generation-item","content":[{"type":"text","text":"Old generation prompt."}]}}}), - ]; - std::fs::write( - &transcript, - session_meta_line.clone() - + &original_suffix - .iter() - .map(ToString::to_string) - .collect::>() - .join("\n") - + "\n", - ) - .unwrap(); - codex::try_admit_codex_jsonl_observations_for_project_with_admission( - &transcript, - &project, - project_id.clone(), - &facade, - None, - ) - .await - .unwrap(); - let project_scope = ObservationScopeV1::Project { - project_id: project_id.clone(), }; - let project_observations = runtime - .replay_observations( - HostAdmissionScope::Project, - ObservationReplayRequest::new(0, 16).unwrap(), - ) - .await - .unwrap(); - let canonical_source = project_observations[0].observation().source().clone(); - let project_cursor = facade - .get_source_cursor(&canonical_source, &project_scope) - .await + rusqlite::Connection::open(&db_path) .unwrap() - .expect("project admission cursor"); - let prefix_scan = try_stream_new_jsonl_raw_strict_with_resume( - &transcript, - Default::default(), - Some(session_meta_line.len() as u64), - 1024 * 1024, - None, - ) - .unwrap(); - assert_eq!(prefix_scan.frames.len(), 1); - let legacy_position = prefix_scan.frames[0].end_offset; - let legacy_source = ObservationSourceIdentityV1::for_provider( - ProviderId::new("codex").unwrap(), - SessionId::new(session_id).unwrap(), - ) - .unwrap(); - let profile_scope = ObservationScopeV1::Profile; - facade - .advance_non_durable_source_cursor( - ObservationCursorAdvance::new( - legacy_source, - profile_scope.clone(), - project_cursor.generation(), - None, - ObservationSourceRangeV1::new(0, legacy_position).unwrap(), - ObservationCoverageReason::UnsupportedFact, - ) - .unwrap() - .with_resume_checkpoint( - prefix_scan.file_identity, - prefix_scan.frames[0].resume_fingerprint, - ), - ObservationCancellation::default(), + .execute( + "DELETE FROM global_schema_migrations \ + WHERE migration = 'observations-unified-identity-v1'", + [], ) - .await .unwrap(); + assert_eq!(observation_count(), 3); - let migration = codex::try_admit_codex_jsonl_observations_for_profile_with_admission( - &transcript, - Some(session_id), - &[], - &facade, - None, - ) - .await - .unwrap(); - assert!(migration.source_deferred); - codex::try_admit_codex_jsonl_observations_for_profile_with_admission( - &transcript, - Some(session_id), - &[], - &facade, - None, - ) - .await - .unwrap(); - let old_cursor = facade - .get_source_cursor(&canonical_source, &profile_scope) - .await - .unwrap() - .expect("ordinary v2 cursor"); - - let replacement_suffix = [ - json!({"timestamp":"2026-09-03T21:08:02.250Z","type":"event_msg","payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":7,"output_tokens":2,"cached_input_tokens":1,"reasoning_output_tokens":0,"total_tokens":9}}}}), - json!({"timestamp":"2026-09-03T21:08:02.300Z","type":"event_msg","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"replacement-generation-item","content":[{"type":"text","text":"Replacement generation prompt."}]}}}), - ]; - std::fs::write( - &transcript, - session_meta_line - + &replacement_suffix - .iter() - .map(ToString::to_string) - .collect::>() - .join("\n") - + "\n", - ) - .unwrap(); - let replacement = codex::try_admit_codex_jsonl_observations_for_profile_with_admission( - &transcript, - Some(session_id), - &[], - &facade, - None, - ) - .await - .unwrap(); - assert!(!replacement.source_deferred); - let replacement_cursor = facade - .get_source_cursor(&canonical_source, &profile_scope) - .await - .unwrap() - .expect("replacement v2 cursor"); - assert_ne!(replacement_cursor.generation(), old_cursor.generation()); - - let observations_before_retry = runtime - .replay_observations( - HostAdmissionScope::Profile, - ObservationReplayRequest::new(0, 32).unwrap(), - ) - .await - .unwrap() - .len(); - let retry = codex::try_admit_codex_jsonl_observations_for_profile_with_admission( - &transcript, - Some(session_id), - &[], - &facade, - None, - ) - .await - .unwrap(); - assert!( - !retry.source_deferred, - "a replacement v2 generation is the sole live cursor authority" - ); - assert_eq!(retry.bytes_consumed, 0); - assert_eq!( - facade - .get_source_cursor(&canonical_source, &profile_scope) - .await - .unwrap() - .expect("stable replacement cursor") - .generation(), - replacement_cursor.generation() - ); + let error = match HostAdmissionTestRuntimeV1::profile(&profile).await { + Ok(_) => panic!("a store holding the retired Codex source must refuse its sessions"), + Err(error) => error, + }; assert_eq!( - runtime - .replay_observations( - HostAdmissionScope::Profile, - ObservationReplayRequest::new(0, 32).unwrap(), - ) - .await - .unwrap() - .len(), - observations_before_retry + error.reset_required_context(), + Some(( + "observations", + "observation rows predate the unified observation identity and cannot be read; \ + reset the profile so ingestion can rebuild them from host transcripts" + )) ); + assert_eq!(observation_count(), 3); } #[tokio::test]