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]