use super::*;
use rusqlite::Connection;
#[path = "tests/since.rs"]
mod since;
fn setup_conn() -> Connection {
let conn = Connection::open_in_memory().unwrap();
crate::migrate::run_migrations(&conn).unwrap();
conn
}
struct TempRoot {
path: PathBuf,
}
impl TempRoot {
fn new(name: &str) -> Self {
let path = std::env::temp_dir().join(format!(
"remem-ingest-{name}-{}-{}",
std::process::id(),
chrono::Utc::now().timestamp_nanos_opt().unwrap_or_default()
));
std::fs::create_dir_all(&path).unwrap();
Self { path }
}
fn scan_root(&self, label: &str) -> ScanRoot {
ScanRoot {
label: label.to_string(),
path: self.path.clone(),
required: true,
}
}
fn write(&self, relative: &str, content: &str) -> PathBuf {
let path = self.path.join(relative);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(&path, content).unwrap();
path
}
}
impl Drop for TempRoot {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.path);
}
}
fn claude_line(cwd: &str, role: &str, text: &str) -> String {
format!(
r#"{{"type":"{role}","cwd":"{cwd}","gitBranch":"main","message":{{"content":[{{"type":"text","text":"{text}"}}]}}}}"#
)
}
fn raw_message_count(conn: &Connection) -> i64 {
conn.query_row("SELECT COUNT(*) FROM raw_messages", [], |row| row.get(0))
.unwrap()
}
fn cursor_count(conn: &Connection) -> i64 {
conn.query_row("SELECT COUNT(*) FROM ingest_cursors", [], |row| row.get(0))
.unwrap()
}
fn run(conn: &Connection, roots: &[ScanRoot]) -> IngestSummary {
run_ingest_sessions(conn, roots, &IngestOptions::default()).unwrap()
}
#[test]
fn double_run_is_idempotent_and_second_run_skips_via_cursor() {
let conn = setup_conn();
let root = TempRoot::new("idempotent");
let cwd = root.path.to_string_lossy().to_string();
root.write(
"proj-a/session-1.jsonl",
&format!(
"{}\n{}\n",
claude_line(&cwd, "user", "first question"),
claude_line(&cwd, "assistant", "first answer")
),
);
let roots = [root.scan_root("local")];
let first = run(&conn, &roots);
assert_eq!(first.scanned, 1);
assert_eq!(first.skipped, 0);
assert_eq!(first.ingested_messages, 2);
assert_eq!(first.failed_files, 0);
assert_eq!(first.exit_code(), 0);
assert_eq!(raw_message_count(&conn), 2);
assert_eq!(cursor_count(&conn), 1);
let second = run(&conn, &roots);
assert_eq!(second.scanned, 1);
assert_eq!(second.skipped, 1, "unchanged cursor must skip the file");
assert_eq!(second.ingested_messages, 0);
assert_eq!(raw_message_count(&conn), 2, "no new rows on second run");
}
#[test]
fn changed_file_is_redrained_and_unique_constraint_dedupes() {
let conn = setup_conn();
let root = TempRoot::new("redrain");
let cwd = root.path.to_string_lossy().to_string();
let file = root.write(
"proj-a/session-1.jsonl",
&format!("{}\n", claude_line(&cwd, "user", "first question")),
);
let roots = [root.scan_root("local")];
let first = run(&conn, &roots);
assert_eq!(first.ingested_messages, 1);
let mut content = std::fs::read_to_string(&file).unwrap();
content.push_str(&claude_line(&cwd, "assistant", "late answer"));
content.push('\n');
std::fs::write(&file, content).unwrap();
let second = run(&conn, &roots);
assert_eq!(second.skipped, 0);
assert_eq!(second.ingested_messages, 1, "only appended message is new");
assert_eq!(raw_message_count(&conn), 2);
}
#[test]
fn corrupt_file_is_isolated_and_batch_continues() {
let conn = setup_conn();
let root = TempRoot::new("corrupt");
let cwd = root.path.to_string_lossy().to_string();
root.write(
"proj-a/broken.jsonl",
"not json at all\n{\"also\": broken\n",
);
root.write(
"proj-a/good.jsonl",
&format!("{}\n", claude_line(&cwd, "user", "healthy message")),
);
let roots = [root.scan_root("local")];
let summary = run(&conn, &roots);
assert_eq!(summary.scanned, 2);
assert_eq!(summary.failed_files, 1);
assert_eq!(summary.ingested_messages, 1, "good file still ingests");
assert_eq!(summary.exit_code(), 1, "partial failure exits non-zero");
let failures: i64 = conn
.query_row("SELECT COUNT(*) FROM raw_ingest_failures", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(failures, 1, "corrupt file recorded in raw_ingest_failures");
assert_eq!(cursor_count(&conn), 1);
}
#[test]
fn hook_drained_transcript_then_batch_produces_no_duplicates() {
let conn = setup_conn();
let root = TempRoot::new("hook-dedup");
let cwd = root.path.to_string_lossy().to_string();
let file = root.write(
"proj-a/session-7.jsonl",
&format!(
"{}\n{}\n",
claude_line(&cwd, "user", "hook question"),
claude_line(&cwd, "assistant", "hook answer")
),
);
let project = crate::project_id::project_from_cwd(&cwd);
let hook_report = raw_archive::drain_transcript(
&conn,
&file.to_string_lossy(),
"session-7",
&project,
Some("main"),
Some(&cwd),
)
.unwrap();
assert_eq!(hook_report.inserted, 2);
assert_eq!(raw_message_count(&conn), 2);
let summary = run(&conn, &[root.scan_root("local")]);
assert_eq!(summary.ingested_messages, 0, "batch adds no duplicate rows");
assert_eq!(raw_message_count(&conn), 2);
}
#[test]
fn active_partial_tail_is_not_a_failure_and_keeps_cursor_behind() {
let conn = setup_conn();
let root = TempRoot::new("partial-tail");
let cwd = root.path.to_string_lossy().to_string();
root.write(
"proj-a/live.jsonl",
&format!(
"{}\n{}",
claude_line(&cwd, "user", "complete line"),
r#"{"type":"assistant","message":{"content":[{"type":"te"#
),
);
let roots = [root.scan_root("local")];
let summary = run(&conn, &roots);
assert_eq!(summary.failed_files, 0, "partial tail is not a failure");
assert_eq!(summary.partial_files, 1);
assert_eq!(summary.ingested_messages, 1);
assert_eq!(summary.exit_code(), 0);
assert_eq!(cursor_count(&conn), 0, "cursor must not advance");
let failures: i64 = conn
.query_row("SELECT COUNT(*) FROM raw_ingest_failures", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(failures, 0, "partial tail is not recorded as a failure");
let second = run(&conn, &roots);
assert_eq!(second.skipped, 0);
assert_eq!(second.ingested_messages, 0);
}
#[test]
fn discovery_excludes_subagents_and_non_jsonl_files() {
let conn = setup_conn();
let root = TempRoot::new("discovery");
let cwd = root.path.to_string_lossy().to_string();
root.write(
"proj-a/keep.jsonl",
&format!("{}\n", claude_line(&cwd, "user", "kept message")),
);
root.write(
"proj-a/subagents/skip.jsonl",
&format!("{}\n", claude_line(&cwd, "user", "subagent message")),
);
root.write("proj-a/notes.txt", "not a transcript\n");
let summary = run(&conn, &[root.scan_root("local")]);
assert_eq!(summary.scanned, 1, "only the top-level jsonl is discovered");
assert_eq!(summary.ingested_messages, 1);
}
#[test]
fn source_root_label_is_stored_and_distinguishes_roots() {
let conn = setup_conn();
let root_a = TempRoot::new("root-a");
let root_b = TempRoot::new("root-b");
let line =
r#"{"type":"user","message":{"content":[{"type":"text","text":"same project name"}]}}"#;
root_a.write("proj-x/session-same.jsonl", &format!("{line}\n"));
root_b.write("proj-x/session-same.jsonl", &format!("{line}\n"));
let summary = run(
&conn,
&[root_a.scan_root("local"), root_b.scan_root("starlight")],
);
assert_eq!(summary.ingested_messages, 2);
let labels: Vec<(String, String)> = {
let mut stmt = conn
.prepare("SELECT source_root, project FROM raw_messages ORDER BY source_root")
.unwrap();
let rows = stmt
.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
.unwrap();
rows.collect::<Result<_, _>>().unwrap()
};
assert_eq!(
labels,
vec![
("local".to_string(), "proj-x".to_string()),
("starlight".to_string(), "proj-x".to_string()),
]
);
}
#[test]
fn cursor_key_includes_source_root_label() {
let conn = setup_conn();
let root = TempRoot::new("cursor-source-root");
let line =
r#"{"type":"user","message":{"content":[{"type":"text","text":"same synced file"}]}}"#;
root.write("proj-x/session-same.jsonl", &format!("{line}\n"));
let first = run(&conn, &[root.scan_root("local")]);
assert_eq!(first.ingested_messages, 1);
assert_eq!(cursor_count(&conn), 1);
let second = run(&conn, &[root.scan_root("starlight")]);
assert_eq!(
second.skipped, 0,
"same path under a different source_root must not use the old cursor"
);
assert_eq!(second.ingested_messages, 1);
assert_eq!(raw_message_count(&conn), 2);
assert_eq!(cursor_count(&conn), 2);
let labels: Vec<String> = {
let mut stmt = conn
.prepare("SELECT source_root FROM raw_messages ORDER BY source_root")
.unwrap();
let rows = stmt.query_map([], |row| row.get(0)).unwrap();
rows.collect::<Result<_, _>>().unwrap()
};
assert_eq!(labels, vec!["local".to_string(), "starlight".to_string()]);
}
#[test]
fn codex_session_meta_id_overrides_rollout_filename() {
let conn = setup_conn();
let root = TempRoot::new("codex-session-id");
root.write(
"2026/06/12/rollout-abc.jsonl",
include_str!("../../../tests/fixtures/codex-rollout-minimal.jsonl"),
);
let summary = run(&conn, &[root.scan_root("local")]);
assert_eq!(summary.failed_files, 0);
assert_eq!(summary.ingested_messages, 2);
let session_id: String = conn
.query_row("SELECT DISTINCT session_id FROM raw_messages", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(session_id, "019eba00-sanitized");
}
#[test]
fn transcript_timestamps_drive_raw_message_window_time() {
let conn = setup_conn();
let root = TempRoot::new("transcript-time");
root.write(
"proj-a/session-1.jsonl",
r#"{"timestamp":"2026-06-12T00:00:01.000Z","type":"user","message":{"content":[{"type":"text","text":"historical question"}]}}"#,
);
let summary = run(&conn, &[root.scan_root("local")]);
assert_eq!(summary.failed_files, 0);
assert_eq!(summary.ingested_messages, 1);
let created_at_epoch: i64 = conn
.query_row("SELECT created_at_epoch FROM raw_messages", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(created_at_epoch, 1_781_222_401);
}
#[test]
fn missing_explicit_root_is_reported_as_failure() {
let conn = setup_conn();
let missing = std::env::temp_dir().join(format!(
"remem-ingest-missing-{}-{}",
std::process::id(),
chrono::Utc::now().timestamp_nanos_opt().unwrap_or_default()
));
let summary = run_ingest_sessions(
&conn,
&[ScanRoot {
label: "remote".to_string(),
path: missing,
required: true,
}],
&IngestOptions::default(),
)
.unwrap();
assert_eq!(summary.failed_files, 1);
assert_eq!(summary.exit_code(), 1);
}
#[test]
fn incomplete_discovery_blocks_all_identity_and_raw_mutation() -> anyhow::Result<()> {
let conn = setup_conn();
let root = TempRoot::new("phase-a-fail-closed");
let cwd = root.path.to_string_lossy().to_string();
root.write(
"proj-a/good.jsonl",
&format!("{}\n", claude_line(&cwd, "user", "must not ingest")),
);
let missing = root.path.join("missing-required-root");
let roots = [
root.scan_root("local"),
ScanRoot {
label: "missing".to_string(),
path: missing,
required: true,
},
];
let summary = run_ingest_sessions(&conn, &roots, &IngestOptions::default())?;
assert_eq!(summary.failed_files, 1);
assert_eq!(summary.ingested_messages, 0);
assert_eq!(raw_message_count(&conn), 0);
assert_eq!(
conn.query_row("SELECT COUNT(*) FROM raw_session_identities", [], |row| {
row.get::<_, i64>(0)
})?,
0
);
assert_eq!(cursor_count(&conn), 0);
Ok(())
}
#[test]
fn incomplete_probe_blocks_all_phase_a_mutation() -> anyhow::Result<()> {
let conn = setup_conn();
let root = TempRoot::new("phase-a-probe-fail");
let cwd = root.path.to_string_lossy().to_string();
root.write(
"proj-a/good.jsonl",
&format!("{}\n", claude_line(&cwd, "user", "must not ingest")),
);
let broken = root.path.join("proj-a/broken.jsonl");
std::fs::write(&broken, [0xff, b'\n'])?;
let summary =
run_ingest_sessions(&conn, &[root.scan_root("local")], &IngestOptions::default())?;
assert_eq!(summary.failed_files, 1);
assert_eq!(raw_message_count(&conn), 0);
assert_eq!(
conn.query_row("SELECT COUNT(*) FROM raw_session_identities", [], |row| {
row.get::<_, i64>(0)
})?,
0
);
Ok(())
}
#[cfg(unix)]
#[test]
fn unreadable_discovery_entry_is_isolated_and_batch_continues() {
use std::os::unix::fs::PermissionsExt;
let conn = setup_conn();
let root = TempRoot::new("discovery-isolation");
let cwd = root.path.to_string_lossy().to_string();
root.write(
"proj-a/good.jsonl",
&format!("{}\n", claude_line(&cwd, "user", "healthy message")),
);
let blocked = root.path.join("proj-a").join("blocked");
std::fs::create_dir_all(&blocked).unwrap();
let original_permissions = std::fs::metadata(&blocked).unwrap().permissions();
std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o000)).unwrap();
let summary = run(&conn, &[root.scan_root("local")]);
std::fs::set_permissions(&blocked, original_permissions).unwrap();
assert_eq!(summary.ingested_messages, 0);
assert_eq!(summary.failed_files, 1);
assert_eq!(raw_message_count(&conn), 0);
}
#[test]
fn completion_failure_rolls_back_raw_rows_ledger_and_cursor() -> anyhow::Result<()> {
let conn = setup_conn();
let root = TempRoot::new("completion-rollback");
let cwd = root.path.to_string_lossy().to_string();
root.write(
"proj-a/session-1.jsonl",
&format!(
"{}\n",
claude_line(&cwd, "user", "rollback all file writes")
),
);
conn.execute_batch(
"CREATE TRIGGER fail_gh871_cursor
BEFORE INSERT ON ingest_cursors
BEGIN
SELECT RAISE(FAIL, 'forced cursor failure');
END;",
)?;
let summary =
run_ingest_sessions(&conn, &[root.scan_root("local")], &IngestOptions::default())?;
assert_eq!(summary.failed_files, 1);
assert_eq!(summary.ingested_messages, 0);
assert_eq!(raw_message_count(&conn), 0);
assert_eq!(cursor_count(&conn), 0);
assert_eq!(
conn.query_row(
"SELECT contract_version FROM raw_session_identities",
[],
|row| row.get::<_, i64>(0)
)?,
0
);
Ok(())
}
#[test]
fn ordinal_replay_conflict_is_sticky_and_preserves_existing_raw_row() -> anyhow::Result<()> {
let conn = setup_conn();
let root = TempRoot::new("ordinal-conflict");
let cwd = root.path.to_string_lossy().to_string();
let file = root.write(
"proj-a/session-1.jsonl",
&format!("{}\n", claude_line(&cwd, "user", "captured value")),
);
let scan_root = root.scan_root("local");
let plan =
crate::ingest::session_identity::probe(&scan_root.label, &scan_root.path, &file, None)?;
let identity_id = crate::ingest::session_identity::upsert_claim(&conn, &plan, 1)?;
crate::ingest::session_identity::resolve_fallback_group(
&conn,
&plan.source_root,
&plan.fallback_session_id,
)?;
conn.execute(
"INSERT INTO raw_messages (
session_id, project, role, content, content_hash, source,
created_at_epoch, source_root, event_time_source,
transcript_identity_id, transcript_record_ordinal
) VALUES (?1, ?2, 'assistant', 'different value', ?3, 'transcript',
1, ?4, 'transcript_event', ?5, 0)",
rusqlite::params![
plan.canonical_session_id,
plan.project,
crate::db::content_identity_hash(b"different value"),
plan.source_root,
identity_id
],
)?;
let summary = run_ingest_sessions(&conn, &[scan_root], &IngestOptions::default())?;
assert_eq!(summary.failed_files, 1);
assert_eq!(summary.ingested_messages, 0);
assert_eq!(cursor_count(&conn), 0);
assert_eq!(
conn.query_row(
"SELECT status || ':' || conflict_reason
FROM raw_session_identities WHERE id = ?1",
[identity_id],
|row| row.get::<_, String>(0)
)?,
"conflict:stable_occurrence_mismatch"
);
assert_eq!(
conn.query_row(
"SELECT role || ':' || content FROM raw_messages
WHERE transcript_identity_id = ?1",
[identity_id],
|row| row.get::<_, String>(0)
)?,
"assistant:different value"
);
Ok(())
}
#[test]
fn fallback_group_conflict_rolls_back_earlier_member_mutations() -> anyhow::Result<()> {
let conn = setup_conn();
let root = TempRoot::new("group-atomicity");
let cwd = root.path.to_string_lossy().to_string();
let first = root.write(
"a/shared.jsonl",
&format!(
"{}\n",
serde_json::json!({
"type": "user",
"sessionId": "canonical-871",
"cwd": cwd,
"timestamp": 100,
"message": {"content": "first group member"}
})
),
);
let second = root.write(
"b/shared.jsonl",
&format!(
"{}\n",
serde_json::json!({
"type": "user",
"sessionId": "canonical-871",
"cwd": cwd,
"timestamp": 101,
"message": {"content": "second group member"}
})
),
);
assert!(
first < second,
"fixture ordering must exercise earlier mutation"
);
let scan_root = root.scan_root("local");
let second_plan =
crate::ingest::session_identity::probe(&scan_root.label, &scan_root.path, &second, None)?;
let second_identity = crate::ingest::session_identity::upsert_claim(&conn, &second_plan, 1)?;
crate::ingest::session_identity::resolve_fallback_group(
&conn,
&second_plan.source_root,
&second_plan.fallback_session_id,
)?;
conn.execute(
"INSERT INTO raw_messages (
session_id, project, role, content, content_hash, source,
created_at_epoch, source_root, event_time_source,
transcript_identity_id, transcript_record_ordinal
) VALUES (?1, ?2, 'assistant', 'preexisting mismatch', ?3, 'transcript',
101, ?4, 'transcript_event', ?5, 0)",
rusqlite::params![
second_plan.canonical_session_id,
second_plan.project,
crate::db::content_identity_hash(b"preexisting mismatch"),
second_plan.source_root,
second_identity
],
)?;
let summary = run_ingest_sessions(&conn, &[scan_root], &IngestOptions::default())?;
assert_eq!(summary.failed_files, 1);
assert_eq!(summary.ingested_messages, 0);
assert_eq!(cursor_count(&conn), 0);
assert_eq!(raw_message_count(&conn), 1);
assert_eq!(
conn.query_row("SELECT content FROM raw_messages", [], |row| row
.get::<_, String>(0))?,
"preexisting mismatch"
);
assert_eq!(
conn.query_row(
"SELECT COUNT(*) FROM raw_session_identities
WHERE fallback_session_id = 'shared' AND status = 'conflict'",
[],
|row| row.get::<_, i64>(0)
)?,
2
);
Ok(())
}
#[test]
fn filename_fallback_promotion_updates_one_existing_occurrence() -> anyhow::Result<()> {
let conn = setup_conn();
let root = TempRoot::new("occurrence-promotion");
let cwd = root.path.to_string_lossy().to_string();
let file = root.write(
"proj-a/fallback-name.jsonl",
&format!(
"{}\n",
serde_json::json!({
"type": "user",
"cwd": cwd,
"timestamp": 100,
"message": {"content": "same occurrence"}
})
),
);
let scan_root = root.scan_root("local");
let first = run_ingest_sessions(
&conn,
std::slice::from_ref(&scan_root),
&IngestOptions::default(),
)?;
assert_eq!(first.ingested_messages, 1);
let identity_id: i64 = conn.query_row(
"SELECT transcript_identity_id FROM raw_messages",
[],
|row| row.get(0),
)?;
std::fs::write(
&file,
format!(
"{}\n",
serde_json::json!({
"type": "user",
"sessionId": "canonical-871",
"cwd": cwd,
"timestamp": 100,
"message": {"content": "same occurrence"}
})
),
)?;
let second = run_ingest_sessions(&conn, &[scan_root], &IngestOptions::default())?;
assert_eq!(second.failed_files, 0);
assert_eq!(
conn.query_row("SELECT COUNT(*) FROM raw_messages", [], |row| {
row.get::<_, i64>(0)
})?,
1
);
assert_eq!(
conn.query_row(
"SELECT transcript_identity_id || ':' || session_id
FROM raw_messages",
[],
|row| row.get::<_, String>(0)
)?,
format!("{identity_id}:canonical-871")
);
Ok(())
}
#[test]
fn raw_sessions_have_created_at_leading_index() {
let conn = setup_conn();
let exists: i64 = conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master
WHERE type = 'index' AND name = 'idx_raw_messages_created_source_project_session'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(exists, 1);
}
#[test]
fn scan_root_parse_accepts_label_path_and_rejects_bad_specs() {
let parsed = ScanRoot::parse("starlight=/tmp/remote-sessions").unwrap();
assert_eq!(parsed.label, "starlight");
assert_eq!(parsed.path, PathBuf::from("/tmp/remote-sessions"));
assert!(ScanRoot::parse("no-separator").is_err());
assert!(ScanRoot::parse("=path-only").is_err());
assert!(ScanRoot::parse("label=").is_err());
}