Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
80 changes: 34 additions & 46 deletions crates/tracedecay-global-db/src/observation_projection/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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());
}
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
// `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.
Expand Down Expand Up @@ -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<T: Clone + Eq>(
field: &'static str,
actual: Option<&T>,
Expand Down Expand Up @@ -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::{
Expand Down Expand Up @@ -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]
Expand Down
4 changes: 4 additions & 0 deletions crates/tracedecay-global-db/src/registered_lcm.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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!(
Expand Down Expand Up @@ -288,6 +289,7 @@ impl RegisteredGlobalDb {
);
let mut response = compression::compress(
&transaction,
&scope,
&publisher,
storage_root,
request,
Expand Down Expand Up @@ -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!(
Expand Down Expand Up @@ -369,6 +372,7 @@ impl RegisteredGlobalDb {
);
let bounded = compression::compress_retained_page(
&transaction,
&scope,
&publisher,
storage_root,
request,
Expand Down
188 changes: 161 additions & 27 deletions crates/tracedecay-global-db/src/session_project_rebind_tests.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -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,
Expand Down Expand Up @@ -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");
Expand All @@ -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
Expand Down Expand Up @@ -383,24 +443,98 @@ 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(),
"catch-up must finish both observations"
);
}

#[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();
Expand Down
Loading
Loading