diff --git a/crates/tracedecay-sessions/src/runtime/lcm/doctor.rs b/crates/tracedecay-sessions/src/runtime/lcm/doctor.rs index 8c524b277..909bad9f3 100644 --- a/crates/tracedecay-sessions/src/runtime/lcm/doctor.rs +++ b/crates/tracedecay-sessions/src/runtime/lcm/doctor.rs @@ -14,6 +14,7 @@ use super::{ }; const MAX_SAMPLES: usize = 20; +const RETENTION_SCAN_MESSAGE_LIMIT: usize = 10_000; const RETENTION_OLD_DAYS: f64 = 30.0; const RETENTION_HEAVY_CHARS: i64 = 128 * 1024; const SQLITE_IN_BATCH_SIZE: usize = 500; @@ -46,15 +47,21 @@ struct RepairRequest<'a> { } pub async fn doctor(conn: &Connection, request: DoctorRequest<'_>) -> Result { - let diagnostics = gather_diagnostics( - conn, - request.storage_root, - request.provider, - request.session_id, - &request.clean_config, - &request.gc_config, - ) - .await?; + let diagnostics = if request.mode == "retention" { + json!({ + "retention": retention_candidates_bounded(conn, request.provider, request.session_id).await?, + }) + } else { + gather_diagnostics( + conn, + request.storage_root, + request.provider, + request.session_id, + &request.clean_config, + &request.gc_config, + ) + .await? + }; let repairs = plan_and_apply_repairs( conn, RepairRequest { @@ -70,7 +77,11 @@ pub async fn doctor(conn: &Connection, request: DoctorRequest<'_>) -> Result, +) -> Result { + let now = current_timestamp(); + let discovery = retention_session_ids(conn, provider, session_id).await?; + let sampled_session_ids = discovery.session_ids; + let per_session_limit = if sampled_session_ids.is_empty() { + 0 + } else { + RETENTION_SCAN_MESSAGE_LIMIT / sampled_session_ids.len() + }; + let mut sessions = BTreeMap::::new(); + let mut messages_analyzed = 0_usize; + let mut incomplete_sessions = Vec::new(); + for sampled_session_id in &sampled_session_ids { + let mut rows = conn + .query( + "SELECT LENGTH(index_text), COALESCE(timestamp, 0) + FROM lcm_raw_messages + WHERE provider = ?1 AND session_id = ?2 + ORDER BY store_id + LIMIT ?3", + params![ + provider, + sampled_session_id.as_str(), + (per_session_limit + 1) as i64 + ], + ) + .await?; + let mut stats = RetentionSessionStats::default(); + while let Some(row) = rows.next().await? { + if stats.message_count == per_session_limit as i64 { + incomplete_sessions.push(sampled_session_id.clone()); + break; + } + let retained_chars: i64 = row.get(0)?; + let timestamp: i64 = row.get(1)?; + stats.message_count += 1; + stats.retained_chars += retained_chars; + stats.first_message_at = Some( + stats + .first_message_at + .map_or(timestamp, |first| first.min(timestamp)), + ); + stats.last_message_at = Some( + stats + .last_message_at + .map_or(timestamp, |last| last.max(timestamp)), + ); + messages_analyzed += 1; + } + if stats.message_count > 0 { + sessions.insert(sampled_session_id.clone(), stats); + } + } + + let complete_session_ids = sessions + .keys() + .filter(|session_id| !incomplete_sessions.contains(session_id)) + .cloned() + .collect::>(); + let sessions_with_summaries = + retention_sessions_with_summaries(conn, provider, &complete_session_ids).await?; + let mut sessions = sessions.into_iter().collect::>(); + sessions.sort_by(|(left_id, left), (right_id, right)| { + right + .retained_chars + .cmp(&left.retained_chars) + .then_with(|| left.last_message_at.cmp(&right.last_message_at)) + .then_with(|| left_id.cmp(right_id)) + }); + let sessions_analyzed = sessions.len(); + let complete_sessions_analyzed = complete_session_ids.len(); + let mut candidates = Vec::new(); + for (session_id, stats) in sessions { + if incomplete_sessions.contains(&session_id) { + continue; + } + let first_message_at = stats.first_message_at.unwrap_or_default(); + let last_message_at = stats.last_message_at.unwrap_or_default(); + let has_summary = sessions_with_summaries.contains(&session_id); + let age_days = if last_message_at > 0 { + (now.saturating_sub(last_message_at) as f64) / 86_400.0 + } else { + 0.0 + }; + let old = age_days >= RETENTION_OLD_DAYS; + let heavy = stats.retained_chars >= RETENTION_HEAVY_CHARS; + let session_only = !has_summary; + if old || heavy || session_only { + candidates.push(json!({ + "session_id": session_id, + "message_count": stats.message_count, + "retained_chars": stats.retained_chars, + "first_message_at": first_message_at, + "last_message_at": last_message_at, + "age_days": age_days, + "old": old, + "heavy": heavy, + "session_only": session_only, + "protected": false, + })); + } + if candidates.len() >= MAX_SAMPLES { + break; + } + } + + Ok(json!({ + "read_only": true, + "sessions_analyzed": sessions_analyzed, + "complete_sessions_analyzed": complete_sessions_analyzed, + "messages_analyzed": messages_analyzed, + "scan_limit": RETENTION_SCAN_MESSAGE_LIMIT, + "per_session_scan_limit": per_session_limit, + "scan_truncated": discovery.truncated || !incomplete_sessions.is_empty(), + "session_scan_truncated": discovery.truncated, + "session_discovery_rows_scanned": discovery.rows_scanned, + "session_discovery_scan_limit": if session_id.is_some() { 0 } else { MAX_SAMPLES + 1 }, + "incomplete_session_count": incomplete_sessions.len(), + "incomplete_sessions": incomplete_sessions, + "candidate_count": candidates.len(), + "candidates": candidates, + })) +} + +struct RetentionSessionDiscovery { + session_ids: Vec, + truncated: bool, + rows_scanned: usize, +} + +async fn retention_session_ids( + conn: &Connection, + provider: &str, + session_id: Option<&str>, +) -> Result { + if let Some(session_id) = session_id { + return Ok(RetentionSessionDiscovery { + session_ids: vec![session_id.to_string()], + truncated: false, + rows_scanned: 0, + }); + } + // Session discovery must stay on the narrow provider/session registry; + // reading distinct ids from lcm_raw_messages can visit its entire payload-heavy store. + let sql = "SELECT session_id + FROM sessions + WHERE provider = ? + ORDER BY session_id + LIMIT ?"; + let values = vec![ + SqlValue::Text(provider.to_string()), + SqlValue::Integer((MAX_SAMPLES + 1) as i64), + ]; + let mut rows = conn.query(sql, values).await?; + let mut session_ids = Vec::new(); + while let Some(row) = rows.next().await? { + session_ids.push(row.get(0)?); + } + let rows_scanned = session_ids.len(); + let truncated = rows_scanned > MAX_SAMPLES; + session_ids.truncate(MAX_SAMPLES); + Ok(RetentionSessionDiscovery { + session_ids, + truncated, + rows_scanned, + }) +} + +#[derive(Default)] +struct RetentionSessionStats { + message_count: i64, + retained_chars: i64, + first_message_at: Option, + last_message_at: Option, +} + +async fn retention_sessions_with_summaries( + conn: &Connection, + provider: &str, + session_ids: &[String], +) -> Result, LcmError> { + let mut sessions_with_summaries = BTreeSet::new(); + for session_id in session_ids { + let mut rows = conn + .query( + "SELECT 1 + FROM lcm_summary_nodes + WHERE provider = ?1 AND session_id = ?2 + LIMIT 1", + params![provider, session_id.as_str()], + ) + .await?; + if rows.next().await?.is_some() { + sessions_with_summaries.insert(session_id.clone()); + } + } + Ok(sessions_with_summaries) +} + #[derive(Default)] struct CleanupSessionCandidate { classes: BTreeSet<&'static str>, @@ -1393,6 +1607,406 @@ mod tests { .map_err(|err| format!("read raw message count for {session_id}: {err}")) } + #[tokio::test] + async fn retention_mode_skips_unrelated_diagnostics() -> Result<(), String> { + let temp = tempfile::tempdir().map_err(|err| format!("create tempdir: {err}"))?; + let db_path = temp.path().join("sessions.db"); + let db = libsql::Builder::new_local(&db_path) + .build() + .await + .map_err(|err| format!("build test database: {err}"))?; + let conn = db + .connect() + .map_err(|err| format!("connect to test database: {err}"))?; + conn.execute_batch( + "CREATE TABLE sessions ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL, + PRIMARY KEY(provider, session_id) + ); + CREATE TABLE lcm_raw_messages ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL, + store_id INTEGER PRIMARY KEY AUTOINCREMENT, + index_text TEXT NOT NULL, + timestamp INTEGER + ); + CREATE TABLE lcm_summary_nodes ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL + ); + INSERT INTO sessions (provider, session_id) + VALUES ('codex', 'retention-only-session'); + INSERT INTO lcm_raw_messages (provider, session_id, index_text, timestamp) + VALUES ('codex', 'retention-only-session', 'retention body', 1);", + ) + .await + .map_err(|err| format!("create retention-only fixture: {err}"))?; + + let result = doctor( + &conn, + DoctorRequest { + storage_root: temp.path(), + db_path: &db_path, + provider: "codex", + session_id: None, + mode: "retention", + apply: false, + clean_config: LcmCleanConfig::default(), + gc_config: LcmGcConfig::default(), + }, + ) + .await + .map_err(|err| format!("run retention doctor: {err}"))?; + + assert_eq!(result["status"], "ok"); + assert_eq!(result["mode"], "retention"); + assert_eq!(result["dry_run"], false); + assert_eq!(result["diagnostics"]["retention"]["read_only"], true); + assert_eq!(result["diagnostics"]["retention"]["candidate_count"], 1); + assert_eq!( + result["diagnostics"].as_object().map(serde_json::Map::len), + Some(1), + "retention mode must not run or report unrelated diagnostics" + ); + assert_eq!(result["repairs"]["planned_actions"], json!([])); + + let mut rows = conn + .query("SELECT COUNT(*) FROM lcm_raw_messages", ()) + .await + .map_err(|err| format!("count retained messages: {err}"))?; + let row = rows + .next() + .await + .map_err(|err| format!("read retained message count: {err}"))? + .ok_or_else(|| "missing retained message count row".to_string())?; + let message_count: i64 = row + .get(0) + .map_err(|err| format!("decode retained message count: {err}"))?; + assert_eq!(message_count, 1); + Ok(()) + } + + #[tokio::test] + async fn diagnose_retention_includes_candidate_beyond_bounded_session_sample() + -> Result<(), String> { + let temp = tempfile::tempdir().map_err(|err| format!("create tempdir: {err}"))?; + let db_path = temp.path().join("sessions.db"); + let db = libsql::Builder::new_local(&db_path) + .build() + .await + .map_err(|err| format!("build test database: {err}"))?; + let conn = db + .connect() + .map_err(|err| format!("connect to test database: {err}"))?; + conn.execute_batch( + "CREATE TABLE sessions ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL, + project_key TEXT NOT NULL, + project_path TEXT NOT NULL, + title TEXT, + started_at INTEGER, + PRIMARY KEY(provider, session_id) + ); + CREATE TABLE session_messages ( + provider TEXT NOT NULL, + message_id TEXT NOT NULL, + session_id TEXT NOT NULL, + role TEXT NOT NULL, + timestamp INTEGER, + ordinal INTEGER NOT NULL, + text TEXT NOT NULL, + metadata_json TEXT, + PRIMARY KEY(provider, message_id) + );", + ) + .await + .map_err(|err| format!("create test session schema: {err}"))?; + schema::ensure_lcm_schema(&conn) + .await + .map_err(|err| format!("create test LCM schema: {err}"))?; + + for ordinal in 1..=MAX_SAMPLES + 1 { + let session_id = format!("session-{ordinal:02}"); + let message_id = format!("message-{ordinal:02}"); + insert_test_clean_candidate(&conn, temp.path(), &session_id, &message_id).await?; + } + conn.execute( + "UPDATE lcm_raw_messages + SET index_text = printf('%.*c', ?1, 'x') + WHERE provider = 'cursor' AND session_id = 'session-21'", + params![RETENTION_HEAVY_CHARS + 1], + ) + .await + .map_err(|err| format!("make late retention candidate largest: {err}"))?; + + let result = doctor( + &conn, + DoctorRequest { + storage_root: temp.path(), + db_path: &db_path, + provider: "cursor", + session_id: None, + mode: "diagnose", + apply: false, + clean_config: LcmCleanConfig::default(), + gc_config: LcmGcConfig::default(), + }, + ) + .await + .map_err(|err| format!("run diagnose doctor: {err}"))?; + + assert!( + result["diagnostics"]["retention"]["candidates"] + .as_array() + .is_some_and(|candidates| candidates + .iter() + .any(|candidate| candidate["session_id"] == "session-21")), + "non-retention diagnostics must retain the largest candidate beyond the bounded first-{MAX_SAMPLES} sessions" + ); + Ok(()) + } + + #[tokio::test] + async fn retention_mode_bounds_raw_message_scan_and_reports_truncation() -> Result<(), String> { + let temp = tempfile::tempdir().map_err(|err| format!("create tempdir: {err}"))?; + let db_path = temp.path().join("sessions.db"); + let db = libsql::Builder::new_local(&db_path) + .build() + .await + .map_err(|err| format!("build test database: {err}"))?; + let conn = db + .connect() + .map_err(|err| format!("connect to test database: {err}"))?; + conn.execute_batch(&format!( + "CREATE TABLE sessions ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL, + PRIMARY KEY(provider, session_id) + ); + CREATE TABLE lcm_raw_messages ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL, + store_id INTEGER PRIMARY KEY AUTOINCREMENT, + index_text TEXT NOT NULL, + timestamp INTEGER + ); + CREATE TABLE lcm_summary_nodes ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL + ); + INSERT INTO sessions (provider, session_id) + VALUES ('codex', 'bounded-retention-session'); + WITH RECURSIVE messages(ordinal) AS ( + VALUES(1) + UNION ALL + SELECT ordinal + 1 FROM messages + WHERE ordinal <= {RETENTION_SCAN_MESSAGE_LIMIT} + ) + INSERT INTO lcm_raw_messages (provider, session_id, index_text, timestamp) + SELECT 'codex', 'bounded-retention-session', 'x', ordinal + FROM messages;" + )) + .await + .map_err(|err| format!("create oversized retention fixture: {err}"))?; + + let result = doctor( + &conn, + DoctorRequest { + storage_root: temp.path(), + db_path: &db_path, + provider: "codex", + session_id: Some("bounded-retention-session"), + mode: "retention", + apply: false, + clean_config: LcmCleanConfig::default(), + gc_config: LcmGcConfig::default(), + }, + ) + .await + .map_err(|err| format!("run bounded retention doctor: {err}"))?; + + let retention = &result["diagnostics"]["retention"]; + assert_eq!(retention["messages_analyzed"], RETENTION_SCAN_MESSAGE_LIMIT); + assert_eq!(retention["scan_limit"], RETENTION_SCAN_MESSAGE_LIMIT); + assert_eq!(retention["scan_truncated"], true); + assert_eq!(retention["sessions_analyzed"], 1); + assert_eq!(retention["candidate_count"], 0); + assert_eq!(retention["complete_sessions_analyzed"], 0); + assert_eq!(retention["incomplete_session_count"], 1); + assert_eq!( + retention["incomplete_sessions"], + json!(["bounded-retention-session"]) + ); + Ok(()) + } + + #[tokio::test] + async fn retention_mode_samples_multiple_sessions_deterministically() -> Result<(), String> { + let temp = tempfile::tempdir().map_err(|err| format!("create tempdir: {err}"))?; + let db_path = temp.path().join("sessions.db"); + let db = libsql::Builder::new_local(&db_path) + .build() + .await + .map_err(|err| format!("build test database: {err}"))?; + let conn = db + .connect() + .map_err(|err| format!("connect to test database: {err}"))?; + conn.execute_batch(&format!( + "CREATE TABLE sessions ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL, + PRIMARY KEY(provider, session_id) + ); + CREATE TABLE lcm_raw_messages ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL, + store_id INTEGER PRIMARY KEY AUTOINCREMENT, + index_text TEXT NOT NULL, + timestamp INTEGER + ); + CREATE TABLE lcm_summary_nodes ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL + ); + INSERT INTO sessions (provider, session_id) + VALUES ('codex', 'a-large-session'), ('codex', 'z-small-session'); + WITH RECURSIVE messages(ordinal) AS ( + VALUES(1) + UNION ALL + SELECT ordinal + 1 FROM messages + WHERE ordinal < {RETENTION_SCAN_MESSAGE_LIMIT} + ) + INSERT INTO lcm_raw_messages (provider, session_id, index_text, timestamp) + SELECT 'codex', 'a-large-session', 'x', ordinal + FROM messages; + INSERT INTO lcm_raw_messages (provider, session_id, index_text, timestamp) + VALUES ('codex', 'z-small-session', 'complete', 1);" + )) + .await + .map_err(|err| format!("create multi-session retention fixture: {err}"))?; + + let first = retention_candidates_bounded(&conn, "codex", None) + .await + .map_err(|err| format!("first retention scan: {err}"))?; + let second = retention_candidates_bounded(&conn, "codex", None) + .await + .map_err(|err| format!("second retention scan: {err}"))?; + + for field in [ + "sessions_analyzed", + "complete_sessions_analyzed", + "messages_analyzed", + "scan_limit", + "per_session_scan_limit", + "scan_truncated", + "session_scan_truncated", + "incomplete_session_count", + "incomplete_sessions", + "candidate_count", + ] { + assert_eq!(first[field], second[field], "{field} must be deterministic"); + } + let candidate_session_ids = |retention: &Value| { + retention["candidates"] + .as_array() + .into_iter() + .flatten() + .map(|candidate| candidate["session_id"].clone()) + .collect::>() + }; + assert_eq!( + candidate_session_ids(&first), + candidate_session_ids(&second), + "bounded candidate sampling must be deterministic" + ); + assert_eq!(first["sessions_analyzed"], 2); + assert_eq!(first["complete_sessions_analyzed"], 1); + assert_eq!(first["incomplete_session_count"], 1); + assert_eq!(first["incomplete_sessions"], json!(["a-large-session"])); + assert!( + first["candidates"] + .as_array() + .unwrap() + .iter() + .any(|candidate| candidate["session_id"] == "z-small-session"), + "a large first session must not starve a later complete session" + ); + Ok(()) + } + + #[tokio::test] + async fn retention_mode_bounds_session_discovery_before_large_raw_sessions() + -> Result<(), String> { + let temp = tempfile::tempdir().map_err(|err| format!("create tempdir: {err}"))?; + let db_path = temp.path().join("sessions.db"); + let db = libsql::Builder::new_local(&db_path) + .build() + .await + .map_err(|err| format!("build test database: {err}"))?; + let conn = db + .connect() + .map_err(|err| format!("connect to test database: {err}"))?; + conn.execute_batch(&format!( + "CREATE TABLE sessions ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL, + PRIMARY KEY(provider, session_id) + ); + CREATE TABLE lcm_raw_messages ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL, + store_id INTEGER PRIMARY KEY AUTOINCREMENT, + index_text TEXT NOT NULL, + timestamp INTEGER + ); + CREATE INDEX idx_lcm_raw_session_order + ON lcm_raw_messages(provider, session_id, store_id); + CREATE TABLE lcm_summary_nodes ( + provider TEXT NOT NULL, + session_id TEXT NOT NULL + ); + WITH RECURSIVE session_numbers(ordinal) AS ( + VALUES(1) + UNION ALL + SELECT ordinal + 1 FROM session_numbers WHERE ordinal < 21 + ) + INSERT INTO sessions (provider, session_id) + SELECT 'codex', printf('session-%02d', ordinal) FROM session_numbers; + WITH RECURSIVE messages(ordinal) AS ( + VALUES(1) + UNION ALL + SELECT ordinal + 1 FROM messages + WHERE ordinal <= {RETENTION_SCAN_MESSAGE_LIMIT} + ) + INSERT INTO lcm_raw_messages (provider, session_id, index_text, timestamp) + SELECT 'codex', 'session-01', 'x', ordinal FROM messages; + INSERT INTO lcm_raw_messages (provider, session_id, index_text, timestamp) + SELECT 'codex', 'session-02', index_text, timestamp + FROM lcm_raw_messages WHERE session_id = 'session-01'; + INSERT INTO lcm_raw_messages (provider, session_id, index_text, timestamp) + SELECT 'codex', 'session-03', index_text, timestamp + FROM lcm_raw_messages WHERE session_id = 'session-01';" + )) + .await + .map_err(|err| format!("create bounded-discovery fixture: {err}"))?; + + let retention = retention_candidates_bounded(&conn, "codex", None) + .await + .map_err(|err| format!("run bounded retention scan: {err}"))?; + + assert_eq!(retention["session_discovery_rows_scanned"], MAX_SAMPLES + 1); + assert_eq!(retention["session_discovery_scan_limit"], MAX_SAMPLES + 1); + assert_eq!(retention["session_scan_truncated"], true); + assert!( + retention["messages_analyzed"].as_u64().unwrap_or_default() + <= RETENTION_SCAN_MESSAGE_LIMIT as u64, + "raw-message analysis must retain its total hard bound" + ); + Ok(()) + } + #[tokio::test] async fn clean_apply_backup_callback_runs_under_immediate_transaction() -> Result<(), String> { let temp = tempfile::tempdir().map_err(|err| format!("create tempdir: {err}"))?;