diff --git a/crates/tracedecay-cli/tests/core_cli_suite/cli_non_interactive_test.rs b/crates/tracedecay-cli/tests/core_cli_suite/cli_non_interactive_test.rs index 93cad3f576..55fc2ac099 100644 --- a/crates/tracedecay-cli/tests/core_cli_suite/cli_non_interactive_test.rs +++ b/crates/tracedecay-cli/tests/core_cli_suite/cli_non_interactive_test.rs @@ -266,26 +266,24 @@ fn sessions_unfinished_lists_workflow_state_evidence() { .await .expect("session fixture write") ); - assert!( - runtime - .upsert_session_message_for_test( - HostAdmissionScope::Project, - &MessageRecordBuilder::new( - "claude", - "message-1", - "session-1", - "assistant", - 1, - "Blocked: waiting on missing deploy credentials", - "message", - ) - .with_source(Some("/tmp/project/transcript.jsonl"), Some(1)) - .with_metadata(Some(r#"{"task_id":"task-7"}"#)) - .build(), + runtime + .upsert_session_message_for_test( + HostAdmissionScope::Project, + &MessageRecordBuilder::new( + "claude", + "message-1", + "session-1", + "assistant", + 1, + "Blocked: waiting on missing deploy credentials", + "message", ) - .await - .expect("session message fixture write") - ); + .with_source(Some("/tmp/project/transcript.jsonl"), Some(1)) + .with_metadata(Some(r#"{"task_id":"task-7"}"#)) + .build(), + ) + .await + .expect("session message fixture write"); // The daemon started below opens this database as a separate process. // Checkpoint and release the writer here so it sees the fixture rows // and can take the single-writer authority, the same discipline the diff --git a/crates/tracedecay-global-db/src/api_types.rs b/crates/tracedecay-global-db/src/api_types.rs index 651c8b82f4..74bea02f62 100644 --- a/crates/tracedecay-global-db/src/api_types.rs +++ b/crates/tracedecay-global-db/src/api_types.rs @@ -4,7 +4,6 @@ use tracedecay_runtime_core::storage::{ProjectStorageLocation, classify_registry pub use tracedecay_sessions::runtime::{ SessionActivityRow, SessionIngestHealth, SessionProviderCoverage, SessionProviderCoverageState, - TranscriptBatch, }; /// Total savings + call count for a project (or all projects when `project` is None). diff --git a/crates/tracedecay-global-db/src/git_correlation_adapter.rs b/crates/tracedecay-global-db/src/git_correlation_adapter.rs index b8ee4df061..ac7d580d77 100644 --- a/crates/tracedecay-global-db/src/git_correlation_adapter.rs +++ b/crates/tracedecay-global-db/src/git_correlation_adapter.rs @@ -325,12 +325,8 @@ impl GitCorrelationSessionStore for RegisteredGlobalDb { #[cfg(test)] mod tests { use super::GlobalDbGitCorrelationStore; - use crate::{ - ParseOffset, TranscriptPersistenceError, - tests::harness::{RegisteredGlobalDbHarness, RegisteredGlobalDbTestRuntime}, - }; + use crate::tests::harness::{RegisteredGlobalDbHarness, RegisteredGlobalDbTestRuntime}; use tracedecay_domain::ProjectId; - use tracedecay_sessions::runtime::SessionRecord; use tracedecay_sessions::runtime::git_correlation::{ CommitRelationFilter, GitCorrelationError, GitRefFilter, GitScopeFilter, SessionsForQuery, SpanObservation, SpanSource, SystemGit, @@ -479,58 +475,10 @@ mod tests { if message.contains("ProjectSessions") )); - let session = SessionRecord { - provider: "codex".to_owned(), - session_id: "profile-git-evidence".to_owned(), - project_key: "user".to_owned(), - project_path: "user".to_owned(), - title: None, - started_at: Some(1), - ended_at: Some(1), - transcript_path: None, - metadata_json: None, - parent_session_id: None, - is_subagent: false, - agent_id: None, - parent_tool_use_id: None, - }; - let error = harness - .registered - .persist_transcript_batch_with_git_evidence_result( - &session, - &[], - "profile-git-evidence.jsonl", - ParseOffset::default(), - ParseOffset::default(), - tracedecay_sessions::runtime::TranscriptGitEvidence::new( - &[], - &[SpanObservation { - provider: "codex".to_owned(), - session_id: session.session_id.clone(), - thread_id: None, - branch: Some("main".to_owned()), - worktree: "/repo".to_owned(), - ts: 1, - source: SpanSource::Ingest, - }], - ), - ) - .await - .expect_err("profile transcript authority must reject project Git evidence"); assert!(matches!( - error, - TranscriptPersistenceError::Storage { operation, source } - if operation == "record transcript git evidence" - && source.to_string().contains("ProjectSessions") + store.record_span_observation(&span_observation(10), 5).await, + Err(GitCorrelationError::Db(message)) + if message.contains("ProjectSessions") )); - assert_eq!( - harness - .registered - .get_session("codex", "profile-git-evidence") - .await - .expect("load profile session"), - None, - "scope rejection must happen before transcript rows commit" - ); } } diff --git a/crates/tracedecay-global-db/src/lib.rs b/crates/tracedecay-global-db/src/lib.rs index fede7091be..d7d6f351ce 100644 --- a/crates/tracedecay-global-db/src/lib.rs +++ b/crates/tracedecay-global-db/src/lib.rs @@ -158,7 +158,7 @@ pub use api_types::{ ProjectStoreContext, ProjectStoreResolution, RegisteredProjectRootInventoryV1, SavingsDay, SavingsTotal, SessionActivityRow, SessionIngestHealth, SessionProviderCoverage, SessionProviderCoverageState, StoreArtifactRecord, StoreArtifactUpsert, StoreInstanceRecord, - StoreInstanceUpsert, TranscriptBatch, registry_context_candidate_roots, + StoreInstanceUpsert, registry_context_candidate_roots, }; pub use support::{ AccountingMode, env_flag, env_value_truthy, global_accounting_enabled, global_accounting_mode, diff --git a/crates/tracedecay-global-db/src/observation_projection/state.rs b/crates/tracedecay-global-db/src/observation_projection/state.rs index d15cb73ae0..e8c250a1dc 100644 --- a/crates/tracedecay-global-db/src/observation_projection/state.rs +++ b/crates/tracedecay-global-db/src/observation_projection/state.rs @@ -1430,9 +1430,9 @@ fn reconcile_metadata( } } } - // Host ingest keeps the first annotation (`merge_session_metadata`). - // A later observation's source, cwd, or hook label is not a different - // session. Session identity stays on provider and session id. + // The first annotation wins: a later observation's source, cwd, or + // hook label is not a different session. Session identity stays on + // provider and session id. Some(_) => {} } } diff --git a/crates/tracedecay-global-db/src/tests/harness.rs b/crates/tracedecay-global-db/src/tests/harness.rs index f2fb28198b..9e2f0f789e 100644 --- a/crates/tracedecay-global-db/src/tests/harness.rs +++ b/crates/tracedecay-global-db/src/tests/harness.rs @@ -749,41 +749,6 @@ impl HostAdmissionTestRuntimeV1 { .await) } - pub async fn upsert_session_message_for_test( - &self, - scope: HostAdmissionScope, - message: &tracedecay_sessions::runtime::SessionMessageRecord, - ) -> tracedecay_domain::errors::Result { - let database = self.session_database_for_test(scope)?; - let session = database - .get_session(&message.provider, &message.session_id) - .await - .map_err( - |error| tracedecay_domain::errors::TraceDecayError::Database { - operation: "seed registered session message fixture".to_owned(), - message: error.to_string(), - }, - )? - .ok_or_else(|| tracedecay_domain::errors::TraceDecayError::Database { - operation: "seed registered session message fixture".to_owned(), - message: format!( - "session {}/{} is unavailable", - message.provider, message.session_id - ), - })?; - Ok(database - .upsert_transcript_batch( - &session, - std::slice::from_ref(message), - &format!( - "global-db-test-message:{}:{}", - message.provider, message.message_id - ), - crate::ParseOffset::default(), - ) - .await) - } - pub async fn session_for_test( &self, scope: HostAdmissionScope, diff --git a/crates/tracedecay-global-db/src/transcript.rs b/crates/tracedecay-global-db/src/transcript.rs index dc2bf27f52..a9424fe42b 100644 --- a/crates/tracedecay-global-db/src/transcript.rs +++ b/crates/tracedecay-global-db/src/transcript.rs @@ -1,8 +1,5 @@ use super::{ParseOffset, RegisteredGlobalDb}; -use tracedecay_sessions::runtime::{ - SessionMessageRecord, SessionRecord, SessionStoreAccess, TranscriptGitEvidence, - TranscriptPersistenceError, -}; +use tracedecay_sessions::runtime::{SessionRecord, SessionStoreAccess, TranscriptPersistenceError}; pub(super) use tracedecay_sessions::runtime::store_access::{ require_expected_offset, set_parse_offset, @@ -25,64 +22,6 @@ impl RegisteredGlobalDb { .await } - #[hotpath::measure(future = true, label = "global_db.transcript.upsert_batch")] - pub async fn upsert_transcript_batch( - &self, - session: &SessionRecord, - messages: &[SessionMessageRecord], - parse_offset_path: &str, - parse_offset: ParseOffset, - ) -> bool { - SessionStoreAccess::new(self) - .upsert_transcript_batch(session, messages, parse_offset_path, parse_offset) - .await - } - - #[hotpath::measure(future = true, label = "global_db.transcript.persist_batch")] - pub async fn persist_transcript_batch_result( - &self, - session: &SessionRecord, - messages: &[SessionMessageRecord], - parse_offset_path: &str, - expected_offset: ParseOffset, - parse_offset: ParseOffset, - ) -> Result<(), TranscriptPersistenceError> { - SessionStoreAccess::new(self) - .persist_transcript_batch_result( - session, - messages, - parse_offset_path, - expected_offset, - parse_offset, - ) - .await - } - - #[hotpath::measure( - future = true, - label = "global_db.transcript.persist_batch_with_git_evidence" - )] - pub async fn persist_transcript_batch_with_git_evidence_result( - &self, - session: &SessionRecord, - messages: &[SessionMessageRecord], - parse_offset_path: &str, - expected_offset: ParseOffset, - parse_offset: ParseOffset, - git_evidence: TranscriptGitEvidence<'_>, - ) -> Result<(), TranscriptPersistenceError> { - SessionStoreAccess::new(self) - .persist_transcript_batch_with_git_evidence_result( - session, - messages, - parse_offset_path, - expected_offset, - parse_offset, - git_evidence, - ) - .await - } - #[hotpath::measure(future = true, label = "global_db.transcript.persist_offset")] pub async fn persist_transcript_offset_result( &self, diff --git a/crates/tracedecay-project/src/test_support/host_admission.rs b/crates/tracedecay-project/src/test_support/host_admission.rs index 4b11e3ec75..c44f92b29d 100644 --- a/crates/tracedecay-project/src/test_support/host_admission.rs +++ b/crates/tracedecay-project/src/test_support/host_admission.rs @@ -557,21 +557,17 @@ impl HostAdmissionTestRuntimeV1 { .await) } + /// Seeds one raw LCM message into its already-registered session. #[doc(hidden)] #[hotpath::skip] pub async fn upsert_session_message_for_test( &self, scope: HostAdmissionScope, message: &tracedecay_sessions::runtime::SessionMessageRecord, - ) -> Result { - let database = self.session_database_for_test(scope)?; - let session = database - .get_session(&message.provider, &message.session_id) - .await - .map_err(|error| TraceDecayError::Database { - operation: "seed registered session message fixture".to_owned(), - message: error.to_string(), - })? + ) -> Result<()> { + let session = self + .session_for_test(scope, &message.provider, &message.session_id) + .await? .ok_or_else(|| TraceDecayError::Database { operation: "seed registered session message fixture".to_owned(), message: format!( @@ -579,17 +575,9 @@ impl HostAdmissionTestRuntimeV1 { message.provider, message.session_id ), })?; - Ok(database - .upsert_transcript_batch( - &session, - std::slice::from_ref(message), - &format!( - "host-admission-test-message:{}:{}", - message.provider, message.message_id - ), - tracedecay_global_db::ParseOffset::default(), - ) - .await) + self.seed_session_messages_for_test(scope, &session, std::slice::from_ref(message)) + .await + .map(|_| ()) } #[doc(hidden)] @@ -622,28 +610,31 @@ impl HostAdmissionTestRuntimeV1 { .await } + /// Seeds one session row and its raw LCM messages, returning each + /// message's raw store id in input order. #[doc(hidden)] #[hotpath::skip] - pub async fn upsert_transcript_batch_for_test( + pub async fn seed_session_messages_for_test( &self, scope: HostAdmissionScope, session: &tracedecay_sessions::runtime::SessionRecord, messages: &[tracedecay_sessions::runtime::SessionMessageRecord], - source: &str, - offset: tracedecay_global_db::ParseOffset, ) -> Result> { - let database = self.session_database_for_test(scope)?; - if !database - .upsert_transcript_batch(session, messages, source, offset) - .await - { + if !self.upsert_session_for_test(scope, session).await? { return Err(TraceDecayError::Database { - operation: "seed registered transcript batch fixture".to_owned(), - message: "registered transcript batch write failed".to_owned(), + operation: "seed registered session fixture".to_owned(), + message: "registered session write failed".to_owned(), }); } + let database = self.session_database_for_test(scope)?; let mut store_ids = Vec::with_capacity(messages.len()); for message in messages { + self.lcm_ingest_raw_message_for_test(scope, message) + .await + .map_err(|error| TraceDecayError::Database { + operation: "seed registered session message fixture".to_owned(), + message: error.to_string(), + })?; let store_id = database .lcm_raw_message_store_id(&message.provider, &message.message_id) .await diff --git a/crates/tracedecay-project/src/test_support/host_admission/session_test_support.rs b/crates/tracedecay-project/src/test_support/host_admission/session_test_support.rs index 2ef9ec8575..0aa439275a 100644 --- a/crates/tracedecay-project/src/test_support/host_admission/session_test_support.rs +++ b/crates/tracedecay-project/src/test_support/host_admission/session_test_support.rs @@ -129,33 +129,6 @@ impl HostAdmissionTestRuntimeV1 { ) } - #[doc(hidden)] - pub async fn set_parse_offset_insert_failure_for_test( - &self, - scope: HostAdmissionScope, - enabled: bool, - ) -> tracedecay_domain::errors::Result<()> { - let statement = if enabled { - "CREATE TRIGGER fail_parse_offset_insert - BEFORE INSERT ON parse_offsets - BEGIN - SELECT RAISE(ABORT, 'late parse offset failure'); - END;" - } else { - "DROP TRIGGER IF EXISTS fail_parse_offset_insert;" - }; - self.session_database_for_test(scope)? - .writer_connection()? - .execute_batch(statement) - .await - .map_err( - |error| tracedecay_domain::errors::TraceDecayError::Database { - operation: "configure registered parse-offset failure".to_owned(), - message: error.to_string(), - }, - ) - } - #[doc(hidden)] pub async fn set_project_parse_offset_for_test( &self, diff --git a/crates/tracedecay-session-memory/src/transcript.rs b/crates/tracedecay-session-memory/src/transcript.rs index 629fab3e30..edbe6a9b5f 100644 --- a/crates/tracedecay-session-memory/src/transcript.rs +++ b/crates/tracedecay-session-memory/src/transcript.rs @@ -3,12 +3,10 @@ use std::borrow::Borrow; use std::path::{Path, PathBuf}; use tracedecay_store::{ - ParseOffset, TranscriptStore, TranscriptStoreError, TranscriptStoreResult, - TranscriptWriteBatch, TranscriptWriteKind, + ParseOffset, TranscriptStore, TranscriptStoreError, TranscriptStoreResult, TranscriptWriteBatch, }; use tracedecay_global_db::{RegisteredGlobalDb, TranscriptPersistenceError}; -use tracedecay_sessions::runtime::TranscriptGitEvidence; use tracedecay_sessions::runtime::store_port::TranscriptIngestStore; /// Transcript-store adapter over an already-open authoritative @@ -76,66 +74,39 @@ where #[hotpath::skip] async fn persist_batch(&self, batch: TranscriptWriteBatch) -> TranscriptStoreResult<()> { - let (cursor_path, kind) = batch.into_parts(); + let (cursor_path, mut expected_offset, next_offset) = batch.into_parts(); let cursor_key = Self::path_text(&cursor_path); - match kind { - TranscriptWriteKind::AdvanceOffset { - expected_offset, - next_offset, - } => { - // Offset-only batches contain no parse products, so advancing - // across a compatible append winner cannot persist stale rows. - // Full batches below must never rewrite their observed cursor: - // their caller has to re-read and reparse after a conflict. - let mut expected_offset = expected_offset; - loop { - match self - .db() - .persist_transcript_offset_result(&cursor_key, expected_offset, next_offset) - .await - { - Ok(()) => return Ok(()), - Err(TranscriptPersistenceError::Conflict { expected, actual }) => { - if actual == next_offset { - return Ok(()); - } - let compatible_successor = actual.file_id != 0 - && actual.file_id == next_offset.file_id - && actual.byte_offset > expected.byte_offset - && actual.mtime >= expected.mtime - && next_offset.byte_offset > actual.byte_offset - && next_offset.mtime >= actual.mtime; - if !compatible_successor { - return Err(Self::persistence_error( - &cursor_path, - TranscriptPersistenceError::Conflict { expected, actual }, - )); - } - expected_offset = actual; - } - Err(error) => { - return Err(Self::persistence_error(&cursor_path, error)); - } + // Offset-only batches contain no parse products, so advancing across a + // compatible append winner cannot persist stale rows. + loop { + match self + .db() + .persist_transcript_offset_result(&cursor_key, expected_offset, next_offset) + .await + { + Ok(()) => return Ok(()), + Err(TranscriptPersistenceError::Conflict { expected, actual }) => { + if actual == next_offset { + return Ok(()); + } + let compatible_successor = actual.file_id != 0 + && actual.file_id == next_offset.file_id + && actual.byte_offset > expected.byte_offset + && actual.mtime >= expected.mtime + && next_offset.byte_offset > actual.byte_offset + && next_offset.mtime >= actual.mtime; + if !compatible_successor { + return Err(Self::persistence_error( + &cursor_path, + TranscriptPersistenceError::Conflict { expected, actual }, + )); } + expected_offset = actual; + } + Err(error) => { + return Err(Self::persistence_error(&cursor_path, error)); } } - TranscriptWriteKind::Upsert { - session, - messages, - expected_offset, - next_offset, - } => self - .db() - .persist_transcript_batch_with_git_evidence_result( - &session, - &messages, - &cursor_key, - expected_offset, - next_offset, - TranscriptGitEvidence::new(&[], &[]), - ) - .await - .map_err(|error| Self::persistence_error(&cursor_path, error)), } } } diff --git a/crates/tracedecay-session-runtime/src/lcm_effects/tests.rs b/crates/tracedecay-session-runtime/src/lcm_effects/tests.rs index fc08915aec..c1bf81f3ff 100644 --- a/crates/tracedecay-session-runtime/src/lcm_effects/tests.rs +++ b/crates/tracedecay-session-runtime/src/lcm_effects/tests.rs @@ -21,7 +21,7 @@ use tracedecay_runtime_core::test_executable::write_executable_script; use tracedecay_sessions::runtime::{SessionMessageRecord, SessionRecord}; use tracedecay_store::{ AnchoredObservationWrite, ObservationProjectionStore, ObservationStore, ObservationWrite, - ParseOffset, SessionRefreshStore as _, build_observation_resolution_authorization_v1, + SessionRefreshStore as _, build_observation_resolution_authorization_v1, build_observation_retrieval_anchor, derive_canonical_projection, }; @@ -775,15 +775,7 @@ async fn transcript_ingest_persists_native_compaction_raw_range() { post_compaction.message_id = "native-range-post-compaction".to_string(); messages.push(post_compaction); - assert!( - db.upsert_transcript_batch( - &session("codex", session_id), - &messages, - "/tmp/codex-native-range-session.jsonl", - ParseOffset::default(), - ) - .await - ); + seed_session_messages(&db, &session("codex", session_id), &messages).await; let first = db .lcm_raw_message_store_id("codex", "native-range-message-1") .await @@ -1081,15 +1073,7 @@ fn native_compaction_requires_exact_selected_raw_membership() { tail.provider = "codex".to_string(); tail.message_id = "native-membership-tail".to_string(); messages.push(tail); - assert!( - db.upsert_transcript_batch( - &session("codex", session_id), - &messages, - "/tmp/codex-native-membership-session.jsonl", - ParseOffset::default(), - ) - .await - ); + seed_session_messages(&db, &session("codex", session_id), &messages).await; let policy_anchor_store_id = db .lcm_raw_message_store_id("codex", "native-membership-message-2") .await @@ -1803,15 +1787,7 @@ async fn mounted_schedulers_share_historical_work_admission() { tail.provider = "codex".to_string(); tail.message_id = format!("{session_id}-tail"); messages.push(tail); - assert!( - db.upsert_transcript_batch( - &session("codex", &session_id), - &messages, - &format!("/tmp/{session_id}.jsonl"), - ParseOffset::default(), - ) - .await - ); + seed_session_messages(&db, &session("codex", &session_id), &messages).await; stores.push((harness, db, session_id)); } @@ -3296,8 +3272,7 @@ fn canonical_record_for_scope(envelope: Value, scope: ObservationScopeV1) -> Can } } -/// Ingests `messages` followed by each record's projected row through the -/// transcript path, then persists and projects each record's observation so +/// Seeds `messages` followed by each record's projected row, then persists and projects each record's observation so /// the row's envelope authority exists exactly as capture leaves it. async fn ingest_canonical( db: &RegisteredGlobalDb, @@ -3308,15 +3283,7 @@ async fn ingest_canonical( let provider = records[0].message.provider.as_str(); let mut batch = messages.to_vec(); batch.extend(records.iter().map(|record| record.message.clone())); - assert!( - db.upsert_transcript_batch( - &session(provider, session_id), - &batch, - &format!("/tmp/{session_id}.jsonl"), - ParseOffset::default(), - ) - .await - ); + seed_session_messages(db, &session(provider, session_id), &batch).await; let store = db.observation_store(); for record in records { let observation = &record.observation; @@ -3427,6 +3394,21 @@ async fn insert_summary_evidence( transaction.commit().await.unwrap(); } +/// Seeds one session row and its raw LCM messages in order. +async fn seed_session_messages( + db: &RegisteredGlobalDb, + session: &SessionRecord, + messages: &[SessionMessageRecord], +) { + assert!(db.upsert_session(session).await); + let storage_root = db.db_path().parent().unwrap(); + for message in messages { + db.lcm_ingest_raw_message(storage_root, message) + .await + .unwrap(); + } +} + async fn ingest_codex_compaction_evidence( db: &RegisteredGlobalDb, session_id: &str, @@ -3449,15 +3431,7 @@ async fn ingest_codex_compaction_evidence( let mut tail = message(session_id, ordinal.saturating_add(1)); tail.provider = "codex".to_string(); tail.message_id = format!("{message_id}-tail"); - assert!( - db.upsert_transcript_batch( - &session("codex", session_id), - &[compaction, tail], - &format!("/tmp/{session_id}.jsonl"), - ParseOffset::default(), - ) - .await - ); + seed_session_messages(db, &session("codex", session_id), &[compaction, tail]).await; } #[cfg(unix)] diff --git a/crates/tracedecay-session-runtime/src/retained/profile.rs b/crates/tracedecay-session-runtime/src/retained/profile.rs index 907d0fa97a..73f40df3c3 100644 --- a/crates/tracedecay-session-runtime/src/retained/profile.rs +++ b/crates/tracedecay-session-runtime/src/retained/profile.rs @@ -547,17 +547,10 @@ mod tests { .project_observation(observation.observation_id()) .await .expect("project canonical observation"); - assert!( - database - .upsert_transcript_batch( - &owning_session, - std::slice::from_ref(&owning_message), - &format!("profile-retained-{message_id}.jsonl"), - tracedecay_global_db::ParseOffset::default(), - ) - .await, - "seed canonical owning transcript", - ); + database + .lcm_ingest_raw_message(database.db_path().parent().unwrap(), &owning_message) + .await + .expect("seed canonical owning raw message"); database .lcm_protect_session_raw_messages(provider, session_id) .await diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs index c26f1d10fe..878c4dc649 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs @@ -1253,7 +1253,6 @@ mod tests { use tracedecay_global_db::tests::harness::RegisteredGlobalDbHarness; use tracedecay_runtime_core::db::engine::params; use tracedecay_sessions::runtime::{SessionMessageRecord, SessionRecord}; - use tracedecay_store::ParseOffset; #[test] fn deterministic_storage_refusals_are_not_retryable() { @@ -1555,16 +1554,14 @@ mod tests { } }) .collect::>(); - assert!( + assert!(database.upsert_session(&session).await); + let storage_root = database.db_path().parent().unwrap(); + for message in &messages { database - .upsert_transcript_batch( - &session, - &messages, - &format!("/tmp/{session_id}.jsonl"), - ParseOffset::default(), - ) + .lcm_ingest_raw_message(storage_root, message) .await - ); + .unwrap(); + } let transaction = database .begin_write_transaction() .await diff --git a/crates/tracedecay-sessions/src/runtime/ingest/failure.rs b/crates/tracedecay-sessions/src/runtime/ingest/failure.rs index 6158b1c3a1..b28d951edc 100644 --- a/crates/tracedecay-sessions/src/runtime/ingest/failure.rs +++ b/crates/tracedecay-sessions/src/runtime/ingest/failure.rs @@ -350,15 +350,6 @@ pub fn classify_transcript_ingest_disposition( source::TranscriptIngestError::Store(TranscriptStoreError::InvalidCursorPath) => { ("transcript_cursor_path_invalid", false, Degraded) } - source::TranscriptIngestError::Store(TranscriptStoreError::InvalidTranscriptPath) => { - ("transcript_path_invalid", false, Degraded) - } - source::TranscriptIngestError::Store(TranscriptStoreError::MissingTranscriptPath { - .. - }) => ("transcript_path_missing", false, Degraded), - source::TranscriptIngestError::Store(TranscriptStoreError::MessageIdentityMismatch { - .. - }) => ("transcript_message_identity_mismatch", false, Degraded), source::TranscriptIngestError::ScanIo { .. } => { ("transcript_source_io_failed", true, Unavailable) } diff --git a/crates/tracedecay-sessions/src/runtime/ingest/scheduler.rs b/crates/tracedecay-sessions/src/runtime/ingest/scheduler.rs index eeb14ff472..4ccb49eda7 100644 --- a/crates/tracedecay-sessions/src/runtime/ingest/scheduler.rs +++ b/crates/tracedecay-sessions/src/runtime/ingest/scheduler.rs @@ -213,7 +213,7 @@ mod tests { use tempfile::TempDir; use tracedecay_store::{ ParseOffset, TranscriptStore, TranscriptStoreError, TranscriptStoreResult, - TranscriptWriteBatch, TranscriptWriteKind, + TranscriptWriteBatch, }; use crate::runtime::hosts::codex; @@ -254,18 +254,7 @@ mod tests { &self, batch: TranscriptWriteBatch, ) -> impl std::future::Future> + Send { - let (cursor_path, kind) = batch.into_parts(); - let (expected, next) = match kind { - TranscriptWriteKind::AdvanceOffset { - expected_offset, - next_offset, - } - | TranscriptWriteKind::Upsert { - expected_offset, - next_offset, - .. - } => (expected_offset, next_offset), - }; + let (cursor_path, expected, next) = batch.into_parts(); let mut offsets = self.offsets.lock().expect("offset lock"); let actual = *offsets.get(&cursor_path).unwrap_or(&ParseOffset::default()); let result = if actual == expected { diff --git a/crates/tracedecay-sessions/src/runtime/mod.rs b/crates/tracedecay-sessions/src/runtime/mod.rs index 72db6dbb63..a2a59e7d15 100644 --- a/crates/tracedecay-sessions/src/runtime/mod.rs +++ b/crates/tracedecay-sessions/src/runtime/mod.rs @@ -49,7 +49,7 @@ pub use shared::SESSION_TRANSCRIPT_STALLED_INGEST_WARNING_BYTES; pub use snapshot_observation::SnapshotCaptureOutcome; pub use store_access::{ SessionActivityRow, SessionIngestHealth, SessionProviderCoverage, SessionProviderCoverageState, - TranscriptBatch, TranscriptGitEvidence, TranscriptPersistenceError, + TranscriptPersistenceError, }; pub use tracedecay_lcm::{SessionMessageType, SessionSearchScope}; diff --git a/crates/tracedecay-sessions/src/runtime/store_access/mod.rs b/crates/tracedecay-sessions/src/runtime/store_access/mod.rs index baefa86d29..03241ded25 100644 --- a/crates/tracedecay-sessions/src/runtime/store_access/mod.rs +++ b/crates/tracedecay-sessions/src/runtime/store_access/mod.rs @@ -21,10 +21,8 @@ pub(crate) use sessions::EXISTING_SESSION_MESSAGE_IDS_SQL; pub(crate) use sessions::SESSION_MESSAGE_ID_LOOKUP_MAX; pub use sessions::SESSION_MESSAGES_AFTER_SQL; pub use sessions::{SqlColumnError, message_record_from_row, session_record_from_row}; -pub use transcript::{ - TranscriptGitEvidence, get_parse_offset, require_expected_offset, set_parse_offset, -}; +pub use transcript::{get_parse_offset, require_expected_offset, set_parse_offset}; pub use types::{ SessionActivityRow, SessionIngestHealth, SessionProviderCoverage, SessionProviderCoverageState, - TranscriptBatch, TranscriptPersistenceError, UNIX_TIMESTAMP_MILLIS_THRESHOLD, + TranscriptPersistenceError, UNIX_TIMESTAMP_MILLIS_THRESHOLD, }; diff --git a/crates/tracedecay-sessions/src/runtime/store_access/transcript.rs b/crates/tracedecay-sessions/src/runtime/store_access/transcript.rs index e277c4063d..abac00b716 100644 --- a/crates/tracedecay-sessions/src/runtime/store_access/transcript.rs +++ b/crates/tracedecay-sessions/src/runtime/store_access/transcript.rs @@ -1,89 +1,9 @@ use tracedecay_runtime_core::db::engine::{Executor, QueryExecutor, Row, params}; -use tracedecay_store::{ParseOffset, SessionMessageRecord, SessionRecord, StoreShardScopeV1}; +use tracedecay_store::{ParseOffset, SessionRecord}; -use tracedecay_lcm::payload::PayloadFileRollback; -use tracedecay_lcm::raw; - -use super::super::git_correlation::{ - CommitSessionRecord, DEFAULT_SPAN_MERGE_GAP_SECS, GitEvidenceBatch, GitEvidenceWriter, - SpanObservation, -}; use super::super::registered_db::{SessionRegisteredDb, SessionStoreAccess, SessionWriteTxn}; use super::super::shared::{durable_project_path_key, path_identity_key}; -use super::codex_goal_reconciliation::find_preceding_codex_goal_response; -use super::types::{TranscriptBatch, TranscriptPersistenceError}; - -/// Git evidence recorded atomically with one transcript write. -#[derive(Debug, Clone, Copy)] -pub struct TranscriptGitEvidence<'a> { - commit_records: &'a [CommitSessionRecord], - span_observations: &'a [SpanObservation], -} - -impl<'a> TranscriptGitEvidence<'a> { - pub const fn new( - commit_records: &'a [CommitSessionRecord], - span_observations: &'a [SpanObservation], - ) -> Self { - Self { - commit_records, - span_observations, - } - } -} - -/// Prepares every privacy-protected raw-message write before the transaction -/// acquires SQLite's single-writer lease. -/// -/// The returned vector preserves batch/message order so the transactional -/// phase can pair each staged payload with its canonical projection row. -#[hotpath::measure(label = "sessions.store.transcript.stage_messages")] -fn stage_full_transcript_messages( - storage_root: &std::path::Path, - batches: &[TranscriptBatch], - payload_rollback: &mut PayloadFileRollback, -) -> Result, TranscriptPersistenceError> { - let message_count = batches.iter().fold(0_usize, |count, batch| { - count.saturating_add(batch.messages.len()) - }); - let mut staged = Vec::with_capacity(message_count); - for batch in batches { - for message in &batch.messages { - staged.push( - raw::stage_raw_message_with_payload_tracked( - storage_root, - message, - payload_rollback, - ) - .map_err(|error| { - TranscriptPersistenceError::storage("upsert LCM raw message", error) - })?, - ); - } - } - Ok(staged) -} - -async fn reconcile_codex_goal_response( - conn: &impl Executor, - current: &SessionMessageRecord, -) -> Result<(), TranscriptPersistenceError> { - let Some(response_message_id) = find_preceding_codex_goal_response(conn, current) - .await - .map_err(|error| { - TranscriptPersistenceError::storage("find preceding Codex goal response", error) - })? - else { - return Ok(()); - }; - conn.execute( - "DELETE FROM lcm_raw_messages WHERE provider = ?1 AND message_id = ?2", - params![current.provider.as_str(), response_message_id.as_str()], - ) - .await - .map(|_| ()) - .map_err(|error| TranscriptPersistenceError::storage("remove paired Codex goal message", error)) -} +use super::types::TranscriptPersistenceError; /// Reads one durable cursor by its canonical key. /// @@ -319,132 +239,6 @@ impl SessionStoreAccess<'_, D> { })) } - fn normalize_session_message_timestamp(timestamp: Option) -> Option { - timestamp.map(|timestamp| { - let magnitude = timestamp.unsigned_abs(); - if magnitude >= 100_000_000_000_000_000 { - timestamp / 1_000_000_000 - } else if magnitude >= 100_000_000_000_000 { - timestamp / 1_000_000 - } else if magnitude >= 100_000_000_000 { - timestamp / 1_000 - } else { - timestamp - } - }) - } - - #[hotpath::skip] - async fn upsert_session_message_in_existing_tx( - conn: &impl Executor, - message: &SessionMessageRecord, - staged: raw::StagedRawMessageIngest, - ) -> Result<(), TranscriptPersistenceError> { - // Clone-on-normalize: message text can be hundreds of kilobytes, and - // most providers already emit in-range timestamps, so the full-record - // copy is paid only when the timestamp actually changes. - let normalized_timestamp = Self::normalize_session_message_timestamp(message.timestamp); - let canonical_message = if normalized_timestamp == message.timestamp { - std::borrow::Cow::Borrowed(message) - } else { - let mut owned = message.clone(); - owned.timestamp = normalized_timestamp; - std::borrow::Cow::Owned(owned) - }; - raw::commit_staged_raw_message(conn, canonical_message.as_ref(), staged) - .await - .map(|_| ()) - .map_err(|error| TranscriptPersistenceError::storage("upsert LCM raw message", error)) - } - - /// Atomically upserts one transcript session + all parsed messages and then - /// advances the parse cursor. Any failure rolls back the entire batch so a - /// follow-up ingest can safely replay from the previous offset. - #[hotpath::skip] - pub async fn upsert_transcript_batch( - &self, - session: &SessionRecord, - messages: &[SessionMessageRecord], - parse_offset_path: &str, - parse_offset: ParseOffset, - ) -> bool { - let Ok(expected_offset) = self.get_parse_offset(parse_offset_path).await else { - return false; - }; - self.persist_transcript_batch_result( - session, - messages, - parse_offset_path, - expected_offset.unwrap_or_default(), - parse_offset, - ) - .await - .is_ok() - } - - #[hotpath::skip] - pub async fn persist_transcript_batch_result( - &self, - session: &SessionRecord, - messages: &[SessionMessageRecord], - parse_offset_path: &str, - expected_offset: ParseOffset, - parse_offset: ParseOffset, - ) -> Result<(), TranscriptPersistenceError> { - let batch = TranscriptBatch { - session: session.clone(), - messages: messages.to_vec(), - }; - self.upsert_transcript_batches_inner( - std::slice::from_ref(&batch), - parse_offset_path, - parse_offset, - expected_offset, - None, - ) - .await - } - - /// Atomically commits parsed transcript rows, the parse cursor, and the - /// Git evidence derived from the batch, in the same transaction. - #[hotpath::skip] - pub async fn persist_transcript_batch_with_git_evidence_result( - &self, - session: &SessionRecord, - messages: &[SessionMessageRecord], - parse_offset_path: &str, - expected_offset: ParseOffset, - parse_offset: ParseOffset, - git_evidence: TranscriptGitEvidence<'_>, - ) -> Result<(), TranscriptPersistenceError> { - if (!git_evidence.commit_records.is_empty() || !git_evidence.span_observations.is_empty()) - && !matches!( - &self.registered_binding().shard_id.scope, - StoreShardScopeV1::ProjectSessions { .. } - ) - { - return Err(TranscriptPersistenceError::storage( - "record transcript git evidence", - std::io::Error::new( - std::io::ErrorKind::InvalidInput, - "transcript Git evidence requires ProjectSessions authority", - ), - )); - } - let batch = TranscriptBatch { - session: session.clone(), - messages: messages.to_vec(), - }; - self.upsert_transcript_batches_inner( - std::slice::from_ref(&batch), - parse_offset_path, - parse_offset, - expected_offset, - Some(git_evidence), - ) - .await - } - #[hotpath::skip] pub async fn persist_transcript_offset_result( &self, @@ -461,99 +255,6 @@ impl SessionStoreAccess<'_, D> { .map_err(|error| TranscriptPersistenceError::storage("commit transcript batch", error)) } - #[hotpath::measure(label = "sessions.store.transcript.write_batches", future = true)] - async fn upsert_transcript_batches_inner( - &self, - batches: &[TranscriptBatch], - parse_offset_path: &str, - parse_offset: ParseOffset, - expected_offset: ParseOffset, - git_evidence: Option>, - ) -> Result<(), TranscriptPersistenceError> { - let storage_root = self - .db_path() - .parent() - .unwrap_or_else(|| std::path::Path::new(".")); - let mut payload_rollback = PayloadFileRollback::begin_cancellation_safe(storage_root); - let staged_messages = - stage_full_transcript_messages(storage_root, batches, &mut payload_rollback)?; - let mut staged_messages = staged_messages.into_iter(); - let transaction = self.begin_transcript_transaction().await?; - - let write_result: Result<(), TranscriptPersistenceError> = async { - // Full batches are one-winner compare-and-swap on the durable - // parse cursor. `actual == next_offset` is not a retry grant: - // a competing writer can share that destination while carrying - // different parse products. Post-commit publication retries - // must not re-enter this CAS with a stale expected cursor. - require_expected_offset(&transaction, parse_offset_path, expected_offset).await?; - for batch in batches { - if !Self::upsert_session_in_existing_tx(&transaction, &batch.session).await { - return Err(TranscriptPersistenceError::message( - "upsert transcript session", - "database write failed", - )); - } - for message in &batch.messages { - if tracedecay_store::codex_goal_context_correlation( - message.kind.as_deref(), - message.metadata_json.as_deref(), - ) - .is_some_and(|correlation| { - correlation.source() - == tracedecay_store::CodexGoalContextSource::ItemCompleted - }) { - reconcile_codex_goal_response(&transaction, message).await?; - } - let staged = staged_messages.next().ok_or_else(|| { - TranscriptPersistenceError::message( - "upsert LCM raw message", - "staged transcript message count did not match the write batch", - ) - })?; - Self::upsert_session_message_in_existing_tx(&transaction, message, staged) - .await?; - } - } - if staged_messages.next().is_some() { - return Err(TranscriptPersistenceError::message( - "upsert LCM raw message", - "staged transcript message count exceeded the write batch", - )); - } - if let Some(evidence) = git_evidence - && (!evidence.commit_records.is_empty() || !evidence.span_observations.is_empty()) - { - let record = |error| { - TranscriptPersistenceError::storage("record transcript git evidence", error) - }; - let mut writer = GitEvidenceWriter::open(&transaction) - .await - .map_err(record)?; - writer - .apply(GitEvidenceBatch { - observations: evidence.span_observations.to_vec(), - commits: evidence.commit_records.to_vec(), - merge_gap_secs: DEFAULT_SPAN_MERGE_GAP_SECS, - ..GitEvidenceBatch::default() - }) - .await - .map_err(record)?; - writer.finish().await.map_err(record)?; - } - set_parse_offset(&transaction, parse_offset_path, parse_offset).await?; - Ok(()) - } - .await; - - write_result?; - transaction.commit().await.map_err(|error| { - TranscriptPersistenceError::storage("commit transcript batch", error) - })?; - payload_rollback.disarm(); - Ok(()) - } - #[hotpath::skip] pub async fn get_parse_offset( &self, @@ -701,12 +402,7 @@ async fn require_expected_pair_offset( #[cfg(test)] mod tests { - use tracedecay_store::{SessionMessageRecord, SessionRecord}; - - use super::{ - PayloadFileRollback, TranscriptBatch, TranscriptPersistenceError, decode_u64_bits_value, - encode_u64_bits, stage_full_transcript_messages, - }; + use super::{decode_u64_bits_value, encode_u64_bits}; /// Every `parse_offsets` column round-trips the whole `u64` domain: the /// Codex corpus epoch stores a 128-bit digest across `byte_offset` and @@ -727,66 +423,4 @@ mod tests { "the upper half maps onto the negative INTEGER range instead of failing" ); } - - #[test] - fn full_transcript_staging_preserves_sanitization_failure_attribution() { - let temp = tempfile::tempdir().unwrap(); - let storage_root = temp.path().join("store"); - std::fs::create_dir(&storage_root).unwrap(); - let batch = TranscriptBatch { - session: SessionRecord { - provider: "claude".to_owned(), - session_id: "session-1".to_owned(), - project_key: "/tmp/project".to_owned(), - project_path: "/tmp/project".to_owned(), - title: None, - started_at: None, - ended_at: None, - transcript_path: None, - metadata_json: None, - parent_session_id: None, - is_subagent: false, - agent_id: None, - parent_tool_use_id: None, - }, - messages: vec![SessionMessageRecord { - provider: "claude".to_owned(), - message_id: "message-1".to_owned(), - session_id: "session-1".to_owned(), - role: "assistant".to_owned(), - timestamp: Some(1), - ordinal: 1, - text: "ordinary content".to_owned(), - kind: None, - model: None, - tool_names: None, - source_path: None, - source_offset: None, - metadata_json: Some("[]".to_owned()), - }], - }; - let mut rollback = PayloadFileRollback::begin_cancellation_safe(&storage_root); - - let error = match stage_full_transcript_messages( - &storage_root, - std::slice::from_ref(&batch), - &mut rollback, - ) { - Ok(_) => panic!("non-object metadata must be refused before transaction acquisition"), - Err(error) => error, - }; - - match error { - TranscriptPersistenceError::Storage { operation, source } => { - assert_eq!(operation, "upsert LCM raw message"); - assert!( - source - .to_string() - .contains("LCM metadata sanitization failed"), - "sanitization cause was lost: {source}" - ); - } - other => panic!("expected attributed sanitization failure, got {other}"), - } - } } diff --git a/crates/tracedecay-sessions/src/runtime/store_access/types.rs b/crates/tracedecay-sessions/src/runtime/store_access/types.rs index d480e8e012..f4cce1ebac 100644 --- a/crates/tracedecay-sessions/src/runtime/store_access/types.rs +++ b/crates/tracedecay-sessions/src/runtime/store_access/types.rs @@ -1,6 +1,6 @@ use std::error::Error; -use tracedecay_store::{ParseOffset, SessionMessageRecord, SessionRecord}; +use tracedecay_store::ParseOffset; use crate::runtime::host_coverage::HostCoverageReason; @@ -60,17 +60,6 @@ pub struct SessionIngestHealth { pub last_ingest_unix: Option, } -/// One transcript session plus its parsed messages, for projection-only -/// multi-session upserts from stores such as Hermes `state.db`. -/// -/// This compatibility DTO remains local because projection-only persistence is -/// intentionally outside the authoritative transcript store contract. -#[derive(Debug, Clone)] -pub struct TranscriptBatch { - pub session: SessionRecord, - pub messages: Vec, -} - #[derive(Debug)] pub enum TranscriptPersistenceError { Conflict { diff --git a/crates/tracedecay-store/src/lib.rs b/crates/tracedecay-store/src/lib.rs index 16836e65da..a9218caa38 100644 --- a/crates/tracedecay-store/src/lib.rs +++ b/crates/tracedecay-store/src/lib.rs @@ -256,5 +256,5 @@ pub use session::{ }; pub use transcript::{ ParseOffset, SessionMessageRecord, SessionRecord, TranscriptStore, TranscriptStoreError, - TranscriptStoreResult, TranscriptWriteBatch, TranscriptWriteKind, + TranscriptStoreResult, TranscriptWriteBatch, }; diff --git a/crates/tracedecay-store/src/transcript.rs b/crates/tracedecay-store/src/transcript.rs index 58446a8656..105d8a0239 100644 --- a/crates/tracedecay-store/src/transcript.rs +++ b/crates/tracedecay-store/src/transcript.rs @@ -48,28 +48,13 @@ pub struct ParseOffset { pub file_id: u64, } -/// Validated authoritative transcript persistence request. +/// Validated cursor advance for parsed transcript input that emitted no +/// messages. #[derive(Debug, Clone, PartialEq, Eq)] pub struct TranscriptWriteBatch { cursor_path: PathBuf, - kind: TranscriptWriteKind, -} - -/// Consumed representation of a validated transcript write. -#[derive(Debug, Clone, PartialEq, Eq)] -pub enum TranscriptWriteKind { - /// Advances the cursor after parsing input that emitted no messages. - AdvanceOffset { - expected_offset: ParseOffset, - next_offset: ParseOffset, - }, - /// Atomically persists a session, its messages, and the next cursor. - Upsert { - session: Box, - messages: Vec, - expected_offset: ParseOffset, - next_offset: ParseOffset, - }, + expected_offset: ParseOffset, + next_offset: ParseOffset, } impl TranscriptWriteBatch { @@ -85,99 +70,15 @@ impl TranscriptWriteBatch { Ok(Self { cursor_path, - kind: TranscriptWriteKind::AdvanceOffset { - expected_offset, - next_offset, - }, + expected_offset, + next_offset, }) } - /// Builds a full atomic session/message/offset write. - pub fn upsert( - session: SessionRecord, - messages: Vec, - expected_offset: ParseOffset, - next_offset: ParseOffset, - ) -> TranscriptStoreResult { - let cursor_path = session - .transcript_path - .as_deref() - .map(PathBuf::from) - .ok_or_else(|| TranscriptStoreError::MissingTranscriptPath { - provider: session.provider.clone(), - session_id: session.session_id.clone(), - })?; - Self::upsert_with_cursor(cursor_path, session, messages, expected_offset, next_offset) - } - - /// Builds a full atomic write whose durable cursor key differs from the - /// session's physical transcript path. - /// - /// Virtual transcript sources use a stable logical cursor while retaining - /// the real source path in [`SessionRecord::transcript_path`]. - pub fn upsert_with_cursor( - cursor_path: PathBuf, - session: SessionRecord, - messages: Vec, - expected_offset: ParseOffset, - next_offset: ParseOffset, - ) -> TranscriptStoreResult { - let session_path = session.transcript_path.as_deref().ok_or_else(|| { - TranscriptStoreError::MissingTranscriptPath { - provider: session.provider.clone(), - session_id: session.session_id.clone(), - } - })?; - if session_path.is_empty() { - return Err(TranscriptStoreError::InvalidTranscriptPath); - } - if cursor_path.as_os_str().is_empty() { - return Err(TranscriptStoreError::InvalidCursorPath); - } - - if let Some(message) = messages.iter().find(|message| { - message.provider != session.provider || message.session_id != session.session_id - }) { - return Err(TranscriptStoreError::MessageIdentityMismatch { - message_id: message.message_id.clone(), - expected_provider: session.provider, - actual_provider: message.provider.clone(), - expected_session_id: session.session_id, - actual_session_id: message.session_id.clone(), - }); - } - - Ok(Self { - cursor_path, - kind: TranscriptWriteKind::Upsert { - session: Box::new(session), - messages, - expected_offset, - next_offset, - }, - }) - } - - /// Returns the durable cursor identity represented by this write. - pub fn cursor_path(&self) -> &Path { - &self.cursor_path - } - - /// Returns the durable cursor that the writer observed before parsing. - pub fn expected_offset(&self) -> ParseOffset { - match &self.kind { - TranscriptWriteKind::AdvanceOffset { - expected_offset, .. - } - | TranscriptWriteKind::Upsert { - expected_offset, .. - } => *expected_offset, - } - } - - /// Consumes this validated request for persistence. - pub fn into_parts(self) -> (PathBuf, TranscriptWriteKind) { - (self.cursor_path, self.kind) + /// Consumes this validated request into its cursor path, the durable + /// cursor the writer observed before parsing, and the cursor to persist. + pub fn into_parts(self) -> (PathBuf, ParseOffset, ParseOffset) { + (self.cursor_path, self.expected_offset, self.next_offset) } } @@ -186,23 +87,6 @@ impl TranscriptWriteBatch { pub enum TranscriptStoreError { #[error("transcript cursor path must not be empty")] InvalidCursorPath, - #[error("transcript path must not be empty")] - InvalidTranscriptPath, - #[error("session {provider}/{session_id} has no transcript path")] - MissingTranscriptPath { - provider: String, - session_id: String, - }, - #[error( - "message {message_id} identity {actual_provider}/{actual_session_id} does not match session {expected_provider}/{expected_session_id}" - )] - MessageIdentityMismatch { - message_id: String, - expected_provider: String, - actual_provider: String, - expected_session_id: String, - actual_session_id: String, - }, #[error( "transcript cursor conflict for {cursor_path:?}: expected {expected:?}, found {actual:?}" )] @@ -233,7 +117,7 @@ pub trait TranscriptStore: Send + Sync { cursor_path: &Path, ) -> impl Future> + Send; - /// Persists one offset-only or full atomic write in the authoritative store. + /// Persists one cursor advance in the authoritative store. fn persist_transcript_batch( &self, batch: TranscriptWriteBatch, @@ -244,73 +128,6 @@ pub trait TranscriptStore: Send + Sync { mod tests { use super::*; - fn session(transcript_path: Option<&str>) -> SessionRecord { - SessionRecord { - provider: "test".into(), - session_id: "session".into(), - project_key: "project".into(), - project_path: "/project".into(), - title: None, - started_at: None, - ended_at: None, - transcript_path: transcript_path.map(str::to_owned), - metadata_json: None, - parent_session_id: None, - is_subagent: false, - agent_id: None, - parent_tool_use_id: None, - } - } - - fn message(provider: &str, session_id: &str) -> SessionMessageRecord { - SessionMessageRecord { - provider: provider.into(), - message_id: "message".into(), - session_id: session_id.into(), - role: "user".into(), - timestamp: None, - ordinal: 0, - text: "hello".into(), - kind: None, - model: None, - tool_names: None, - source_path: None, - source_offset: None, - metadata_json: None, - } - } - - #[test] - fn upsert_with_cursor_rejects_an_empty_cursor_path() { - let batch = TranscriptWriteBatch::upsert_with_cursor( - PathBuf::new(), - session(Some("/physical/store.db")), - Vec::new(), - ParseOffset::default(), - ParseOffset::default(), - ); - - assert!(matches!( - batch, - Err(TranscriptStoreError::InvalidCursorPath) - )); - } - - #[test] - fn upsert_requires_a_session_transcript_path() { - let batch = TranscriptWriteBatch::upsert( - session(None), - Vec::new(), - ParseOffset::default(), - ParseOffset::default(), - ); - - assert!(matches!( - batch, - Err(TranscriptStoreError::MissingTranscriptPath { .. }) - )); - } - #[test] fn advance_offset_rejects_an_empty_cursor_path() { let batch = TranscriptWriteBatch::advance_offset( @@ -324,49 +141,4 @@ mod tests { Err(TranscriptStoreError::InvalidCursorPath) )); } - - #[test] - fn upsert_rejects_an_empty_transcript_path() { - let batch = TranscriptWriteBatch::upsert( - session(Some("")), - Vec::new(), - ParseOffset::default(), - ParseOffset::default(), - ); - - assert!(matches!( - batch, - Err(TranscriptStoreError::InvalidTranscriptPath) - )); - } - - #[test] - fn upsert_rejects_a_foreign_message_provider() { - let batch = TranscriptWriteBatch::upsert( - session(Some("session.jsonl")), - vec![message("other", "session")], - ParseOffset::default(), - ParseOffset::default(), - ); - - assert!(matches!( - batch, - Err(TranscriptStoreError::MessageIdentityMismatch { .. }) - )); - } - - #[test] - fn upsert_rejects_a_foreign_message_session() { - let batch = TranscriptWriteBatch::upsert( - session(Some("session.jsonl")), - vec![message("test", "other")], - ParseOffset::default(), - ParseOffset::default(), - ); - - assert!(matches!( - batch, - Err(TranscriptStoreError::MessageIdentityMismatch { .. }) - )); - } } diff --git a/crates/tracedecay/src/daemon/scheduler.rs b/crates/tracedecay/src/daemon/scheduler.rs index 819230c229..963eb270a5 100644 --- a/crates/tracedecay/src/daemon/scheduler.rs +++ b/crates/tracedecay/src/daemon/scheduler.rs @@ -1687,17 +1687,10 @@ mod global_retention_tests { source_offset: None, metadata_json: None, }; - assert!( - database - .upsert_transcript_batch( - &session, - std::slice::from_ref(&message), - "global-retention-fixture", - tracedecay_global_db::ParseOffset::default(), - ) - .await, - "project the registered retention fixture message" - ); + database + .lcm_ingest_raw_message(database.db_path().parent().unwrap(), &message) + .await + .expect("seed the registered retention fixture message"); let transaction = database .begin_write_transaction() diff --git a/crates/tracedecay/src/mcp/server/host_admission_tests.rs b/crates/tracedecay/src/mcp/server/host_admission_tests.rs index 61e453ddf0..dd60b02d0c 100644 --- a/crates/tracedecay/src/mcp/server/host_admission_tests.rs +++ b/crates/tracedecay/src/mcp/server/host_admission_tests.rs @@ -1490,7 +1490,7 @@ async fn credential_canary_receipt_analytics_and_git_span_survive_database_reope .expect("seed protected session") ); test_runtime - .upsert_transcript_batch_for_test( + .seed_session_messages_for_test( HostAdmissionScope::Project, &session, std::slice::from_ref(&SessionMessageRecord { @@ -1508,8 +1508,6 @@ async fn credential_canary_receipt_analytics_and_git_span_survive_database_reope source_offset: None, metadata_json: None, }), - &format!("host-admission-test-message:hermes:{protected}"), - tracedecay_global_db::ParseOffset::default(), ) .await .expect("seed protected transcript"); diff --git a/crates/tracedecay/tests/automation_runner_test/support/fixtures.rs b/crates/tracedecay/tests/automation_runner_test/support/fixtures.rs index b5aed1127c..f30e5d65db 100644 --- a/crates/tracedecay/tests/automation_runner_test/support/fixtures.rs +++ b/crates/tracedecay/tests/automation_runner_test/support/fixtures.rs @@ -9,7 +9,6 @@ use tracedecay_automation_runtime::automation::run_ledger::{ AutomationRunLedgerRecord, read_run_artifact_payload, }; use tracedecay_domain::FactOwnerV1; -use tracedecay_global_db::ParseOffset; use tracedecay_project::project::{TraceDecay, TraceDecayOpenOptions}; use tracedecay_project::test_support::host_admission::HostAdmissionTestRuntimeV1; use tracedecay_runtime_core::tracedecay::current_timestamp; @@ -119,17 +118,10 @@ pub(crate) async fn seed_project_session_activity_at(cg: &TraceDecay, timestamp: source_offset: None, metadata_json: None, }; - assert!( - sessions - .upsert_transcript_batch( - &session, - std::slice::from_ref(&message), - &format!("combined-review-activity:{timestamp}"), - ParseOffset::default(), - ) - .await, - "activity fixture must persist a timestamped message" - ); + sessions + .lcm_ingest_raw_message(sessions.db_path().parent().unwrap(), &message) + .await + .expect("activity fixture must persist a timestamped message"); } #[cfg(feature = "test-transport")] @@ -197,11 +189,9 @@ pub(crate) async fn seed_search_underuse_session_evidence(cg: &TraceDecay) { source_offset: None, metadata_json: Some(json!({ "cmd": "rg automation src" }).to_string()), }; - assert!( - db.upsert_session_message_for_test(HostAdmissionScope::Project, &message) - .await - .unwrap() - ); + db.upsert_session_message_for_test(HostAdmissionScope::Project, &message) + .await + .unwrap(); } /// Seeds one session message at `timestamp` so the scheduler observes LCM @@ -279,11 +269,9 @@ pub(crate) async fn seed_session_message_in_db( .source .map(|source| json!({ "source": source }).to_string()), }; - assert!( - db.upsert_session_message_for_test(HostAdmissionScope::Project, &message) - .await - .unwrap() - ); + db.upsert_session_message_for_test(HostAdmissionScope::Project, &message) + .await + .unwrap(); } #[derive(Debug, Clone)] diff --git a/crates/tracedecay/tests/common/mod.rs b/crates/tracedecay/tests/common/mod.rs index 9494ab83ca..cc7628dbb8 100644 --- a/crates/tracedecay/tests/common/mod.rs +++ b/crates/tracedecay/tests/common/mod.rs @@ -1438,7 +1438,7 @@ impl LcmTestRuntime { self.runtime .upsert_session_message_for_test(HostAdmissionScope::Profile, message) .await - .unwrap_or(false) + .is_ok() } pub async fn lcm_load_raw_message( diff --git a/crates/tracedecay/tests/dashboard_api_test/analytics.rs b/crates/tracedecay/tests/dashboard_api_test/analytics.rs index 604ea5c967..9b6687922a 100644 --- a/crates/tracedecay/tests/dashboard_api_test/analytics.rs +++ b/crates/tracedecay/tests/dashboard_api_test/analytics.rs @@ -199,12 +199,10 @@ async fn seed_session_store(runtime: &DashboardTestRuntimeV1, project: &Path) { ]; for row in rows { - assert!( - runtime - .upsert_session_message_for_test(HostAdmissionScope::Project, &row) - .await - .expect("seed analytics session message") - ); + runtime + .upsert_session_message_for_test(HostAdmissionScope::Project, &row) + .await + .expect("seed analytics session message"); } } diff --git a/crates/tracedecay/tests/dashboard_api_test/delivery.rs b/crates/tracedecay/tests/dashboard_api_test/delivery.rs index 3bce463bbf..a01b6fdcc5 100644 --- a/crates/tracedecay/tests/dashboard_api_test/delivery.rs +++ b/crates/tracedecay/tests/dashboard_api_test/delivery.rs @@ -41,7 +41,6 @@ use tracedecay_domain::{ ObservationSourceIdentityV1, ProjectId, ProviderId, RefId, RepositoryId, SessionId, SourceSpan, SymbolOccurrenceId, UtcMicros, WorktreeId, }; -use tracedecay_global_db::ParseOffset; use tracedecay_mcp::handlers::dashboard_delivery::DashboardDeliveryReadAdapter; use tracedecay_tool_catalog::{CapabilityId, UseCaseId}; @@ -783,13 +782,7 @@ fn delivery_overview_counts_agent_tool_calls_for_sessions_on_the_live_branch() { .collect(); fixture .host_runtime - .upsert_transcript_batch_for_test( - HostAdmissionScope::Project, - session, - &messages, - &format!("agent-usage-fixture:{}", session.session_id), - ParseOffset::default(), - ) + .seed_session_messages_for_test(HostAdmissionScope::Project, session, &messages) .await .expect("seed agent usage transcript"); } diff --git a/crates/tracedecay/tests/dashboard_api_test/loom.rs b/crates/tracedecay/tests/dashboard_api_test/loom.rs index 7f8ac7872a..e1c87a578a 100644 --- a/crates/tracedecay/tests/dashboard_api_test/loom.rs +++ b/crates/tracedecay/tests/dashboard_api_test/loom.rs @@ -7,7 +7,6 @@ use serde_json::json; use tracedecay_dashboard_api::{ DashboardGitCorrelationReadFutureV1, DashboardGitCorrelationReadPortV1, }; -use tracedecay_global_db::ParseOffset; use tracedecay_sessions::runtime::git_correlation::{ DEFAULT_SPAN_MERGE_GAP_SECS, SpanObservation, SpanSource, }; @@ -298,13 +297,11 @@ fn loom_temporal_serves_recorded_tool_and_pull_request_events_in_recorded_time_o for (record, messages) in batches { fixture .host_runtime - .upsert_transcript_batch_for_test( - HostAdmissionScope::Project, - record, - &messages, - &format!("loom-events:{}", record.session_id), - ParseOffset::default(), - ) + .seed_session_messages_for_test( + HostAdmissionScope::Project, + record, + &messages, + ) .await .unwrap_or_else(|error| panic!("seed {}: {error}", record.session_id)); } @@ -376,7 +373,7 @@ fn loom_temporal_serves_an_empty_event_stream_as_complete_zero_coverage() { }; fixture .host_runtime - .upsert_transcript_batch_for_test( + .seed_session_messages_for_test( HostAdmissionScope::Project, &quiet, &[loom_message( @@ -388,8 +385,6 @@ fn loom_temporal_serves_an_empty_event_stream_as_complete_zero_coverage() { None, "hello", )], - "loom-events:quiet", - ParseOffset::default(), ) .await .unwrap_or_else(|error| panic!("seed quiet session: {error}")); @@ -632,12 +627,10 @@ fn loom_temporal_serves_one_bounded_page_of_a_large_history() { large_history_session(&fixture.host_runtime, &fixture.project_root, index); fixture .host_runtime - .upsert_transcript_batch_for_test( + .seed_session_messages_for_test( HostAdmissionScope::Project, &session, &large_history_messages(&session), - &format!("large-history:{}", session.session_id), - ParseOffset::default(), ) .await .unwrap_or_else(|error| panic!("seed {}: {error}", session.session_id)); diff --git a/crates/tracedecay/tests/dashboard_api_test/runtime.rs b/crates/tracedecay/tests/dashboard_api_test/runtime.rs index 443d3c1181..feb1c5cd12 100644 --- a/crates/tracedecay/tests/dashboard_api_test/runtime.rs +++ b/crates/tracedecay/tests/dashboard_api_test/runtime.rs @@ -470,13 +470,14 @@ impl DashboardTestRuntimeV1 { Ok(self.database(scope)?.upsert_session(session).await) } + /// Seeds one raw LCM message into its already-registered session. pub(crate) async fn upsert_session_message_for_test( &self, scope: HostAdmissionScope, message: &SessionMessageRecord, - ) -> Result { - let database = self.database(scope)?; - let session = database + ) -> Result<()> { + let session = self + .database(scope)? .get_session(&message.provider, &message.session_id) .await .map_err(|error| TraceDecayError::Database { @@ -490,37 +491,43 @@ impl DashboardTestRuntimeV1 { message.provider, message.session_id ), })?; - Ok(database - .upsert_transcript_batch( - &session, - std::slice::from_ref(message), - &format!( - "dashboard-test-message:{}:{}", - message.provider, message.message_id - ), - tracedecay_global_db::ParseOffset::default(), - ) - .await) + self.seed_session_messages_for_test(scope, &session, std::slice::from_ref(message)) + .await + .map(|_| ()) } - pub(crate) async fn upsert_transcript_batch_for_test( + /// Seeds one session row and its raw LCM messages, returning each + /// message's raw store id in input order. + pub(crate) async fn seed_session_messages_for_test( &self, scope: HostAdmissionScope, session: &SessionRecord, messages: &[SessionMessageRecord], - source: &str, - offset: tracedecay_global_db::ParseOffset, ) -> Result> { let database = self.database(scope)?; - if !database - .upsert_transcript_batch(session, messages, source, offset) - .await - { + if !database.upsert_session(session).await { return Err(TraceDecayError::Database { - operation: "seed dashboard test transcript batch".to_owned(), - message: "registered transcript batch write failed".to_owned(), + operation: "seed dashboard test session".to_owned(), + message: "registered session write failed".to_owned(), }); } + let storage_root = + database + .db_path() + .parent() + .ok_or_else(|| TraceDecayError::Database { + operation: "seed dashboard test session message".to_owned(), + message: "registered session database has no storage root".to_owned(), + })?; + for message in messages { + database + .lcm_ingest_raw_message(storage_root, message) + .await + .map_err(|error| TraceDecayError::Database { + operation: "seed dashboard test session message".to_owned(), + message: error.to_string(), + })?; + } let mut store_ids = Vec::with_capacity(messages.len()); for message in messages { let raw = load_registered_raw_message(database, &message.provider, &message.message_id) diff --git a/crates/tracedecay/tests/dashboard_api_test/savings.rs b/crates/tracedecay/tests/dashboard_api_test/savings.rs index 219766a445..e7544be6a5 100644 --- a/crates/tracedecay/tests/dashboard_api_test/savings.rs +++ b/crates/tracedecay/tests/dashboard_api_test/savings.rs @@ -18,7 +18,6 @@ use serde_json::Value; use std::sync::Arc; use tempfile::TempDir; use tracedecay::dashboard; -use tracedecay_global_db::ParseOffset; use tracedecay_runtime_core::config::ProfileRoot; use tracedecay_sessions::admission::HostAdmissionScope; use tracedecay_sessions::runtime::SessionRecord; @@ -101,30 +100,22 @@ impl SavingsSeed<'_> { async fn upsert_session_message( &self, message: &tracedecay_sessions::runtime::SessionMessageRecord, - ) -> bool { + ) { self.0 .upsert_session_message_for_test(HostAdmissionScope::Project, message) .await .expect("seed savings session message") } - async fn upsert_transcript_batch( + async fn seed_session_messages( &self, session: &SessionRecord, messages: &[tracedecay_sessions::runtime::SessionMessageRecord], - source: &str, - offset: ParseOffset, - ) -> bool { + ) { self.0 - .upsert_transcript_batch_for_test( - HostAdmissionScope::Project, - session, - messages, - source, - offset, - ) + .seed_session_messages_for_test(HostAdmissionScope::Project, session, messages) .await - .is_ok() + .expect("seed savings session messages"); } } @@ -163,8 +154,7 @@ async fn seed_global_db(runtime: &DashboardTestRuntimeV1, project: &Path, day_st )) .await ); - assert!( - gdb.upsert_session_message(&message( + gdb.upsert_session_message(&message( "m-usage-1", "sess-usage", "assistant", @@ -178,8 +168,7 @@ async fn seed_global_db(runtime: &DashboardTestRuntimeV1, project: &Path, day_st ), }, )) - .await - ); + .await; // S2: no usage anywhere → estimated (chars/4, user→input, assistant→output). assert!( @@ -191,36 +180,32 @@ async fn seed_global_db(runtime: &DashboardTestRuntimeV1, project: &Path, day_st )) .await ); - assert!( - gdb.upsert_session_message(&message( - "m-est-1", - "sess-estimated", - "user", - 1, - TEXT_USER, - MessageDetails { - timestamp: day_start + 210, - model: Some("gpt-5.5-high"), - metadata_json: None, - }, - )) - .await - ); - assert!( - gdb.upsert_session_message(&message( - "m-est-2", - "sess-estimated", - "assistant", - 2, - TEXT_ASSISTANT, - MessageDetails { - timestamp: day_start + 220, - model: Some("gpt-5.5-high"), - metadata_json: None, - }, - )) - .await - ); + gdb.upsert_session_message(&message( + "m-est-1", + "sess-estimated", + "user", + 1, + TEXT_USER, + MessageDetails { + timestamp: day_start + 210, + model: Some("gpt-5.5-high"), + metadata_json: None, + }, + )) + .await; + gdb.upsert_session_message(&message( + "m-est-2", + "sess-estimated", + "assistant", + 2, + TEXT_ASSISTANT, + MessageDetails { + timestamp: day_start + 220, + model: Some("gpt-5.5-high"), + metadata_json: None, + }, + )) + .await; // S3: no model id recorded at all → "unknown model" row, never priced. assert!( @@ -232,21 +217,19 @@ async fn seed_global_db(runtime: &DashboardTestRuntimeV1, project: &Path, day_st )) .await ); - assert!( - gdb.upsert_session_message(&message( - "m-unknown-1", - "sess-unknown", - "assistant", - 1, - TEXT_UNKNOWN, - MessageDetails { - timestamp: day_start + 310, - model: None, - metadata_json: None, - }, - )) - .await - ); + gdb.upsert_session_message(&message( + "m-unknown-1", + "sess-unknown", + "assistant", + 1, + TEXT_UNKNOWN, + MessageDetails { + timestamp: day_start + 310, + model: None, + metadata_json: None, + }, + )) + .await; // S4: usage (OpenAI field names) + a usage-less message → mixed. assert!( @@ -258,36 +241,32 @@ async fn seed_global_db(runtime: &DashboardTestRuntimeV1, project: &Path, day_st )) .await ); - assert!( - gdb.upsert_session_message(&message( - "m-mixed-1", - "sess-mixed", - "assistant", - 1, - TEXT_ASSISTANT, - MessageDetails { - timestamp: day_start + 410, - model: Some("claude-opus-4-8-thinking-max"), - metadata_json: Some(r#"{"usage":{"prompt_tokens":500,"completion_tokens":700}}"#,), - }, - )) - .await - ); - assert!( - gdb.upsert_session_message(&message( - "m-mixed-2", - "sess-mixed", - "assistant", - 2, - TEXT_MIXED, - MessageDetails { - timestamp: day_start + 420, - model: Some("claude-opus-4-8-thinking-max"), - metadata_json: None, - }, - )) - .await - ); + gdb.upsert_session_message(&message( + "m-mixed-1", + "sess-mixed", + "assistant", + 1, + TEXT_ASSISTANT, + MessageDetails { + timestamp: day_start + 410, + model: Some("claude-opus-4-8-thinking-max"), + metadata_json: Some(r#"{"usage":{"prompt_tokens":500,"completion_tokens":700}}"#), + }, + )) + .await; + gdb.upsert_session_message(&message( + "m-mixed-2", + "sess-mixed", + "assistant", + 2, + TEXT_MIXED, + MessageDetails { + timestamp: day_start + 420, + model: Some("claude-opus-4-8-thinking-max"), + metadata_json: None, + }, + )) + .await; // S5: transcript metadata carries the shape Codex backfill writes. It is // still content metadata, not provider-usage billing authority. @@ -300,8 +279,7 @@ async fn seed_global_db(runtime: &DashboardTestRuntimeV1, project: &Path, day_st )) .await ); - assert!( - gdb.upsert_session_message(&message( + gdb.upsert_session_message(&message( "m-codex-1", "sess-codex", "assistant", @@ -315,25 +293,22 @@ async fn seed_global_db(runtime: &DashboardTestRuntimeV1, project: &Path, day_st ), }, )) - .await - ); - assert!( - gdb.upsert_session_message( - &MessageRecordBuilder::new( - "cursor", - "m-codex-summary", - "sess-codex", - "assistant", - 2, - "Synthetic Codex compaction placeholder that is not real model output.", - "summary", - ) - .with_timestamp(Some(day_start + 520)) - .with_model(Some("gpt-5.3-codex-high")) - .build() + .await; + gdb.upsert_session_message( + &MessageRecordBuilder::new( + "cursor", + "m-codex-summary", + "sess-codex", + "assistant", + 2, + "Synthetic Codex compaction placeholder that is not real model output.", + "summary", ) - .await - ); + .with_timestamp(Some(day_start + 520)) + .with_model(Some("gpt-5.3-codex-high")) + .build(), + ) + .await; } async fn seed_daily_limit_regression( @@ -379,15 +354,7 @@ async fn seed_daily_limit_regression( }, )); - assert!( - gdb.upsert_transcript_batch( - &daily_session, - &messages, - "daily-limit-regression.jsonl", - ParseOffset::default(), - ) - .await - ); + gdb.seed_session_messages(&daily_session, &messages).await; } async fn start_fixture(seed: FixtureSeed) -> Fixture { diff --git a/crates/tracedecay/tests/hooks_lsp_suite/hint_settlement_test.rs b/crates/tracedecay/tests/hooks_lsp_suite/hint_settlement_test.rs index 69ea6a066e..da350b8e45 100644 --- a/crates/tracedecay/tests/hooks_lsp_suite/hint_settlement_test.rs +++ b/crates/tracedecay/tests/hooks_lsp_suite/hint_settlement_test.rs @@ -120,8 +120,7 @@ impl ProjectSettlementFixture { let Some(tools) = tools else { return; }; - let inserted = self - .runtime + self.runtime .upsert_session_message_for_test( HostAdmissionScope::Project, &SessionMessageRecord { @@ -142,7 +141,6 @@ impl ProjectSettlementFixture { ) .await .expect("upsert project session message"); - assert!(inserted, "session message should upsert"); } async fn settle(&self, now_secs: i64) -> HintOutcomeStats { diff --git a/crates/tracedecay/tests/mcp_suite/git_correlation_test.rs b/crates/tracedecay/tests/mcp_suite/git_correlation_test.rs index 61df6b69ac..1bfad95b49 100644 --- a/crates/tracedecay/tests/mcp_suite/git_correlation_test.rs +++ b/crates/tracedecay/tests/mcp_suite/git_correlation_test.rs @@ -277,15 +277,13 @@ async fn sessions_for_distinguishes_empty_correlation_index_from_no_match() { .await .unwrap_or_else(|e| panic!("seed session: {e}")) ); - assert!( - runtime - .upsert_session_message_for_test( - HostAdmissionScope::Project, - &message("s1", "s1-m1", 1_050, "work on main"), - ) - .await - .unwrap_or_else(|e| panic!("seed session message: {e}")) - ); + runtime + .upsert_session_message_for_test( + HostAdmissionScope::Project, + &message("s1", "s1-m1", 1_050, "work on main"), + ) + .await + .unwrap_or_else(|e| panic!("seed session message: {e}")); let server = McpServer::new_with_host_admission_test_runtime_for_test( cg, None, diff --git a/crates/tracedecay/tests/mcp_suite/support.rs b/crates/tracedecay/tests/mcp_suite/support.rs index 7cfaf4750f..ea68b38566 100644 --- a/crates/tracedecay/tests/mcp_suite/support.rs +++ b/crates/tracedecay/tests/mcp_suite/support.rs @@ -1461,29 +1461,27 @@ async fn seed_lcm_message_with_role( .await .unwrap() ); - assert!( - runtime - .upsert_session_message_for_test( - HostAdmissionScope::Project, - &SessionMessageRecord { - provider: provider.to_string(), - message_id: message_id.to_string(), - session_id: session_id.to_string(), - role: role.to_string(), - timestamp: Some(ordinal + 1), - ordinal, - text: text.into(), - kind: Some(kind.to_string()), - model: Some("test-model".to_string()), - tool_names: None, - source_path: Some(format!("{session_id}.jsonl")), - source_offset: Some(0), - metadata_json: None, - }, - ) - .await - .unwrap() - ); + runtime + .upsert_session_message_for_test( + HostAdmissionScope::Project, + &SessionMessageRecord { + provider: provider.to_string(), + message_id: message_id.to_string(), + session_id: session_id.to_string(), + role: role.to_string(), + timestamp: Some(ordinal + 1), + ordinal, + text: text.into(), + kind: Some(kind.to_string()), + model: Some("test-model".to_string()), + tool_names: None, + source_path: Some(format!("{session_id}.jsonl")), + source_offset: Some(0), + metadata_json: None, + }, + ) + .await + .unwrap(); } #[cfg(feature = "test-transport")] @@ -1629,13 +1627,10 @@ pub(crate) async fn seed_temporal_lcm_tool_result_message( .await .unwrap() .expect("canonical tool result must project to the compatibility store"); - assert!( - runtime - .upsert_session_message_for_test(HostAdmissionScope::Project, &projected) - .await - .unwrap(), - "canonical compatibility output must apply the bounded payload policy" - ); + runtime + .upsert_session_message_for_test(HostAdmissionScope::Project, &projected) + .await + .expect("canonical compatibility output must apply the bounded payload policy"); projection } diff --git a/crates/tracedecay/tests/session_suite/git_backfill.rs b/crates/tracedecay/tests/session_suite/git_backfill.rs index b606de7a23..3cb6c53fa7 100644 --- a/crates/tracedecay/tests/session_suite/git_backfill.rs +++ b/crates/tracedecay/tests/session_suite/git_backfill.rs @@ -202,22 +202,18 @@ async fn open_seeded_db(repo: &Path) -> (TempDir, HostAdmissionTestRuntimeV1, St .await .unwrap() ); - assert!( - db.upsert_session_message_for_test( - HostAdmissionScope::Project, - &message("s_switch", "m1", T_BASE + 50), - ) - .await - .unwrap() - ); - assert!( - db.upsert_session_message_for_test( - HostAdmissionScope::Project, - &message("s_switch", "m2", T_BASE + 850), - ) - .await - .unwrap() - ); + db.upsert_session_message_for_test( + HostAdmissionScope::Project, + &message("s_switch", "m1", T_BASE + 50), + ) + .await + .unwrap(); + db.upsert_session_message_for_test( + HostAdmissionScope::Project, + &message("s_switch", "m2", T_BASE + 850), + ) + .await + .unwrap(); // s_main only overlaps the first main stretch. assert!( @@ -228,14 +224,12 @@ async fn open_seeded_db(repo: &Path) -> (TempDir, HostAdmissionTestRuntimeV1, St .await .unwrap() ); - assert!( - db.upsert_session_message_for_test( - HostAdmissionScope::Project, - &message("s_main", "m3", T_BASE + 200), - ) - .await - .unwrap() - ); + db.upsert_session_message_for_test( + HostAdmissionScope::Project, + &message("s_main", "m3", T_BASE + 200), + ) + .await + .unwrap(); (tmp, db, project) } @@ -740,14 +734,12 @@ async fn backfill_skips_non_worktree_sessions() { .await .unwrap() ); - assert!( - db.upsert_session_message_for_test( - HostAdmissionScope::Project, - &message("s_orphan", "m1", T_BASE + 50), - ) - .await - .unwrap() - ); + db.upsert_session_message_for_test( + HostAdmissionScope::Project, + &message("s_orphan", "m1", T_BASE + 50), + ) + .await + .unwrap(); let git = FakeGit { timeline: vec![], diff --git a/crates/tracedecay/tests/session_suite/global_db.rs b/crates/tracedecay/tests/session_suite/global_db.rs index d5f5b6dfae..cbe18e5f28 100644 --- a/crates/tracedecay/tests/session_suite/global_db.rs +++ b/crates/tracedecay/tests/session_suite/global_db.rs @@ -157,7 +157,7 @@ impl RegisteredSessionTestExt for HostAdmissionTestRuntimeV1 { async fn upsert_session_message(&self, message: &SessionMessageRecord) -> bool { self.upsert_session_message_for_test(HostAdmissionScope::Profile, message) .await - .unwrap_or(false) + .is_ok() } async fn lcm_load_raw_message( diff --git a/crates/tracedecay/tests/session_suite/lcm_compression/mod.rs b/crates/tracedecay/tests/session_suite/lcm_compression/mod.rs index c14bdf56be..60efcdd2f8 100644 --- a/crates/tracedecay/tests/session_suite/lcm_compression/mod.rs +++ b/crates/tracedecay/tests/session_suite/lcm_compression/mod.rs @@ -109,13 +109,7 @@ async fn insert_registered_raw_messages( }) .collect::>(); runtime - .upsert_transcript_batch_for_test( - HostAdmissionScope::Profile, - &session, - &messages, - "lcm-compression-test-fixture", - tracedecay_global_db::ParseOffset::default(), - ) + .seed_session_messages_for_test(HostAdmissionScope::Profile, &session, &messages) .await .expect("registered raw message fixture") } diff --git a/crates/tracedecay/tests/session_suite/lcm_dag.rs b/crates/tracedecay/tests/session_suite/lcm_dag.rs index 126154f122..9c474c6404 100644 --- a/crates/tracedecay/tests/session_suite/lcm_dag.rs +++ b/crates/tracedecay/tests/session_suite/lcm_dag.rs @@ -47,7 +47,7 @@ impl ProfileLcmFixture for HostAdmissionTestRuntimeV1 { ) -> bool { self.upsert_session_message_for_test(HostAdmissionScope::Profile, message) .await - .unwrap_or(false) + .is_ok() } async fn lcm_insert_summary_node( @@ -219,9 +219,6 @@ async fn summary_node_preserves_source_lineage_and_expands_sources() { let db = registered_lcm_runtime(&tmp).await; let store_ids = insert_raw_messages(&db, "cursor", "session-1", &["alpha", "beta", "gamma"]).await; - let mut first_source = raw_message("cursor", "session-1-message-1", "session-1", 1, "alpha"); - first_source.timestamp = Some(1_715_000_001_000_000); - assert!(db.upsert_session_message(&first_source).await); let node = db .lcm_insert_summary_node(summary_draft( @@ -1030,11 +1027,9 @@ async fn summary_grep_denies_dirty_raw_sources_before_convergence() { 1, "revised source", ); - assert!( - db.upsert_session_message_for_test(HostAdmissionScope::Profile, &revised) - .await - .expect("revise canonical raw source") - ); + db.upsert_session_message_for_test(HostAdmissionScope::Profile, &revised) + .await + .expect("revise canonical raw source"); // The source-change trigger closes retrieval immediately, before the // background convergence worker updates generation-bound availability. let availability = db diff --git a/crates/tracedecay/tests/session_suite/lcm_query/mod.rs b/crates/tracedecay/tests/session_suite/lcm_query/mod.rs index a9ba025fb4..22b75d669b 100644 --- a/crates/tracedecay/tests/session_suite/lcm_query/mod.rs +++ b/crates/tracedecay/tests/session_suite/lcm_query/mod.rs @@ -1,5 +1,4 @@ use tempfile::TempDir; -use tracedecay_global_db::ParseOffset; use tracedecay_lcm::{ LCM_SCHEMA_VERSION, LcmContentSlice, LcmDescribeRequest, LcmDescribeTarget, LcmError, LcmExpandQueryRequest, LcmExpandRequest, LcmExpandTarget, LcmGcConfig, LcmGrepRequest, @@ -27,12 +26,10 @@ trait ProfileLcmFixture { async fn upsert_session_message(&self, message: &SessionMessageRecord) -> bool; - async fn upsert_transcript_batch( + async fn seed_session_messages( &self, session: &SessionRecord, messages: &[SessionMessageRecord], - source: &str, - offset: ParseOffset, ) -> bool; async fn lcm_insert_summary_node( @@ -58,25 +55,17 @@ impl ProfileLcmFixture for HostAdmissionTestRuntimeV1 { async fn upsert_session_message(&self, message: &SessionMessageRecord) -> bool { self.upsert_session_message_for_test(HostAdmissionScope::Profile, message) .await - .unwrap_or(false) + .is_ok() } - async fn upsert_transcript_batch( + async fn seed_session_messages( &self, session: &SessionRecord, messages: &[SessionMessageRecord], - source: &str, - offset: ParseOffset, ) -> bool { - self.upsert_transcript_batch_for_test( - HostAdmissionScope::Profile, - session, - messages, - source, - offset, - ) - .await - .is_ok() + self.seed_session_messages_for_test(HostAdmissionScope::Profile, session, messages) + .await + .is_ok() } async fn lcm_insert_summary_node( @@ -159,15 +148,9 @@ async fn insert_raw_messages( raw_message(provider, &message_id, session_id, (idx + 1) as i64, content) }) .collect(); - db.upsert_transcript_batch_for_test( - HostAdmissionScope::Profile, - &session, - &messages, - &format!("session-lcm-query-{provider}-{session_id}.jsonl"), - ParseOffset::default(), - ) - .await - .expect("registered transcript fixture should write") + db.seed_session_messages_for_test(HostAdmissionScope::Profile, &session, &messages) + .await + .expect("registered transcript fixture should write") } async fn replace_inline_content_without_updating_hash( diff --git a/crates/tracedecay/tests/session_suite/lcm_query/sessions.rs b/crates/tracedecay/tests/session_suite/lcm_query/sessions.rs index a2413a2540..8b6be84c28 100644 --- a/crates/tracedecay/tests/session_suite/lcm_query/sessions.rs +++ b/crates/tracedecay/tests/session_suite/lcm_query/sessions.rs @@ -71,15 +71,7 @@ async fn recent_sessions_uses_store_order_for_null_timestamp_activity() { "ingested later without source timestamp", ); message.timestamp = None; - assert!( - db.upsert_transcript_batch( - &session, - &[message], - "session-lcm-query-cursor-null-timestamp-session.jsonl", - ParseOffset::default(), - ) - .await - ); + assert!(db.seed_session_messages(&session, &[message]).await); let sessions = db .lcm_recent_sessions_for_test(None, 1) diff --git a/crates/tracedecay/tests/session_suite/lcm_raw.rs b/crates/tracedecay/tests/session_suite/lcm_raw.rs index 4718b68c88..dc5b762e9b 100644 --- a/crates/tracedecay/tests/session_suite/lcm_raw.rs +++ b/crates/tracedecay/tests/session_suite/lcm_raw.rs @@ -160,11 +160,9 @@ async fn search_uses_bounded_projection_but_load_recovers_raw() { "x".repeat(tracedecay_lcm::MAX_DERIVED_TEXT_CHARS * 5) ); let message = sample_message("cursor", "message-1", "session-1", &oversized); - assert!( - db.upsert_session_message_for_test(HostAdmissionScope::Profile, &message) - .await - .unwrap() - ); + db.upsert_session_message_for_test(HostAdmissionScope::Profile, &message) + .await + .unwrap(); let results = db .search_session_messages_for_test( diff --git a/crates/tracedecay/tests/session_suite/lcm_summary_lineage_review.rs b/crates/tracedecay/tests/session_suite/lcm_summary_lineage_review.rs index 168191c304..39ef3038a8 100644 --- a/crates/tracedecay/tests/session_suite/lcm_summary_lineage_review.rs +++ b/crates/tracedecay/tests/session_suite/lcm_summary_lineage_review.rs @@ -51,7 +51,7 @@ impl ProfileLcmFixture for HostAdmissionTestRuntimeV1 { ) -> bool { self.upsert_session_message_for_test(HostAdmissionScope::Profile, message) .await - .unwrap_or(false) + .is_ok() } async fn lcm_publish_immutable_summary( diff --git a/crates/tracedecay/tests/session_suite/message_search_eval_test.rs b/crates/tracedecay/tests/session_suite/message_search_eval_test.rs index e34000db1a..ef4bf10d13 100644 --- a/crates/tracedecay/tests/session_suite/message_search_eval_test.rs +++ b/crates/tracedecay/tests/session_suite/message_search_eval_test.rs @@ -83,13 +83,14 @@ async fn seed_corpus(db: &HostAdmissionTestRuntimeV1, fixture: &Value) { .with_tool_names(message["tool_names"].as_str()) .with_source(Some("/tmp/project/transcript.jsonl"), Some(ordinal as i64)) .build(); - assert!( - db.upsert_session_message_for_test(HostAdmissionScope::Project, &record) - .await - .expect("seed registered session message"), - "seed message {}", - message["id"].as_str().unwrap_or("?") - ); + db.upsert_session_message_for_test(HostAdmissionScope::Project, &record) + .await + .unwrap_or_else(|error| { + panic!( + "seed message {}: {error}", + message["id"].as_str().unwrap_or("?") + ) + }); } } } diff --git a/crates/tracedecay/tests/session_suite/transcript_store.rs b/crates/tracedecay/tests/session_suite/transcript_store.rs index 52f4b33b34..c5174559a4 100644 --- a/crates/tracedecay/tests/session_suite/transcript_store.rs +++ b/crates/tracedecay/tests/session_suite/transcript_store.rs @@ -4,8 +4,6 @@ use tracedecay_project::test_support::host_admission::HostAdmissionTestRuntimeV1 use tracedecay_sessions::admission::HostAdmissionScope; use tracedecay_store::{TranscriptStore, TranscriptStoreError, TranscriptWriteBatch}; -use crate::common::{global_message as sample_message, global_session as sample_session}; - async fn profile_runtime(tmp: &TempDir) -> HostAdmissionTestRuntimeV1 { HostAdmissionTestRuntimeV1::profile(tmp.path().join(".tracedecay")) .await @@ -47,689 +45,6 @@ async fn store_counts( } } -fn summary_message( - provider: &str, - message_id: &str, - session_id: &str, -) -> tracedecay_sessions::runtime::SessionMessageRecord { - let mut summary = sample_message( - provider, - message_id, - session_id, - "Compacted transcript summary.", - ); - summary.ordinal = 2; - summary.kind = Some("summary".to_string()); - summary.metadata_json = Some( - serde_json::json!({ - "source": "codex_context_compacted", - "summary_body": "plaintext", - "codex_compaction_depth": 1 - }) - .to_string(), - ); - summary -} - -#[tokio::test] -async fn transcript_batch_survives_restart_and_replay_is_idempotent() { - let tmp = TempDir::new().unwrap(); - let transcript_path = tmp.path().join("cursor-restart.jsonl"); - let mut session = sample_session("cursor", "restart-session", "project-a"); - session.transcript_path = Some(transcript_path.to_string_lossy().to_string()); - let mut second = sample_message( - "cursor", - "restart-message-2", - "restart-session", - "Second durable restart message.", - ); - second.ordinal = 2; - let messages = vec![ - sample_message( - "cursor", - "restart-message-1", - "restart-session", - "First durable restart message.", - ), - second, - ]; - let offset = ParseOffset { - byte_offset: 512, - mtime: 1_800_000_000, - file_id: 42, - }; - - let db = profile_runtime(&tmp).await; - db.transcript_store_for_test(HostAdmissionScope::Profile) - .unwrap() - .persist_transcript_batch( - TranscriptWriteBatch::upsert( - session.clone(), - messages.clone(), - ParseOffset::default(), - offset, - ) - .unwrap(), - ) - .await - .unwrap(); - drop(db); - - let reopened = profile_runtime(&tmp).await; - assert_eq!( - reopened - .parse_offset_for_test( - HostAdmissionScope::Profile, - transcript_path.to_string_lossy().as_ref(), - ) - .await - .unwrap(), - Some(offset) - ); - reopened - .transcript_store_for_test(HostAdmissionScope::Profile) - .unwrap() - .persist_transcript_batch( - TranscriptWriteBatch::upsert(session, messages, offset, offset).unwrap(), - ) - .await - .unwrap(); - - assert_eq!( - store_counts(&reopened, "cursor", "restart-session", &transcript_path).await, - StoreCounts { - sessions: 1, - raw_messages: 2, - raw_fts: 2, - all_raw_fts: 2, - summaries: 0, - cursors: 1, - } - ); -} - -#[tokio::test] -async fn late_cursor_failure_rolls_back_every_transcript_write_then_retries() { - let tmp = TempDir::new().unwrap(); - let transcript_path = tmp.path().join("late-cursor-failure.jsonl"); - let payload_dir = tracedecay_lcm::payload::payload_dir(&tmp.path().join(".tracedecay")); - std::fs::create_dir_all(&payload_dir).unwrap(); - let sentinel_path = payload_dir.join("preexisting.payload"); - std::fs::write(&sentinel_path, "must survive rollback").unwrap(); - let mut session = sample_session("codex", "atomic-session", "project-a"); - session.transcript_path = Some(transcript_path.to_string_lossy().to_string()); - let mut source = sample_message( - "codex", - "source-message", - "atomic-session", - &format!("oversized tool payload\n{}", "P".repeat(300_000)), - ); - source.role = "tool".to_string(); - source.kind = Some("tool_result".to_string()); - let mut summary = sample_message( - "codex", - "summary-message", - "atomic-session", - "Compacted transcript summary.", - ); - summary.ordinal = 2; - summary.kind = Some("summary".to_string()); - summary.metadata_json = Some( - serde_json::json!({ - "source": "codex_context_compacted", - "summary_body": "plaintext", - "codex_compaction_depth": 1 - }) - .to_string(), - ); - let batch = TranscriptWriteBatch::upsert( - session, - vec![source, summary], - ParseOffset::default(), - ParseOffset { - byte_offset: 384, - mtime: 1_800_000_100, - file_id: 43, - }, - ) - .unwrap(); - - let db = profile_runtime(&tmp).await; - db.set_parse_offset_insert_failure_for_test(HostAdmissionScope::Profile, true) - .await - .unwrap(); - - let error = db - .transcript_store_for_test(HostAdmissionScope::Profile) - .unwrap() - .persist_transcript_batch(batch.clone()) - .await - .expect_err("the late cursor write must fail the batch"); - assert!(matches!(error, TranscriptStoreError::Storage { .. })); - assert_eq!( - db.parse_offset_for_test( - HostAdmissionScope::Profile, - transcript_path.to_string_lossy().as_ref(), - ) - .await - .unwrap(), - None - ); - - assert_eq!( - std::fs::read_to_string(&sentinel_path).unwrap(), - "must survive rollback" - ); - let remaining_payload_files = std::fs::read_dir(&payload_dir) - .unwrap() - .map(|entry| entry.unwrap().file_name()) - .collect::>(); - assert_eq!( - remaining_payload_files, - vec![std::ffi::OsString::from("preexisting.payload")] - ); - - assert_eq!( - store_counts(&db, "codex", "atomic-session", &transcript_path).await, - StoreCounts { - sessions: 0, - raw_messages: 0, - raw_fts: 0, - all_raw_fts: 0, - summaries: 0, - cursors: 0, - } - ); - - db.set_parse_offset_insert_failure_for_test(HostAdmissionScope::Profile, false) - .await - .unwrap(); - drop(db); - - let reopened = profile_runtime(&tmp).await; - reopened - .transcript_store_for_test(HostAdmissionScope::Profile) - .unwrap() - .persist_transcript_batch(batch) - .await - .unwrap(); - let raw = reopened - .lcm_load_raw_message_for_test("codex", "source-message") - .await - .expect("oversized message must persist on retry"); - let payload_ref = raw.payload_ref.expect("oversized message must externalize"); - assert!(payload_dir.join(payload_ref).is_file()); - assert_eq!( - std::fs::read_to_string(&sentinel_path).unwrap(), - "must survive rollback" - ); - assert_eq!( - store_counts(&reopened, "codex", "atomic-session", &transcript_path).await, - StoreCounts { - sessions: 1, - raw_messages: 2, - raw_fts: 2, - all_raw_fts: 2, - summaries: 0, - cursors: 1, - } - ); -} - -#[tokio::test] -async fn invalid_batch_mutates_no_transcript_state() { - let tmp = TempDir::new().unwrap(); - let transcript_path = tmp.path().join("invalid-batch.jsonl"); - let db = profile_runtime(&tmp).await; - let mut session = sample_session("cursor", "expected-session", "project-a"); - session.transcript_path = Some(transcript_path.to_string_lossy().to_string()); - let error = TranscriptWriteBatch::upsert( - session, - vec![sample_message( - "cursor", - "invalid-message", - "invalid-session", - "must never be persisted", - )], - ParseOffset::default(), - ParseOffset { - byte_offset: 77, - mtime: 1_715_000_350, - file_id: 12, - }, - ) - .expect_err("a mismatched message identity must be rejected"); - - assert!(matches!( - error, - TranscriptStoreError::MessageIdentityMismatch { .. } - )); - assert_eq!( - store_counts(&db, "cursor", "expected-session", &transcript_path).await, - StoreCounts { - sessions: 0, - raw_messages: 0, - raw_fts: 0, - all_raw_fts: 0, - summaries: 0, - cursors: 0, - } - ); -} - -#[tokio::test] -async fn stale_higher_batch_is_rejected_until_reparsed_from_durable_cursor() { - let tmp = TempDir::new().unwrap(); - let transcript_path = tmp.path().join("concurrent.jsonl"); - let db = profile_runtime(&tmp).await; - let store = db - .transcript_store_for_test(HostAdmissionScope::Profile) - .unwrap(); - let mut session = sample_session("cursor", "concurrent-session", "project-a"); - session.transcript_path = Some(transcript_path.to_string_lossy().to_string()); - let first_message = sample_message( - "cursor", - "concurrent-first-message", - "concurrent-session", - "first committed transcript message", - ); - let mut second_message = sample_message( - "cursor", - "concurrent-second-message", - "concurrent-session", - "second committed transcript message", - ); - second_message.ordinal = 2; - let mut summary = summary_message("cursor", "concurrent-summary", "concurrent-session"); - summary.ordinal = 3; - let first_offset = ParseOffset { - byte_offset: 100, - mtime: 1_000, - file_id: 7, - }; - let second_offset = ParseOffset { - byte_offset: 200, - mtime: 2_000, - file_id: 7, - }; - - store - .persist_transcript_batch( - TranscriptWriteBatch::upsert( - session.clone(), - vec![first_message.clone()], - ParseOffset::default(), - first_offset, - ) - .unwrap(), - ) - .await - .unwrap(); - let stale_higher_error = store - .persist_transcript_batch( - TranscriptWriteBatch::upsert( - session.clone(), - vec![ - first_message.clone(), - second_message.clone(), - summary.clone(), - ], - ParseOffset::default(), - second_offset, - ) - .unwrap(), - ) - .await - .expect_err("a pre-parsed batch must not change its observed cursor and retry"); - assert!(matches!( - stale_higher_error, - TranscriptStoreError::Conflict { - expected, - actual, - .. - } if expected == ParseOffset::default() && actual == first_offset - )); - - assert_eq!( - db.parse_offset_for_test( - HostAdmissionScope::Profile, - transcript_path.to_string_lossy().as_ref(), - ) - .await - .unwrap(), - Some(first_offset) - ); - assert!( - db.session_message_for_test( - HostAdmissionScope::Profile, - "cursor", - "concurrent-second-message", - ) - .await - .unwrap() - .is_none(), - "the stale parse products must roll back with the cursor conflict" - ); - - // A runtime boundary may re-read the winner and reparse the suffix/full - // source. That fresh batch carries the actually observed durable cursor. - store - .persist_transcript_batch( - TranscriptWriteBatch::upsert( - session, - vec![first_message, second_message, summary], - first_offset, - second_offset, - ) - .unwrap(), - ) - .await - .expect("a freshly parsed batch may advance from the durable winner"); - - let mut stale_session = sample_session("cursor", "concurrent-session", "project-a"); - stale_session.transcript_path = Some(transcript_path.to_string_lossy().to_string()); - let stale_error = store - .persist_transcript_batch( - TranscriptWriteBatch::upsert( - stale_session, - vec![sample_message( - "cursor", - "concurrent-stale-message", - "concurrent-session", - "stale transcript message", - )], - ParseOffset::default(), - first_offset, - ) - .unwrap(), - ) - .await - .expect_err("a lower stale cursor must not replace the converged maximum"); - assert!(matches!( - stale_error, - TranscriptStoreError::Conflict { - actual, - .. - } if actual == second_offset - )); - assert_eq!( - store_counts(&db, "cursor", "concurrent-session", &transcript_path).await, - StoreCounts { - sessions: 1, - raw_messages: 3, - raw_fts: 3, - all_raw_fts: 3, - summaries: 0, - cursors: 1, - } - ); -} - -#[tokio::test] -async fn concurrent_full_batches_converge_without_split_brain_or_partial_writes() { - let tmp = TempDir::new().unwrap(); - let transcript_path = tmp.path().join("concurrent-full-batches.jsonl"); - let db = profile_runtime(&tmp).await; - let first_store = db - .transcript_store_for_test(HostAdmissionScope::Profile) - .unwrap(); - let second_store = db - .transcript_store_for_test(HostAdmissionScope::Profile) - .unwrap(); - let mut session = sample_session("cursor", "concurrent-full-session", "project-a"); - session.transcript_path = Some(transcript_path.to_string_lossy().to_string()); - let first_message = sample_message( - "cursor", - "concurrent-full-message-1", - "concurrent-full-session", - "first concurrent transcript message", - ); - let mut second_message = sample_message( - "cursor", - "concurrent-full-message-2", - "concurrent-full-session", - "second concurrent transcript message", - ); - second_message.ordinal = 2; - let mut summary = summary_message( - "cursor", - "concurrent-full-summary", - "concurrent-full-session", - ); - summary.ordinal = 3; - let first_offset = ParseOffset { - byte_offset: 100, - mtime: 1_000, - file_id: 7, - }; - let higher_offset = ParseOffset { - byte_offset: 200, - mtime: 2_000, - file_id: 7, - }; - let first_batch = TranscriptWriteBatch::upsert( - session.clone(), - vec![first_message.clone()], - ParseOffset::default(), - first_offset, - ) - .unwrap(); - let competing_batch = TranscriptWriteBatch::upsert( - session.clone(), - vec![second_message.clone()], - ParseOffset::default(), - first_offset, - ) - .unwrap(); - - let (first_result, competing_result) = tokio::join!( - first_store.persist_transcript_batch(first_batch), - second_store.persist_transcript_batch(competing_batch), - ); - - let conflict_actual = match (first_result, competing_result) { - (Ok(()), Err(TranscriptStoreError::Conflict { actual, .. })) - | (Err(TranscriptStoreError::Conflict { actual, .. }), Ok(())) => actual, - outcomes => panic!("exactly one concurrent full batch must commit, got {outcomes:?}"), - }; - assert_eq!(conflict_actual, first_offset); - assert_eq!( - db.parse_offset_for_test( - HostAdmissionScope::Profile, - transcript_path.to_string_lossy().as_ref(), - ) - .await - .unwrap(), - Some(first_offset) - ); - assert_eq!( - store_counts(&db, "cursor", "concurrent-full-session", &transcript_path,).await, - StoreCounts { - sessions: 1, - raw_messages: 1, - raw_fts: 1, - all_raw_fts: 1, - summaries: 0, - cursors: 1, - } - ); - - first_store - .persist_transcript_batch( - TranscriptWriteBatch::upsert( - session.clone(), - vec![first_message, second_message, summary], - conflict_actual, - higher_offset, - ) - .unwrap(), - ) - .await - .expect("a freshly parsed batch must advance from the returned durable cursor"); - assert_eq!( - db.parse_offset_for_test( - HostAdmissionScope::Profile, - transcript_path.to_string_lossy().as_ref(), - ) - .await - .unwrap(), - Some(higher_offset) - ); - - let stale_error = first_store - .persist_transcript_batch( - TranscriptWriteBatch::upsert( - session.clone(), - vec![sample_message( - "cursor", - "concurrent-full-stale-message", - "concurrent-full-session", - "stale owner must not mutate state", - )], - first_offset, - ParseOffset { - byte_offset: 150, - mtime: 1_500, - file_id: 7, - }, - ) - .unwrap(), - ) - .await - .expect_err("a stale owner behind the durable maximum must be rejected"); - assert!(matches!( - stale_error, - TranscriptStoreError::Conflict { actual, .. } if actual == higher_offset - )); - - let competing_error = second_store - .persist_transcript_batch( - TranscriptWriteBatch::upsert( - session, - vec![sample_message( - "cursor", - "concurrent-full-competing-message", - "concurrent-full-session", - "competing file identity must not mutate state", - )], - ParseOffset::default(), - ParseOffset { - byte_offset: 400, - mtime: 4_000, - file_id: 8, - }, - ) - .unwrap(), - ) - .await - .expect_err("a competing file identity must be rejected"); - assert!(matches!( - competing_error, - TranscriptStoreError::Conflict { actual, .. } if actual == higher_offset - )); - - // Transcript persistence never projects summary nodes: a `kind = "summary"` - // transcript message is durable raw evidence, and `session_summary_nodes` rows - // are only ever written by the immutable-summary publication path - // (`lcm_publish_immutable_summary_guarded`, which LCM compression reaches - // through `dag::insert_summary_node`). See - // `transcript_summary_message_keeps_native_compaction_evidence` for the - // other half of that contract. - assert_eq!( - store_counts(&db, "cursor", "concurrent-full-session", &transcript_path,).await, - StoreCounts { - sessions: 1, - raw_messages: 3, - raw_fts: 3, - all_raw_fts: 3, - summaries: 0, - cursors: 1, - } - ); -} - -/// `summaries: 0` above is only safe because ingest keeps the evidence the -/// summarization pipeline recognizes. `native_summary_evidence` scans -/// `session_messages` for a non-empty body plus `kind = "summary"` and the -/// provider's metadata discriminators, and only a recognized row reaches -/// compression, the sole production writer of `session_summary_nodes`. Dropping -/// or rewriting either column at persist time would strand every host -/// compaction with no failing count to show for it. -#[tokio::test] -async fn transcript_summary_message_keeps_native_compaction_evidence() { - let tmp = TempDir::new().unwrap(); - let transcript_path = tmp.path().join("codex-compaction.jsonl"); - let db = profile_runtime(&tmp).await; - let store = db - .transcript_store_for_test(HostAdmissionScope::Profile) - .unwrap(); - let mut session = sample_session("codex", "codex-compaction-session", "project-a"); - session.transcript_path = Some(transcript_path.to_string_lossy().to_string()); - store - .persist_transcript_batch( - TranscriptWriteBatch::upsert( - session, - vec![summary_message( - "codex", - "codex-compaction-summary", - "codex-compaction-session", - )], - ParseOffset::default(), - ParseOffset { - byte_offset: 120, - mtime: 1_000, - file_id: 11, - }, - ) - .unwrap(), - ) - .await - .unwrap(); - - assert_eq!( - store_counts(&db, "codex", "codex-compaction-session", &transcript_path).await, - StoreCounts { - sessions: 1, - raw_messages: 1, - raw_fts: 1, - all_raw_fts: 1, - summaries: 0, - cursors: 1, - } - ); - - let snapshot_path = tmp.path().join("codex-compaction-snapshot.db"); - db.snapshot_session_database_for_test(HostAdmissionScope::Profile, &snapshot_path) - .await - .unwrap(); - let connection = rusqlite::Connection::open(&snapshot_path).unwrap(); - let (text, kind, metadata_json) = connection - .query_row( - "SELECT COALESCE(content, placeholder_text, ''), kind, metadata_json - FROM lcm_raw_messages - WHERE provider = ?1 AND session_id = ?2", - ("codex", "codex-compaction-session"), - |row| { - Ok(( - row.get::<_, String>(0)?, - row.get::<_, Option>(1)?, - row.get::<_, Option>(2)?, - )) - }, - ) - .unwrap(); - assert!( - !text.trim().is_empty(), - "the evidence scan skips blank bodies, so the summary text must survive persist" - ); - assert_eq!(kind.as_deref(), Some("summary")); - let metadata: serde_json::Value = - serde_json::from_str(&metadata_json.expect("summary metadata must survive persist")) - .unwrap(); - assert_eq!(metadata["source"], "codex_context_compacted"); - assert_eq!(metadata["summary_body"], "plaintext"); -} - #[tokio::test] async fn concurrent_empty_advances_converge_to_highest_compatible_offset_without_rows() { let tmp = TempDir::new().unwrap();