diff --git a/crates/tracedecay-global-db/src/observation_projection/state.rs b/crates/tracedecay-global-db/src/observation_projection/state.rs index 2fffbb103c..686d512f6c 100644 --- a/crates/tracedecay-global-db/src/observation_projection/state.rs +++ b/crates/tracedecay-global-db/src/observation_projection/state.rs @@ -1280,11 +1280,6 @@ pub(super) fn reconcile_session_rows_detailed( if actual.provider != expected.provider || actual.session_id != expected.session_id { return Err(SessionReconcileConflict("identity")); } - // LCM may insert this host session before the rollout projection, with no - // project. That row is not a second root: the observation replaces it. - if session_project_is_unscoped_placeholder(actual) { - return Ok(expected.clone()); - } // `expected` is the projection being applied now. A typed project id that // changed (re-enroll/reset, or a cwd that now resolves to another // registered project) moves this host session onto that current id. @@ -1373,16 +1368,6 @@ fn project_key_is_directory(session: &SessionRecord) -> bool { session.project_key == session.project_path } -/// An LCM foreign-key shell: placeholder project fields, no transcript, no metadata. -fn session_project_is_unscoped_placeholder(session: &SessionRecord) -> bool { - session.transcript_path.is_none() - && session.metadata_json.is_none() - && tracedecay_lcm::compression::lcm_unscoped_session_project( - &session.project_key, - &session.project_path, - ) -} - fn reconcile_optional( field: &'static str, actual: Option<&T>, @@ -1581,13 +1566,14 @@ mod reconcile_tests { CanonicalObservationEnvelopeV1, ComponentVersion, ObservationId, ObservationIdentityMaterialV1, ObservationOrderingDomainV1, ObservationScopeV1, ObservationSourceGenerationV1, ObservationSourceIdentityV1, ObservationSourceRangeV1, - PayloadReferenceV1, RetentionClass, SanitizationReceiptId, SanitizationReceiptRefV1, - SanitizationReceiptV1, SanitizerDispositionV1, SensitivityV1, + PayloadReferenceV1, ProjectId, RetentionClass, SanitizationReceiptId, + SanitizationReceiptRefV1, SanitizationReceiptV1, SanitizerDispositionV1, SensitivityV1, }; #[cfg(unix)] use tracedecay_runtime_core::db::engine::params; use tracedecay_store::{ ObservationProjection, ProjectionStoreError, SessionMessageRecord, SessionRecord, + session_project_fields, }; use super::{ @@ -1812,35 +1798,37 @@ mod reconcile_tests { } #[test] - fn lcm_placeholder_project_is_replaced_by_the_rollout_session() { - let mut stored = record("lcm-active-context"); - stored.title = Some("LCM active context".to_owned()); - let mut rollout = record("/Volumes/bigssd/projects/core"); - rollout.title = Some("rollout".to_owned()); - rollout.transcript_path = Some("/Volumes/bigssd/projects/core/session.jsonl".to_owned()); - - let merged = reconcile_session_rows_detailed(&stored, &rollout) - .expect("an LCM placeholder must take the rollout project"); - - assert_eq!(merged.project_key, "/Volumes/bigssd/projects/core"); - assert_eq!(merged.project_path, "/Volumes/bigssd/projects/core"); - assert_eq!(merged.title.as_deref(), Some("rollout")); - assert_eq!( - merged.transcript_path.as_deref(), - Some("/Volumes/bigssd/projects/core/session.jsonl") - ); - } - - #[test] - fn unknown_project_placeholder_is_replaced_by_the_rollout_session() { - let stored = record("unknown"); - let rollout = record("/work/repo"); - - let merged = reconcile_session_rows_detailed(&stored, &rollout) - .expect("an unknown project shell must take the rollout project"); - - assert_eq!(merged.project_key, "/work/repo"); - assert_eq!(merged.project_path, "/work/repo"); + fn a_scope_shell_row_takes_the_rollout_working_directory() { + for (scope, key) in [ + ( + ObservationScopeV1::Project { + project_id: ProjectId::new("project.core").unwrap(), + }, + "project.core", + ), + (ObservationScopeV1::Profile, "user"), + ] { + let (shell_key, shell_path) = session_project_fields(&scope); + let mut shell = record(&shell_path); + shell.project_key = shell_key; + shell.started_at = Some(5); + let mut rollout = record("/work/repo"); + rollout.project_key = key.to_owned(); + rollout.title = Some("rollout".to_owned()); + rollout.transcript_path = Some("/work/repo/session.jsonl".to_owned()); + + let merged = reconcile_session_rows_detailed(&shell, &rollout) + .expect("the shell row and the rollout name one session"); + + assert_eq!(merged.project_key, key); + assert_eq!(merged.project_path, "/work/repo"); + assert_eq!(merged.title.as_deref(), Some("rollout")); + assert_eq!(merged.started_at, Some(1)); + assert_eq!( + merged.transcript_path.as_deref(), + Some("/work/repo/session.jsonl") + ); + } } #[test] diff --git a/crates/tracedecay-global-db/src/registered_lcm.rs b/crates/tracedecay-global-db/src/registered_lcm.rs index fab1954a49..ecc1781481 100644 --- a/crates/tracedecay-global-db/src/registered_lcm.rs +++ b/crates/tracedecay-global-db/src/registered_lcm.rs @@ -260,6 +260,7 @@ impl RegisteredGlobalDb { F: FnOnce() -> Result<(), LcmError>, { check_execution(control)?; + let scope = SessionStoreAccess::new(self).lcm_session_scope()?; let storage_root = self.lcm_storage_root()?; let session_id = SessionId::new(request.session_id.clone()).map_err(|error| { LcmError::Db(format!( @@ -288,6 +289,7 @@ impl RegisteredGlobalDb { ); let mut response = compression::compress( &transaction, + &scope, &publisher, storage_root, request, @@ -330,6 +332,7 @@ impl RegisteredGlobalDb { F: FnOnce() -> Result<(), LcmError>, { check_execution(control)?; + let scope = SessionStoreAccess::new(self).lcm_session_scope()?; let storage_root = self.lcm_storage_root()?; let session_id = SessionId::new(request.session_id.clone()).map_err(|error| { LcmError::Db(format!( @@ -369,6 +372,7 @@ impl RegisteredGlobalDb { ); let bounded = compression::compress_retained_page( &transaction, + &scope, &publisher, storage_root, request, diff --git a/crates/tracedecay-global-db/src/session_project_rebind_tests.rs b/crates/tracedecay-global-db/src/session_project_rebind_tests.rs index b15a2f6526..9d6dd6b18f 100644 --- a/crates/tracedecay-global-db/src/session_project_rebind_tests.rs +++ b/crates/tracedecay-global-db/src/session_project_rebind_tests.rs @@ -1,7 +1,7 @@ //! A host session is keyed by provider and session id. LCM may insert that -//! row before the rollout projection, and a later project binding must replace -//! the placeholder. A genuine collision stays on that queue row and must not -//! stop later sessions. +//! row before the rollout projection, under its store's scope; the rollout +//! then refines its working directory. A genuine collision stays on that +//! queue row and must not stop later sessions. use serde_json::{Value, json}; use tempfile::TempDir; @@ -14,12 +14,14 @@ use tracedecay_domain::{ RetentionClass, SanitizationReceiptId, SanitizationReceiptRefV1, SanitizationReceiptV1, SanitizerDispositionV1, SensitivityV1, SessionId, UtcMicros, }; +use tracedecay_lcm::{LcmCompressionRequest, LcmSummarizerMode}; use tracedecay_runtime_core::db::engine::params; use tracedecay_sessions::runtime::shared::durable_project_path_key; use tracedecay_store::{ AnchoredObservationWrite, ObservationPersistOutcome, ObservationProjectionStore, ObservationStore, ObservationWrite, }; +use tracedecay_temporal_query::execution::ExecutionControl; use crate::tests::harness::{ HostAdmissionScope, HostAdmissionTestRuntimeV1, reopen_project_sessions_database, @@ -310,8 +312,57 @@ async fn changed_project_key_rebinds_without_blocking_other_sessions() { ); } +fn lcm_request(provider: &str, session_id: &str) -> LcmCompressionRequest { + LcmCompressionRequest { + provider: provider.to_owned(), + session_id: session_id.to_owned(), + messages: vec![json!({"id": "lcm-active-1", "role": "user", "content": "active turn"})], + current_tokens: Some(100), + focus_topic: None, + ignore_session_patterns: Vec::new(), + stateless_session_patterns: Vec::new(), + ignore_message_patterns: Vec::new(), + expected_current_frontier_store_id: None, + threshold_tokens: None, + max_assembly_tokens: None, + leaf_chunk_tokens: None, + max_source_messages: None, + summary_fan_in: None, + incremental_max_depth: None, + fresh_tail_count: None, + dynamic_leaf_chunk_enabled: None, + dynamic_leaf_chunk_max: None, + context_length: None, + reserve_tokens_floor: None, + summarizer: LcmSummarizerMode::Noop, + } +} + +async fn lcm_compress(database: &RegisteredGlobalDb, provider: &str, session_id: &str) { + let control = ExecutionControl::new(None); + let response = database + .lcm_compress_guarded(&lcm_request(provider, session_id), &control, || Ok(())) + .await + .unwrap(); + assert_eq!(response.status, "ok"); +} + +async fn stored_project( + runtime: &HostAdmissionTestRuntimeV1, + scope: HostAdmissionScope, + provider: &str, + session_id: &str, +) -> (String, String) { + let session = runtime + .session_for_test(scope, provider, session_id) + .await + .unwrap() + .expect("LCM compression stores the session it ingested"); + (session.project_key, session.project_path) +} + #[tokio::test] -async fn lcm_ensure_session_then_rollout_keeps_the_real_project() { +async fn lcm_compression_before_the_rollout_stores_the_session_under_its_store_scope() { let tmp = TempDir::new().unwrap(); let profile = tmp.path().join("profile"); let project_root = tmp.path().join("repo"); @@ -327,32 +378,41 @@ async fn lcm_ensure_session_then_rollout_keeps_the_real_project() { let database = runtime .registered_database(HostAdmissionScope::Project) .unwrap(); - let transaction = database - .runtime_database() - .begin_write_transaction("ensure lcm session before rollout projection") - .await - .unwrap(); - tracedecay_lcm::compression::ensure_session(&transaction, "codex", CODEX_SESSION) - .await - .unwrap(); - transaction.commit().await.unwrap(); + lcm_compress(database, "codex", CODEX_SESSION).await; + lcm_compress( + runtime + .registered_database(HostAdmissionScope::Profile) + .unwrap(), + "codex", + "session.profile-only", + ) + .await; - let (project_key, project_path) = - session_project_binding(database, "codex", CODEX_SESSION).await; assert_eq!( - project_key, - tracedecay_lcm::compression::LCM_UNKNOWN_PROJECT_KEY, - "LCM must not store a fake project key" + stored_project( + &runtime, + HostAdmissionScope::Project, + "codex", + CODEX_SESSION + ) + .await, + ("project.core".to_owned(), "project.core".to_owned()) ); assert_eq!( - project_path, - tracedecay_lcm::compression::LCM_UNKNOWN_PROJECT_KEY + stored_project( + &runtime, + HostAdmissionScope::Profile, + "codex", + "session.profile-only" + ) + .await, + ("user".to_owned(), "user".to_owned()) ); assert!( session_title(database, "codex", CODEX_SESSION) .await .is_none(), - "the foreign-key shell must not invent a session title" + "LCM must not invent a session title" ); let store = runtime @@ -383,17 +443,23 @@ async fn lcm_ensure_session_then_rollout_keeps_the_real_project() { .await; drain_projection_queue(&store).await; - let (project_key, project_path) = - session_project_binding(database, "codex", CODEX_SESSION).await; - assert_eq!(project_key, "project.core"); - assert_eq!(project_path, durable_project_path_key(&cwd)); + assert_eq!( + stored_project( + &runtime, + HostAdmissionScope::Project, + "codex", + CODEX_SESSION + ) + .await, + ("project.core".to_owned(), durable_project_path_key(&cwd)) + ); assert!( projected_text_contains(database, "codex", CODEX_SESSION, ORIGINAL_TEXT).await, - "the rollout message must project onto the placeholder session" + "the rollout message must project onto the session LCM created" ); assert!( projected_text_contains(database, "cursor", CURSOR_SESSION, CURSOR_TEXT).await, - "a later session must project after the placeholder is replaced" + "a later session must project after the rollout refines the session" ); assert!( store.next_queued_observation().await.unwrap().is_none(), @@ -401,6 +467,74 @@ async fn lcm_ensure_session_then_rollout_keeps_the_real_project() { ); } +#[tokio::test] +async fn a_rollout_starting_after_lcm_compression_keeps_its_own_start_time() { + const ROLLOUT_STARTED_AT: i64 = 4_102_444_800; + let tmp = TempDir::new().unwrap(); + let profile = tmp.path().join("profile"); + let project_root = tmp.path().join("repo"); + std::fs::create_dir_all(&project_root).unwrap(); + let cwd = project_root.to_string_lossy().into_owned(); + let runtime = HostAdmissionTestRuntimeV1::project( + &profile, + &project_root, + ProjectId::new("project.start").unwrap(), + ) + .await + .unwrap(); + lcm_compress( + runtime + .registered_database(HostAdmissionScope::Project) + .unwrap(), + "codex", + CODEX_SESSION, + ) + .await; + assert_eq!( + runtime + .session_for_test(HostAdmissionScope::Project, "codex", CODEX_SESSION) + .await + .unwrap() + .expect("LCM compression stores the session it ingested") + .started_at, + None, + "LCM did not observe the session start" + ); + + let store = runtime + .observation_store(HostAdmissionScope::Project) + .unwrap(); + let mut rollout_session = session_fact(&cwd); + let CanonicalObservationFactV1::Session { started_at, .. } = &mut rollout_session else { + unreachable!("session_fact builds a session fact") + }; + *started_at = Some(ROLLOUT_STARTED_AT); + persist( + &store, + observation( + "codex", + CODEX_SESSION, + project_scope("project.start"), + "record.codex.late-start", + ORIGINAL_TEXT, + "receipt.codex.late-start", + Some(rollout_session), + ), + ) + .await; + drain_projection_queue(&store).await; + + assert_eq!( + runtime + .session_for_test(HostAdmissionScope::Project, "codex", CODEX_SESSION) + .await + .unwrap() + .expect("the rollout projects onto the LCM session") + .started_at, + Some(ROLLOUT_STARTED_AT) + ); +} + #[tokio::test] async fn genuine_session_collision_is_recorded_and_later_sessions_project() { let tmp = TempDir::new().unwrap(); diff --git a/crates/tracedecay-lcm/src/compression.rs b/crates/tracedecay-lcm/src/compression.rs index ff1abb87d6..463e647609 100644 --- a/crates/tracedecay-lcm/src/compression.rs +++ b/crates/tracedecay-lcm/src/compression.rs @@ -5,8 +5,9 @@ use serde_json::{Map, Value, json}; use crate::message_storage_text; use crate::retrieval_content::projected_content_hash; +use tracedecay_domain::ObservationScopeV1; use tracedecay_runtime_core::db::engine::{Executor, QueryExecutor, Value as SqlValue, params}; -use tracedecay_store::SessionMessageRecord; +use tracedecay_store::{SessionMessageRecord, session_project_fields}; use super::compression_decision::{ self, AssemblyCapInput, CompressionPlanInput, CondensationCandidateDecision, @@ -155,6 +156,7 @@ pub async fn lifecycle_state( /// while pressure is still unrelieved. pub async fn record_session_boundary( conn: &impl Executor, + scope: &ObservationScopeV1, request: LcmSessionBoundaryRequest, ) -> Result { match compression_decision::boundary_transition_decision( @@ -165,7 +167,7 @@ pub async fn record_session_boundary( Ok(session_boundary_response(false, "not_compression_boundary")) } compression_decision::BoundaryTransitionDecision::CarryOver { old_session_id } => { - link_session_boundary(conn, &request, &old_session_id).await + link_session_boundary(conn, scope, &request, &old_session_id).await } compression_decision::BoundaryTransitionDecision::StartCooldown { boundary_skip_at } => { conn.execute( @@ -197,10 +199,11 @@ pub async fn record_session_boundary( /// lifecycle rows retain their original owner. async fn link_session_boundary( conn: &impl Executor, + scope: &ObservationScopeV1, request: &LcmSessionBoundaryRequest, old_session_id: &str, ) -> Result { - ensure_session(conn, &request.provider, &request.session_id).await?; + ensure_session(conn, scope, &request.provider, &request.session_id).await?; let old_state = lifecycle_state_or_default(conn, &request.provider, old_session_id, old_session_id).await?; // The link carries only frozen lifecycle coordinates. Authority rows keep @@ -367,6 +370,7 @@ fn canonical_replay_messages(raw_messages: &[LcmRawMessage]) -> Vec { #[hotpath::measure(label = "sessions.lcm.compress", future = true)] pub async fn compress( conn: &impl Executor, + scope: &ObservationScopeV1, publisher: &impl dag::LcmSummaryPublicationPort, storage_root: &Path, request: &LcmCompressionRequest, @@ -374,6 +378,7 @@ pub async fn compress( ) -> Result { let response = compress_inner( conn, + scope, publisher, storage_root, request, @@ -397,6 +402,7 @@ pub async fn compress( /// same lifecycle CAS without materializing a mega-session in one future. pub async fn compress_retained_page( conn: &impl Executor, + scope: &ObservationScopeV1, publisher: &impl dag::LcmSummaryPublicationPort, storage_root: &Path, request: &LcmCompressionRequest, @@ -405,6 +411,7 @@ pub async fn compress_retained_page( ) -> Result { let response = compress_inner( conn, + scope, publisher, storage_root, request, @@ -420,6 +427,7 @@ pub async fn compress_retained_page( async fn compress_inner( conn: &impl Executor, + scope: &ObservationScopeV1, publisher: &impl dag::LcmSummaryPublicationPort, storage_root: &Path, request: &LcmCompressionRequest, @@ -459,7 +467,7 @@ async fn compress_inner( }); } - ensure_session(conn, &request.provider, &request.session_id).await?; + ensure_session(conn, scope, &request.provider, &request.session_id).await?; let ingested = ingest_active_messages( conn, storage_root, @@ -2390,45 +2398,27 @@ async fn update_active_replay_metadata( Ok(()) } -/// Project fields on a session row LCM inserts only so foreign keys resolve. -/// Not a project identity. A later rollout projection overwrites them. -pub const LCM_UNKNOWN_PROJECT_KEY: &str = "unknown"; - -/// Historical placeholder `ensure_session` used to write as if it were a project. -/// Live stores still hold it; projection treats it as unscoped, same as -/// [`LCM_UNKNOWN_PROJECT_KEY`]. -pub const LCM_LEGACY_PLACEHOLDER_PROJECT_KEY: &str = "lcm-active-context"; - -/// Whether both project fields are an LCM placeholder rather than a project. -pub fn lcm_unscoped_session_project(project_key: &str, project_path: &str) -> bool { - project_key == project_path - && matches!( - project_key, - LCM_UNKNOWN_PROJECT_KEY | LCM_LEGACY_PLACEHOLDER_PROJECT_KEY | "" - ) -} - -/// Inserts the session row LCM foreign keys require without claiming a project. +/// Inserts the session row LCM foreign keys require. /// +/// `scope` is the session store's own scope, so the row carries the same +/// project fields the rollout projection writes for that store; the +/// projection's working directory later refines `project_path`. +/// `started_at` stays unset: LCM never observes the session start, and the +/// rollout's start must not lose to the time of this insert. /// `INSERT OR IGNORE` leaves a row the rollout projection already created. -/// Callers that run first store [`LCM_UNKNOWN_PROJECT_KEY`] in both project -/// fields so the projector can replace them with the real project. -pub async fn ensure_session( +async fn ensure_session( conn: &impl Executor, + scope: &ObservationScopeV1, provider: &str, session_id: &str, ) -> Result<(), LcmError> { + let (project_key, project_path) = session_project_fields(scope); conn.execute( "INSERT OR IGNORE INTO sessions ( - provider, session_id, project_key, project_path, started_at + provider, session_id, project_key, project_path ) - VALUES (?1, ?2, ?3, ?4, unixepoch())", - params![ - provider, - session_id, - LCM_UNKNOWN_PROJECT_KEY, - LCM_UNKNOWN_PROJECT_KEY, - ], + VALUES (?1, ?2, ?3, ?4)", + params![provider, session_id, project_key, project_path], ) .await?; Ok(()) diff --git a/crates/tracedecay-lcm/src/contracts.rs b/crates/tracedecay-lcm/src/contracts.rs index 1fd12b7d98..88697999b9 100644 --- a/crates/tracedecay-lcm/src/contracts.rs +++ b/crates/tracedecay-lcm/src/contracts.rs @@ -506,6 +506,8 @@ pub enum LcmError { actual_to: Option, }, LifecycleStateNotFound, + /// An LCM write targeted a store that holds no sessions. + NotASessionStore, Cancelled, DeadlineExceeded, BudgetExhausted, @@ -618,6 +620,9 @@ impl std::fmt::Display for LcmError { Self::LifecycleStateNotFound => { write!(f, "payload database error: lifecycle state not found") } + Self::NotASessionStore => { + write!(f, "LCM writes require a profile or project session store") + } Self::Cancelled => write!(f, "LCM payload verification was cancelled"), Self::DeadlineExceeded => { write!(f, "LCM payload verification deadline was exceeded") diff --git a/crates/tracedecay-lcm/src/schema.rs b/crates/tracedecay-lcm/src/schema.rs index 5017910bb5..a8c4c46c2a 100644 --- a/crates/tracedecay-lcm/src/schema.rs +++ b/crates/tracedecay-lcm/src/schema.rs @@ -24,9 +24,10 @@ use super::util; /// so a host record's parts share that row. There are no LCM summary /// tables. Every summary read joins the canonical `session_summary_nodes` / /// `session_summary_sources` authority (session temporal schema) through -/// [`SUMMARY_VISIBLE_SQL`]. Stores at an older version require a profile -/// reset. -pub const LCM_SCHEMA_VERSION: i64 = 13; +/// [`SUMMARY_VISIBLE_SQL`]. Session rows LCM creates carry their store's +/// scope as project fields (version 14); older stores hold rows with an +/// invented project and require a profile reset. +pub const LCM_SCHEMA_VERSION: i64 = 14; /// Visibility rule for every LCM summary read, over a `session_summary_nodes` /// row aliased `n`: a summary surfaces iff its availability in the session's @@ -789,7 +790,7 @@ mod tests { } #[tokio::test] - async fn shipped_v13_profile_remains_admissible_for_projector_rebuild() -> Result<(), String> { + async fn shipped_v13_profile_refuses_typed_without_upgrading() -> Result<(), String> { let temp = tempfile::tempdir().map_err(|error| error.to_string())?; let conn = TestConnection::open(&temp.path().join("sessions.db")); conn.execute_batch( @@ -804,11 +805,39 @@ mod tests { .await .map_err(|error| error.to_string())?; + let refusal = |error: LcmError| match error { + LcmError::ProfileResetRequired { + found_version, + required_version, + } => Ok((found_version, required_version)), + other => Err(format!("expected a typed reset refusal, got {other}")), + }; assert_eq!( - require_admissible_lcm_schema(&*conn) - .await - .map_err(|error| error.to_string())?, - LcmSchemaAdmission::Current + refusal( + require_admissible_lcm_schema(&*conn) + .await + .expect_err("a shipped v13 store must refuse admission") + )?, + (Some(13), 14) + ); + assert_eq!( + refusal( + ensure_lcm_schema(&conn) + .await + .expect_err("opening a shipped v13 store must refuse, not upgrade it") + )?, + (Some(13), 14) + ); + assert_eq!( + util::fetch_i64( + &*conn, + "SELECT version FROM session_schema_migrations WHERE name = 'lcm'", + (), + "migration marker version", + ) + .await + .map_err(|error| error.to_string())?, + 13 ); Ok(()) } diff --git a/crates/tracedecay-session-runtime/src/lcm_summary_convergence.rs b/crates/tracedecay-session-runtime/src/lcm_summary_convergence.rs index 2035b34f12..316ea75618 100644 --- a/crates/tracedecay-session-runtime/src/lcm_summary_convergence.rs +++ b/crates/tracedecay-session-runtime/src/lcm_summary_convergence.rs @@ -515,6 +515,7 @@ fn error_code(error: &LcmError) -> &'static str { LcmError::StaleRawProtectionSource { .. } => "stale_raw_protection_source", LcmError::StaleSummarySourceRange { .. } => "stale_summary_source_range", LcmError::LifecycleStateNotFound => "lifecycle_state_not_found", + LcmError::NotASessionStore => "not_a_session_store", LcmError::Cancelled => "cancelled", LcmError::DeadlineExceeded => "deadline_exceeded", LcmError::BudgetExhausted => "budget_exhausted", diff --git a/crates/tracedecay-sessions/src/runtime/store_access/lcm.rs b/crates/tracedecay-sessions/src/runtime/store_access/lcm.rs index b59412df65..eb82c9b477 100644 --- a/crates/tracedecay-sessions/src/runtime/store_access/lcm.rs +++ b/crates/tracedecay-sessions/src/runtime/store_access/lcm.rs @@ -1,5 +1,6 @@ use std::path::Path; +use tracedecay_domain::ObservationScopeV1; use tracedecay_runtime_core::db::DatabaseEngineReadSnapshot; use tracedecay_runtime_core::db::engine::{QueryExecutor, params}; @@ -294,6 +295,15 @@ impl<'a, D: SessionRegisteredDb + Sync> SessionStoreAccess<'a, D> { .await } + /// The session scope LCM writes this store's session rows under. + pub fn lcm_session_scope(&self) -> Result { + self.registered_binding() + .shard_id + .scope + .session_scope() + .ok_or(LcmError::NotASessionStore) + } + #[hotpath::skip] pub async fn lcm_session_boundary_guarded( &self, @@ -303,11 +313,12 @@ impl<'a, D: SessionRegisteredDb + Sync> SessionStoreAccess<'a, D> { where F: FnOnce() -> Result<(), LcmError>, { + let scope = self.lcm_session_scope()?; let transaction = self .begin_write_transaction() .await .map_err(|error| LcmError::Db(error.to_string()))?; - let response = compression::record_session_boundary(&transaction, request).await?; + let response = compression::record_session_boundary(&transaction, &scope, request).await?; before_commit()?; SessionWriteTxn::commit(transaction).await?; Ok(response) diff --git a/crates/tracedecay-store/src/canonical_projection.rs b/crates/tracedecay-store/src/canonical_projection.rs index f8ba42449e..190cf759bf 100644 --- a/crates/tracedecay-store/src/canonical_projection.rs +++ b/crates/tracedecay-store/src/canonical_projection.rs @@ -51,6 +51,19 @@ fn rendering_message_semantics( } } +/// The `(project_key, project_path)` a session row in `scope` carries before +/// any observation names its working directory. Every writer of a session row +/// uses this, so the projection's later cwd only refines the path. +pub fn session_project_fields(scope: &ObservationScopeV1) -> (String, String) { + match scope { + ObservationScopeV1::Profile => ("user".to_owned(), "user".to_owned()), + ObservationScopeV1::Project { project_id } => ( + project_id.as_str().to_owned(), + project_id.as_str().to_owned(), + ), + } +} + pub fn derive_canonical_projection( observation: &DurableObservationV1, ) -> ProjectionStoreResult { @@ -119,13 +132,7 @@ fn derive_canonical_projection_for( } let provider = envelope.provider().as_str().to_owned(); let session_id = envelope.relations().session_id().as_str().to_owned(); - let (project_key, fallback_project_path) = match observation.scope() { - ObservationScopeV1::Profile => ("user".to_owned(), "user".to_owned()), - ObservationScopeV1::Project { project_id } => ( - project_id.as_str().to_owned(), - project_id.as_str().to_owned(), - ), - }; + let (project_key, fallback_project_path) = session_project_fields(observation.scope()); let timestamp = projected .as_ref() .and_then(|projected| projected.timestamp) diff --git a/crates/tracedecay-store/src/lib.rs b/crates/tracedecay-store/src/lib.rs index 78d7ebc08c..a7194263dd 100644 --- a/crates/tracedecay-store/src/lib.rs +++ b/crates/tracedecay-store/src/lib.rs @@ -33,7 +33,7 @@ pub mod transcript; pub use canonical_projection::{ EDITED_FILES_KEY, SPAWNED_SESSIONS_KEY, TOOL_USE_ID_KEY, canonical_fact_text, - derive_canonical_projection, message_metadata_with_envelope, + derive_canonical_projection, message_metadata_with_envelope, session_project_fields, stored_message_is_shipped_release_rendering, workflow_semantic_kind, }; pub use codex_goal_context::{ diff --git a/crates/tracedecay-store/src/runtime/identity.rs b/crates/tracedecay-store/src/runtime/identity.rs index f5643fe260..7aa4085985 100644 --- a/crates/tracedecay-store/src/runtime/identity.rs +++ b/crates/tracedecay-store/src/runtime/identity.rs @@ -4,6 +4,7 @@ use std::sync::Arc; use serde::{Deserialize, Deserializer, Serialize}; use sha2::{Digest, Sha256}; +use tracedecay_domain::ObservationScopeV1; use tracedecay_domain::canonical_text::{encode_tagged_lowercase_hex, is_canonical_text}; pub use tracedecay_domain::{ AuthorityEpoch, BrainId, BrainNodeId, LocatorDigest, ProjectId, RefId, RepositoryId, @@ -173,6 +174,22 @@ impl StoreShardScopeV1 { } } + /// The observation scope whose sessions this shard stores; `None` for a + /// shard that holds no sessions. + pub fn session_scope(&self) -> Option { + match self { + Self::ProfileSessions => Some(ObservationScopeV1::Profile), + Self::ProjectSessions { project_id } => Some(ObservationScopeV1::Project { + project_id: project_id.clone(), + }), + Self::Profile + | Self::ProfileMemory + | Self::RemoteNode { .. } + | Self::Project { .. } + | Self::Code { .. } => None, + } + } + pub fn is_mutable(&self) -> bool { match self { Self::Profile @@ -583,6 +600,44 @@ mod tests { ); } + #[test] + fn only_session_shards_name_a_session_scope() { + let project_id = id::("project.identity"); + assert_eq!( + StoreShardScopeV1::ProfileSessions.session_scope(), + Some(ObservationScopeV1::Profile) + ); + assert_eq!( + StoreShardScopeV1::ProjectSessions { + project_id: project_id.clone(), + } + .session_scope(), + Some(ObservationScopeV1::Project { + project_id: project_id.clone(), + }) + ); + + for scope in [ + StoreShardScopeV1::Profile, + StoreShardScopeV1::ProfileMemory, + StoreShardScopeV1::RemoteNode { + node_id: id::("node.identity"), + }, + StoreShardScopeV1::Project { + project_id: project_id.clone(), + }, + StoreShardScopeV1::Code { + project_id: project_id.clone(), + repository_id: id::("repository.identity"), + scope: CodeShardScopeV1::Worktree { + worktree_id: id::("worktree.identity"), + }, + }, + ] { + assert_eq!(scope.session_scope(), None, "{scope:?} holds no sessions"); + } + } + #[test] fn canonical_locator_digest_binds_the_exact_absolute_path() { let fixture = host_temp_root("tracedecay-store-locator-fixture"); diff --git a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs index 7001fa8c95..9925445dc8 100644 --- a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs +++ b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs @@ -635,19 +635,19 @@ fn stale_session_stores_refuse_sessions_only_until_their_scoped_reset() { } #[test] -fn project_session_store_at_another_lcm_schema_version_refuses_sessions_only() { +fn session_stores_at_shipped_lcm_schema_13_refuse_sessions_only_until_their_scoped_reset() { refused_session_stores_serve_code_until_their_scoped_reset(&SessionStoreRefusal { age: |db| { execute_once( db, - "UPDATE session_schema_migrations SET version = 12 WHERE name = 'lcm'", + "UPDATE session_schema_migrations SET version = 13 WHERE name = 'lcm'", ); }, - ages_profile_store: false, + ages_profile_store: true, authority: "LCM", - found_version: json!(12), - required_version: json!(13), - reason: "LCM profile schema 12 is incompatible with required schema 13; reset the profile", + found_version: json!(13), + required_version: json!(14), + reason: "LCM profile schema 13 is incompatible with required schema 14; reset the profile", session_tool: "tracedecay_lcm_grep", session_tool_args: || json!({ "query": "probe", "format": "json" }), });