From 75ffb4e26cc0d5164183788aeefec6a542e58bb4 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Wed, 30 Sep 2026 07:28:25 +0000 Subject: [PATCH] simplify(sessions): delete the test-only transcript upsert port Production transcript ingest only advances cursors through TranscriptWriteBatch::advance_offset; session and message rows arrive through observation admission and projection. The Upsert write kind and its global-db/session-store chain (upsert_transcript_batch, persist_transcript_batch_result, the Git-evidence variant, batch message staging, write-time timestamp normalization) were reachable only from tests. Delete that chain, the unreachable TranscriptStoreError variants and their ingest failure codes, TranscriptBatch, and TranscriptGitEvidence. TranscriptWriteBatch now carries only the offset pair. Fixture helpers seed session rows plus raw LCM messages through the existing upsert_session and lcm_ingest_raw_message primitives and return typed errors instead of a swallowed false. Tests that only protected the Upsert write (restart/replay, rollback, identity validation, full-batch CAS, write-time microsecond normalization) are removed; the profile-scope Git-evidence refusal now exercises the production span writer. --- .../cli_non_interactive_test.rs | 36 +- crates/tracedecay-global-db/src/api_types.rs | 1 - .../src/git_correlation_adapter.rs | 60 +- crates/tracedecay-global-db/src/lib.rs | 2 +- .../src/observation_projection/state.rs | 6 +- .../tracedecay-global-db/src/tests/harness.rs | 35 - crates/tracedecay-global-db/src/transcript.rs | 63 +- .../src/test_support/host_admission.rs | 51 +- .../host_admission/session_test_support.rs | 27 - .../src/transcript.rs | 89 +-- .../src/lcm_effects/tests.rs | 70 +- .../src/retained/profile.rs | 15 +- .../worker.rs | 15 +- .../src/runtime/ingest/failure.rs | 9 - .../src/runtime/ingest/scheduler.rs | 15 +- crates/tracedecay-sessions/src/runtime/mod.rs | 2 +- .../src/runtime/store_access/mod.rs | 6 +- .../src/runtime/store_access/transcript.rs | 372 +--------- .../src/runtime/store_access/types.rs | 13 +- crates/tracedecay-store/src/lib.rs | 2 +- crates/tracedecay-store/src/transcript.rs | 250 +------ crates/tracedecay/src/daemon/scheduler.rs | 15 +- .../src/mcp/server/host_admission_tests.rs | 4 +- .../support/fixtures.rs | 32 +- crates/tracedecay/tests/common/mod.rs | 2 +- .../tests/dashboard_api_test/analytics.rs | 10 +- .../tests/dashboard_api_test/delivery.rs | 9 +- .../tests/dashboard_api_test/loom.rs | 21 +- .../tests/dashboard_api_test/runtime.rs | 53 +- .../tests/dashboard_api_test/savings.rs | 211 +++--- .../hooks_lsp_suite/hint_settlement_test.rs | 4 +- .../tests/mcp_suite/git_correlation_test.rs | 16 +- crates/tracedecay/tests/mcp_suite/support.rs | 55 +- .../tests/session_suite/git_backfill.rs | 56 +- .../tests/session_suite/global_db.rs | 2 +- .../session_suite/lcm_compression/mod.rs | 8 +- .../tracedecay/tests/session_suite/lcm_dag.rs | 13 +- .../tests/session_suite/lcm_query/mod.rs | 35 +- .../tests/session_suite/lcm_query/sessions.rs | 10 +- .../tracedecay/tests/session_suite/lcm_raw.rs | 8 +- .../lcm_summary_lineage_review.rs | 2 +- .../session_suite/message_search_eval_test.rs | 15 +- .../tests/session_suite/transcript_store.rs | 685 ------------------ 43 files changed, 362 insertions(+), 2043 deletions(-) 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();