From c580ed61301d41654d5a6582ba5aa6308739d228 Mon Sep 17 00:00:00 2001 From: Dmitry Starkov <21260939+starkdmi@users.noreply.github.com> Date: Fri, 21 Aug 2026 22:09:37 +0300 Subject: [PATCH] Repair Claude archive collection and cut import work to what changed `conversation collect --provider claude_code` failed outright with a foreign key violation. Claude records the conversation identity per line, and a resumed session records its parent's identifier before its own, so trace edits reconstructed mid-file referenced a conversation that was never written. Edits are now re-bound to the conversation their file is finally written as. The import revision moves to archive.v8 so files imported under the old reconstruction are read again, which means the first run after this change re-imports every source. The store did work proportional to the statements it issued rather than to the content that changed: - Re-importing replaced byte-identical rows, deleting and reinserting each one along with its full-text index entry. A row is kept when every persisted field matches, and base64 is decoded only once a write is known to be needed. - Items, content parts, and their reconciliation were written one statement at a time. Once a transaction has touched FTS5, SQLite flushes that index at every statement boundary, so these are now batched, and reconciling stale source records is two statements per conversation against a staged table, joined in an order that keeps it from walking every item recorded for the source. - Archive files are reconstructed on several threads and written in file order, bounded by how much reconstructed content is held at once. A file that fails still leaves every file before it stored. Commit durability is relaxed for the imports alone, which rebuild from the provider's own files. The code-change refresh that follows carries metrics forward for commits Git no longer rescans, so it keeps the store's normal durability. Measured against the previous build on the same archives, back to back: OpenCode 384s to 24s, a 120-file Codex fixture 103s to 6s, and re-imports write no rows where they previously rewrote every one. Imported content is byte-identical to the previous build's. --- crates/statsai-adapters/src/archive.rs | 91 +- crates/statsai-adapters/src/lib.rs | 2 +- crates/statsai-core/src/code_changes/types.rs | 21 + crates/statsai-store/src/archive.rs | 914 +++++++++++++----- crates/statsai-store/src/lib.rs | 113 +++ crates/statsai/src/main.rs | 552 +++++++++-- 6 files changed, 1359 insertions(+), 334 deletions(-) diff --git a/crates/statsai-adapters/src/archive.rs b/crates/statsai-adapters/src/archive.rs index 0258a48..1731bc9 100644 --- a/crates/statsai-adapters/src/archive.rs +++ b/crates/statsai-adapters/src/archive.rs @@ -593,6 +593,11 @@ fn collect_claude_file( let mut line_number = 0usize; let mut pending_mutations = HashMap::new(); let source_record_path = canonical_display(path); + // Claude records the conversation identity per line, and a resumed session + // carries its parent's identifier before its own. Edits reconstructed + // before the last identity is known are re-bound to it below, so they + // always reference the conversation this file is finally written as. + let trace_edits_start = trace_edits.len(); loop { let status = read_bounded_jsonl_line(&mut reader, &mut line_bytes, MAX_JSONL_RECORD_BYTES)?; if status == BoundedLineRead::Eof { @@ -772,7 +777,11 @@ fn collect_claude_file( } } mark_unresolved_mutations(&pending_mutations, diagnostics, trace_coverage); - Ok(builder.finish()) + let conversation = builder.finish(); + for edit in &mut trace_edits[trace_edits_start..] { + edit.rebind_conversation(&conversation.conversation_id); + } + Ok(conversation) } fn claude_archive_native_id( @@ -1289,7 +1298,11 @@ struct ItemInput<'a> { } fn item_from_value(input: ItemInput<'_>) -> (ArchiveItem, u64) { - let fingerprint = hash_text(&input.content.to_string()); + // Rendered once and reused: the identifier hashes the same JSON that a + // tool call or result goes on to store, and rendering a large tool payload + // twice is pure duplicate work. + let rendered = input.content.to_string(); + let fingerprint = hash_text(&rendered); let item_id = archive_item_id( input.provider, input.conversation_native_id, @@ -1307,7 +1320,7 @@ fn item_from_value(input: ItemInput<'_>) -> (ArchiveItem, u64) { if !input.content.is_null() { let text = match input.content { Value::String(value) => value.clone(), - value => value.to_string(), + _ => rendered.clone(), }; if !text.trim().is_empty() && text != "null" { push_text_part(&item_id, ArchiveContentKind::Json, text, &mut parts); @@ -1339,7 +1352,7 @@ fn item_from_value(input: ItemInput<'_>) -> (ArchiveItem, u64) { } else if parts.is_empty() && !input.content.is_null() { let text = match input.content { Value::String(value) => value.clone(), - value => value.to_string(), + _ => rendered, }; if !text.trim().is_empty() && text != "null" { push_text_part(&item_id, ArchiveContentKind::Json, text, &mut parts); @@ -1582,7 +1595,9 @@ fn extract_content_parts( return; } if matches!(content_type, "tool_use" | "tool_call" | "tool_result") { - let compact = Value::Object(object.clone()).to_string(); + // Serialized borrowed: cloning the block first copied every + // nested value only to render and drop it. + let compact = serde_json::to_string(object).unwrap_or_default(); push_text_part(item_id, ArchiveContentKind::Json, compact, parts); return; } @@ -2156,6 +2171,72 @@ mod tests { ); } + /// A resumed session records its parent's identifier before its own, so the + /// conversation this file becomes is known only at the end. Every edit must + /// still reference it, or the store rejects the import outright. + #[test] + fn claude_resumed_session_binds_every_edit_to_the_written_conversation() { + let dir = tempdir().unwrap(); + let project = dir.path().join("projects").join("workspace"); + std::fs::create_dir_all(&project).unwrap(); + let path = project.join("resumed.jsonl"); + let mut file = File::create(&path).unwrap(); + // The parent session's edit is reconstructed before the record that + // reveals this file is really the resumed session. + for record in [ + serde_json::json!({ + "sessionId": "parent-session", + "type": "assistant", + "message": {"role": "assistant", "content": [{ + "type": "tool_use", + "id": "call-1", + "name": "Write", + "input": {"file_path": "src/lib.rs", "content": "one\ntwo\n"} + }]} + }), + serde_json::json!({ + "sessionId": "parent-session", + "type": "user", + "message": {"role": "user", "content": [{ + "type": "tool_result", + "tool_use_id": "call-1", + "content": "File created successfully at: src/lib.rs" + }]} + }), + serde_json::json!({ + "sessionId": "resumed-session", + "type": "user", + "message": {"role": "user", "content": [{"type": "text", "text": "carry on"}]} + }), + ] { + writeln!(file, "{record}").unwrap(); + } + drop(file); + + let scan = collect_claude(&source(CLAUDE_CODE_PROVIDER, dir.path()), None).unwrap(); + + assert_eq!(scan.conversations.len(), 1); + let conversation = &scan.conversations[0]; + assert_eq!(conversation.native_conversation_id, "resumed-session"); + assert!(!scan.trace_edits.is_empty(), "the write must be measured"); + // Every edit references the conversation that was actually written, and + // no two edits collapsed onto one identifier while being re-bound. + assert!(scan + .trace_edits + .iter() + .all(|edit| edit.conversation_id == conversation.conversation_id)); + let distinct_ids = scan + .trace_edits + .iter() + .map(|edit| edit.trace_edit_id.as_str()) + .collect::>(); + assert_eq!(distinct_ids.len(), scan.trace_edits.len()); + // Re-reading the same file reproduces the identifiers, so a re-import + // replaces the edits instead of duplicating them. + let repeated = collect_claude(&source(CLAUDE_CODE_PROVIDER, dir.path()), None).unwrap(); + assert_eq!(repeated.trace_edits, scan.trace_edits); + } + #[test] fn codex_collects_visible_reasoning_and_exact_embedded_image() { let dir = tempdir().unwrap(); diff --git a/crates/statsai-adapters/src/lib.rs b/crates/statsai-adapters/src/lib.rs index bd3d0c9..50dadf7 100644 --- a/crates/statsai-adapters/src/lib.rs +++ b/crates/statsai-adapters/src/lib.rs @@ -200,7 +200,7 @@ pub struct AdapterScan { pub verified_source_state: Option, } -pub trait ProviderAdapter { +pub trait ProviderAdapter: Send + Sync { fn id(&self) -> &'static str; fn version(&self) -> &'static str; fn provider(&self) -> &'static str; diff --git a/crates/statsai-core/src/code_changes/types.rs b/crates/statsai-core/src/code_changes/types.rs index 8acac3d..0901bf7 100644 --- a/crates/statsai-core/src/code_changes/types.rs +++ b/crates/statsai-core/src/code_changes/types.rs @@ -207,6 +207,27 @@ pub struct TraceEdit { pub deleted_line_fingerprints: Vec, } +impl TraceEdit { + /// Re-files this edit under the conversation that received its items. + /// + /// A transcript can reveal the conversation it belongs to only once the + /// whole file is read — a resumed session records its parent's identifier + /// before its own — so an edit is reconstructed against the identity that + /// was current at the time and re-bound here. Without this the edit + /// references a conversation that is never written, and the store rejects + /// it. The new identifier is derived from the old one, so it stays + /// deterministic: re-importing the same file replaces the edit instead of + /// duplicating it. + pub fn rebind_conversation(&mut self, conversation_id: &str) { + if self.conversation_id == conversation_id { + return; + } + let fingerprint = crate::hash_text(&format!("{}:{}", self.trace_edit_id, conversation_id)); + self.trace_edit_id = format!("edit_{}", &fingerprint[..24]); + self.conversation_id = conversation_id.to_string(); + } +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct ParsedMutation { pub edits: Vec, diff --git a/crates/statsai-store/src/archive.rs b/crates/statsai-store/src/archive.rs index 78c9ab4..0dfdb5d 100644 --- a/crates/statsai-store/src/archive.rs +++ b/crates/statsai-store/src/archive.rs @@ -19,8 +19,10 @@ use std::path::Path; // reclassified read-only shell calls and stopped treating echoed file content // as a failed mutation; v7 reads whole-file creation from the tool result, so // files a `Write` created become counted additions with matchable fingerprints -// instead of unclassified lines. -const ARCHIVE_IMPORT_REVISION: &str = "archive.v7"; +// instead of unclassified lines; v8 binds a trace edit to the conversation its +// file is written as, so edits a resumed session recorded against its parent +// stop being attributed to that parent. +const ARCHIVE_IMPORT_REVISION: &str = "archive.v8"; const UNSCOPED_MISSING_CONTENT_SCOPE: &str = "unscoped"; #[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize)] @@ -78,8 +80,111 @@ struct ContentRetentionQuality { #[derive(Debug, Clone, PartialEq, Eq)] struct RetainedContentPart { + /// Item the stored row belongs to. Retention is read for a whole + /// conversation at once, so an authoritative item must be able to tell its + /// own obsolete content apart from its neighbours'. + item_id: String, content_hash: String, quality: ContentRetentionQuality, + /// Everything persisted about the content that its hash does not cover. + /// + /// The hash is taken over the content alone, and identically over text and + /// over bytes, so it says nothing about how the content is described or + /// even which column holds it. Everything a read would reconstruct is + /// compared here rather than inferred from what a provider happens to + /// produce today. + metadata: RetainedContentMetadata, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +struct RetainedContentMetadata { + ordinal: u64, + kind: String, + mime_type: Option, + name: Option, + external_uri: Option, + original_bytes: u64, + stored_as_text: bool, + stored_as_binary: bool, +} + +impl RetainedContentPart { + /// Whether the stored row would read back as exactly this part. + fn already_stores( + &self, + part: &ArchiveContentPart, + quality: ContentRetentionQuality, + item_id: &str, + ) -> bool { + self.item_id == item_id + && self.content_hash == part.content_hash + && self.quality == quality + && self.metadata + == RetainedContentMetadata { + ordinal: part.ordinal, + kind: part.kind.as_str().to_string(), + mime_type: part.mime_type.clone(), + name: part.name.clone(), + external_uri: part.external_uri.clone(), + original_bytes: part.original_bytes, + stored_as_text: part.text.is_some(), + stored_as_binary: part.data_base64.is_some(), + } + } +} + +/// Rows per batched write. +/// +/// Large enough to amortize the full-text index flush that SQLite performs at +/// every statement boundary, small enough to stay far inside the bound on host +/// parameters per statement. +const ARCHIVE_WRITE_BATCH_ROWS: usize = 64; +/// Content bytes after which a batch is issued regardless of its row count. +const ARCHIVE_WRITE_BATCH_BYTES: usize = 8 * 1024 * 1024; +const ARCHIVE_DELETE_BATCH_ROWS: usize = 256; + +/// Item columns that are derived rather than stored directly on [`ArchiveItem`]. +struct EncodedItem { + kind: &'static str, + role: Option<&'static str>, + created_at: Option, + model_json: Option, + usage_json: Option, +} + +struct PendingContentPart<'a> { + part: &'a ArchiveContentPart, + item_id: &'a str, + kind: &'static str, + binary_content: Option>, +} + +impl PendingContentPart<'_> { + fn stored_bytes(&self) -> usize { + self.part.text.as_ref().map_or(0, String::len) + + self.binary_content.as_ref().map_or(0, |bytes| bytes.len()) + } +} + +/// Appends `rows` parenthesised groups of `columns` placeholders. +fn append_row_placeholders(sql: &mut String, rows: usize, columns: usize) { + for row in 0..rows { + if row > 0 { + sql.push(','); + } + if columns == 1 { + sql.push('?'); + continue; + } + sql.push('('); + for column in 0..columns { + if column > 0 { + sql.push(','); + } + sql.push('?'); + } + sql.push(')'); + } } impl Store { @@ -268,220 +373,18 @@ impl Store { )?; result.conversations += 1; - for source_record_id in &conversation.discarded_source_record_ids { - self.conn.execute( - r#" - DELETE FROM archive_content_parts - WHERE item_id IN ( - SELECT archive_items.item_id - FROM archive_items - JOIN archive_conversations USING (conversation_id) - WHERE archive_items.source_record_id = ?1 - AND archive_conversations.provider = ?2 - AND archive_conversations.source_id = ?3 - ) - "#, - params![ - source_record_id, - &conversation.provider, - &conversation.source_id.0, - ], - )?; - self.conn.execute( - r#" - DELETE FROM archive_items - WHERE item_id IN ( - SELECT archive_items.item_id - FROM archive_items - JOIN archive_conversations USING (conversation_id) - WHERE archive_items.source_record_id = ?1 - AND archive_conversations.provider = ?2 - AND archive_conversations.source_id = ?3 - ) - "#, - params![ - source_record_id, - &conversation.provider, - &conversation.source_id.0, - ], - )?; - } - - for item in &conversation.items { - if let Some(source_record_id) = item.source_record_id.as_deref() { - self.conn.execute( - r#" - DELETE FROM archive_content_parts - WHERE item_id IN ( - SELECT archive_items.item_id - FROM archive_items - JOIN archive_conversations USING (conversation_id) - WHERE archive_items.source_record_id = ?1 - AND archive_items.item_id <> ?2 - AND archive_conversations.provider = ?3 - AND archive_conversations.source_id = ?4 - ) - "#, - params![ - source_record_id, - &item.item_id, - &conversation.provider, - &conversation.source_id.0, - ], - )?; - self.conn.execute( - r#" - DELETE FROM archive_items - WHERE item_id IN ( - SELECT archive_items.item_id - FROM archive_items - JOIN archive_conversations USING (conversation_id) - WHERE archive_items.source_record_id = ?1 - AND archive_items.item_id <> ?2 - AND archive_conversations.provider = ?3 - AND archive_conversations.source_id = ?4 - ) - "#, - params![ - source_record_id, - &item.item_id, - &conversation.provider, - &conversation.source_id.0, - ], - )?; - } - let model_json = item.model.as_ref().map(serde_json::to_string).transpose()?; - let usage_json = item.usage.as_ref().map(serde_json::to_string).transpose()?; - self.conn.execute( - r#" - INSERT INTO archive_items - (item_id, conversation_id, native_item_id, source_record_id, ordinal, - kind, role, created_at, model_json, tool_name, tool_call_id, status, - usage_json) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13) - ON CONFLICT(item_id) DO UPDATE SET - conversation_id = excluded.conversation_id, - native_item_id = COALESCE(excluded.native_item_id, archive_items.native_item_id), - source_record_id = COALESCE(excluded.source_record_id, archive_items.source_record_id), - ordinal = excluded.ordinal, - kind = excluded.kind, - role = COALESCE(excluded.role, archive_items.role), - created_at = COALESCE(excluded.created_at, archive_items.created_at), - model_json = COALESCE(excluded.model_json, archive_items.model_json), - tool_name = COALESCE(excluded.tool_name, archive_items.tool_name), - tool_call_id = COALESCE(excluded.tool_call_id, archive_items.tool_call_id), - status = COALESCE(excluded.status, archive_items.status), - usage_json = COALESCE(excluded.usage_json, archive_items.usage_json) - "#, - params![ - &item.item_id, - &conversation.conversation_id, - &item.native_item_id, - &item.source_record_id, - item.ordinal, - item.kind.as_str(), - item.role.map(ArchiveRole::as_str), - item.created_at.map(|value| value.to_rfc3339()), - model_json, - &item.tool_name, - &item.tool_call_id, - &item.status, - usage_json, - ], - )?; - result.items += 1; - - let mut retained_parts = self.archive_content_retention(&item.item_id)?; - let incoming_part_ids = item - .parts - .iter() - .map(|part| part.content_id.as_str()) - .collect::>(); - for part in &item.parts { - let binary_content = part - .data_base64 - .as_deref() - .map(|encoded| { - BASE64.decode(encoded).with_context(|| { - format!("decode archive content {}", part.content_id) - }) - }) - .transpose()?; - let incoming_quality = ContentRetentionQuality { - materialized: part.text.is_some() || binary_content.is_some(), - external: part.external_uri.is_some(), - stored_bytes: part.text.as_ref().map_or(0, |text| text.len() as u64) - + binary_content - .as_ref() - .map_or(0, |bytes| bytes.len() as u64), - truncated: part.truncated, - }; - if retained_parts - .get(&part.content_id) - .is_some_and(|existing| { - !incoming_content_should_replace( - existing, - &part.content_hash, - incoming_quality, - ) - }) - { - continue; - } - if retained_parts.contains_key(&part.content_id) { - self.conn.execute( - "DELETE FROM archive_content_parts WHERE content_id = ?1", - params![&part.content_id], - )?; - } - result.binary_bytes += binary_content - .as_ref() - .map_or(0, |bytes| bytes.len() as u64); - self.conn.execute( - r#" - INSERT INTO archive_content_parts - (content_id, item_id, ordinal, kind, mime_type, name, text_content, - binary_content, external_uri, content_hash, original_bytes, truncated) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12) - "#, - params![ - &part.content_id, - &item.item_id, - part.ordinal, - part.kind.as_str(), - &part.mime_type, - &part.name, - &part.text, - binary_content, - &part.external_uri, - &part.content_hash, - part.original_bytes, - part.truncated, - ], - )?; - result.content_parts += 1; - retained_parts.insert( - part.content_id.clone(), - RetainedContentPart { - content_hash: part.content_hash.clone(), - quality: incoming_quality, - }, - ); - } - if item.parts_authoritative { - let obsolete_part_ids = retained_parts - .keys() - .filter(|content_id| !incoming_part_ids.contains(content_id.as_str())) - .cloned() - .collect::>(); - for content_id in obsolete_part_ids { - self.conn.execute( - "DELETE FROM archive_content_parts WHERE content_id = ?1", - params![content_id], - )?; - } - } - } + // Reconciliation, item writes, and content writes each run as a + // handful of statements for the whole conversation rather than a + // few per item. Once a transaction has touched the full-text index, + // SQLite flushes that index at every statement boundary, so the + // cost of an import followed the number of statements it issued far + // more closely than the amount of content that had actually + // changed. + self.stage_conversation_records(conversation)?; + self.delete_replaced_source_records(conversation)?; + let mut retained_parts = self.staged_content_retention()?; + self.insert_archive_items(conversation, &mut result)?; + self.write_archive_content_parts(conversation, &mut retained_parts, &mut result)?; for superseded_id in &conversation.superseded_conversation_ids { if superseded_id == &conversation.conversation_id { @@ -602,6 +505,396 @@ impl Store { Ok(result) } + /// Stages the source records and item identifiers this conversation brings. + /// + /// Reconciliation needs to ask "which stored items does this import + /// replace" for every incoming record at once. Staging the batch turns two + /// statements per item into two statements per conversation. + /// + /// A discarded record stages an empty item identifier, which no real item + /// can hold, so every stored copy of that record is replaced by nothing. + /// An item without a source record stages a NULL, which joins to nothing + /// and so replaces nothing, while still carrying its identifier for the + /// content lookup. + fn stage_conversation_records(&self, conversation: &ArchiveConversation) -> Result<()> { + self.conn.execute("DELETE FROM incoming_records", [])?; + const DISCARDED: &str = ""; + let staged = conversation + .discarded_source_record_ids + .iter() + .map(|source_record_id| (Some(source_record_id), DISCARDED)) + .chain( + conversation + .items + .iter() + .map(|item| (item.source_record_id.as_ref(), item.item_id.as_str())), + ) + .collect::>(); + for chunk in staged.chunks(ARCHIVE_WRITE_BATCH_ROWS) { + let mut sql = + String::from("INSERT INTO incoming_records (source_record_id, item_id) VALUES "); + append_row_placeholders(&mut sql, chunk.len(), 2); + let mut bindings: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(chunk.len() * 2); + for (source_record_id, item_id) in chunk { + bindings.push(source_record_id); + bindings.push(item_id); + } + self.conn + .prepare_cached(&sql)? + .execute(bindings.as_slice())?; + } + Ok(()) + } + + /// Removes the stored items that the staged records supersede. + /// + /// The joins are written as `CROSS JOIN` to pin the order: the staged + /// records are the small side and must drive, so each one probes the stored + /// items by source record. Left to its own estimates the planner walks + /// every item recorded for the source instead, which turns importing a + /// source into quadratic work. + fn delete_replaced_source_records(&self, conversation: &ArchiveConversation) -> Result<()> { + for table in ["archive_content_parts", "archive_items"] { + self.conn + .prepare_cached(&format!( + r#" + DELETE FROM {table} + WHERE item_id IN ( + SELECT stored.item_id + FROM incoming_records incoming + CROSS JOIN archive_items stored + ON stored.source_record_id = incoming.source_record_id + CROSS JOIN archive_conversations conversation + ON conversation.conversation_id = stored.conversation_id + WHERE stored.item_id <> incoming.item_id + AND conversation.provider = ?1 + AND conversation.source_id = ?2 + ) + "# + ))? + .execute(params![&conversation.provider, &conversation.source_id.0])?; + } + Ok(()) + } + + /// Content already stored for the staged items. + fn staged_content_retention(&self) -> Result> { + let mut statement = self.conn.prepare_cached( + r#" + SELECT part.content_id, + part.text_content IS NOT NULL OR part.binary_content IS NOT NULL, + part.external_uri IS NOT NULL, + COALESCE(length(CAST(part.text_content AS BLOB)), 0) + + COALESCE(length(part.binary_content), 0), + part.truncated, + part.content_hash, + part.item_id, + part.kind, + part.mime_type, + part.name, + part.original_bytes, + part.ordinal, + part.external_uri, + part.text_content IS NOT NULL, + part.binary_content IS NOT NULL + FROM incoming_records incoming + CROSS JOIN archive_content_parts part ON part.item_id = incoming.item_id + "#, + )?; + let rows = statement.query_map([], |row| { + Ok(( + row.get::<_, String>(0)?, + RetainedContentPart { + quality: ContentRetentionQuality { + materialized: row.get(1)?, + external: row.get(2)?, + stored_bytes: row.get(3)?, + truncated: row.get(4)?, + }, + content_hash: row.get(5)?, + item_id: row.get(6)?, + metadata: RetainedContentMetadata { + kind: row.get(7)?, + mime_type: row.get(8)?, + name: row.get(9)?, + original_bytes: row.get(10)?, + ordinal: row.get(11)?, + external_uri: row.get(12)?, + stored_as_text: row.get(13)?, + stored_as_binary: row.get(14)?, + }, + }, + )) + })?; + rows.collect::>() + .map_err(Into::into) + } + + fn insert_archive_items( + &self, + conversation: &ArchiveConversation, + result: &mut ArchiveWriteResult, + ) -> Result<()> { + let encoded = conversation + .items + .iter() + .map(|item| { + Ok(EncodedItem { + kind: item.kind.as_str(), + role: item.role.map(ArchiveRole::as_str), + created_at: item.created_at.map(|value| value.to_rfc3339()), + model_json: item.model.as_ref().map(serde_json::to_string).transpose()?, + usage_json: item.usage.as_ref().map(serde_json::to_string).transpose()?, + }) + }) + .collect::>>()?; + + for (items, encoded) in conversation + .items + .chunks(ARCHIVE_WRITE_BATCH_ROWS) + .zip(encoded.chunks(ARCHIVE_WRITE_BATCH_ROWS)) + { + let mut sql = String::from( + "INSERT INTO archive_items + (item_id, conversation_id, native_item_id, source_record_id, ordinal, + kind, role, created_at, model_json, tool_name, tool_call_id, status, + usage_json) VALUES ", + ); + append_row_placeholders(&mut sql, items.len(), 13); + sql.push_str( + " ON CONFLICT(item_id) DO UPDATE SET + conversation_id = excluded.conversation_id, + native_item_id = COALESCE(excluded.native_item_id, archive_items.native_item_id), + source_record_id = COALESCE(excluded.source_record_id, archive_items.source_record_id), + ordinal = excluded.ordinal, + kind = excluded.kind, + role = COALESCE(excluded.role, archive_items.role), + created_at = COALESCE(excluded.created_at, archive_items.created_at), + model_json = COALESCE(excluded.model_json, archive_items.model_json), + tool_name = COALESCE(excluded.tool_name, archive_items.tool_name), + tool_call_id = COALESCE(excluded.tool_call_id, archive_items.tool_call_id), + status = COALESCE(excluded.status, archive_items.status), + usage_json = COALESCE(excluded.usage_json, archive_items.usage_json)", + ); + let mut bindings: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(items.len() * 13); + for (item, encoded) in items.iter().zip(encoded) { + bindings.push(&item.item_id); + bindings.push(&conversation.conversation_id); + bindings.push(&item.native_item_id); + bindings.push(&item.source_record_id); + bindings.push(&item.ordinal); + bindings.push(&encoded.kind); + bindings.push(&encoded.role); + bindings.push(&encoded.created_at); + bindings.push(&encoded.model_json); + bindings.push(&item.tool_name); + bindings.push(&item.tool_call_id); + bindings.push(&item.status); + bindings.push(&encoded.usage_json); + } + self.conn + .prepare_cached(&sql)? + .execute(bindings.as_slice())?; + } + result.items += conversation.items.len() as u64; + Ok(()) + } + + /// Writes the content this import adds or improves, and drops what an + /// authoritative item no longer lists. + fn write_archive_content_parts( + &self, + conversation: &ArchiveConversation, + retained_parts: &mut HashMap, + result: &mut ArchiveWriteResult, + ) -> Result<()> { + let mut replaced = Vec::new(); + let mut obsolete = Vec::new(); + let mut pending = Vec::new(); + let mut pending_bytes = 0usize; + for item in &conversation.items { + let incoming_part_ids = item + .parts + .iter() + .map(|part| part.content_id.as_str()) + .collect::>(); + for part in &item.parts { + // Decoding is deferred until the row is known to need writing: + // on a re-import most parts are already stored byte for byte, + // and decoding them only to discard the result is the most + // expensive thing this loop can do to an archive of images. + let incoming_quality = ContentRetentionQuality { + materialized: part.text.is_some() || part.data_base64.is_some(), + external: part.external_uri.is_some(), + stored_bytes: part.text.as_ref().map_or(0, |text| text.len() as u64) + + part.data_base64.as_deref().map_or(0, base64_decoded_len), + truncated: part.truncated, + }; + match retained_parts.get(&part.content_id) { + // The stored row already is the incoming row, down to the + // metadata the content hash does not cover. Rewriting it + // would delete and reinsert identical bytes and drag the + // full-text index through both. + Some(existing) + if existing.already_stores(part, incoming_quality, &item.item_id) => + { + continue + } + Some(existing) + if !incoming_content_should_replace( + existing, + &part.content_hash, + incoming_quality, + ) => + { + continue + } + _ => {} + } + let binary_content = part + .data_base64 + .as_deref() + .map(|encoded| { + BASE64 + .decode(encoded) + .with_context(|| format!("decode archive content {}", part.content_id)) + }) + .transpose()?; + if retained_parts.contains_key(&part.content_id) { + replaced.push(part.content_id.clone()); + } + result.binary_bytes += binary_content + .as_ref() + .map_or(0, |bytes| bytes.len() as u64); + result.content_parts += 1; + retained_parts.insert( + part.content_id.clone(), + RetainedContentPart { + item_id: item.item_id.clone(), + content_hash: part.content_hash.clone(), + quality: incoming_quality, + metadata: RetainedContentMetadata { + ordinal: part.ordinal, + kind: part.kind.as_str().to_string(), + mime_type: part.mime_type.clone(), + name: part.name.clone(), + external_uri: part.external_uri.clone(), + original_bytes: part.original_bytes, + stored_as_text: part.text.is_some(), + stored_as_binary: part.data_base64.is_some(), + }, + }, + ); + pending_bytes += binary_content.as_ref().map_or(0, Vec::len); + pending.push(PendingContentPart { + part, + item_id: item.item_id.as_str(), + kind: part.kind.as_str(), + binary_content, + }); + // Decoded content is written out as it accumulates rather than + // once the conversation has been walked: a conversation can + // reference artifacts far larger than the transcript naming + // them, and holding every decoded copy at once is what makes + // peak memory a property of the conversation instead of the + // batch. + if pending.len() >= ARCHIVE_WRITE_BATCH_ROWS + || pending_bytes >= ARCHIVE_WRITE_BATCH_BYTES + { + self.delete_content_parts(&replaced)?; + self.insert_content_parts(&pending)?; + replaced.clear(); + pending.clear(); + pending_bytes = 0; + } + } + if item.parts_authoritative { + obsolete.extend( + retained_parts + .iter() + .filter(|(content_id, retained)| { + retained.item_id == item.item_id + && !incoming_part_ids.contains(content_id.as_str()) + }) + .map(|(content_id, _)| content_id.clone()), + ); + } + } + // Each batch deletes the rows it replaces before writing them, so a + // replacement never races the insert that supersedes it. Content this + // import did not list is only dropped once every row that survives is + // in, which keeps an interrupted transaction from being the one that + // removed content without writing its replacement. + self.delete_content_parts(&replaced)?; + self.insert_content_parts(&pending)?; + self.delete_content_parts(&obsolete)?; + for content_id in &obsolete { + retained_parts.remove(content_id); + } + Ok(()) + } + + fn delete_content_parts(&self, content_ids: &[String]) -> Result<()> { + for chunk in content_ids.chunks(ARCHIVE_DELETE_BATCH_ROWS) { + let mut sql = String::from("DELETE FROM archive_content_parts WHERE content_id IN ("); + append_row_placeholders(&mut sql, chunk.len(), 1); + sql.push(')'); + let bindings = chunk + .iter() + .map(|content_id| content_id as &dyn rusqlite::ToSql) + .collect::>(); + self.conn + .prepare_cached(&sql)? + .execute(bindings.as_slice())?; + } + Ok(()) + } + + fn insert_content_parts(&self, pending: &[PendingContentPart<'_>]) -> Result<()> { + let mut start = 0; + while start < pending.len() { + // Batches are bounded by content size as well as row count so that + // a run of large artifacts cannot build an unbounded parameter list. + let mut end = start; + let mut batch_bytes = 0usize; + while end < pending.len() && end - start < ARCHIVE_WRITE_BATCH_ROWS { + batch_bytes += pending[end].stored_bytes(); + end += 1; + if batch_bytes >= ARCHIVE_WRITE_BATCH_BYTES { + break; + } + } + let chunk = &pending[start..end]; + let mut sql = String::from( + "INSERT INTO archive_content_parts + (content_id, item_id, ordinal, kind, mime_type, name, text_content, + binary_content, external_uri, content_hash, original_bytes, truncated) + VALUES ", + ); + append_row_placeholders(&mut sql, chunk.len(), 12); + let mut bindings: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(chunk.len() * 12); + for entry in chunk { + bindings.push(&entry.part.content_id); + bindings.push(&entry.item_id); + bindings.push(&entry.part.ordinal); + bindings.push(&entry.kind); + bindings.push(&entry.part.mime_type); + bindings.push(&entry.part.name); + bindings.push(&entry.part.text); + bindings.push(&entry.binary_content); + bindings.push(&entry.part.external_uri); + bindings.push(&entry.part.content_hash); + bindings.push(&entry.part.original_bytes); + bindings.push(&entry.part.truncated); + } + self.conn + .prepare_cached(&sql)? + .execute(bindings.as_slice())?; + start = end; + } + Ok(()) + } + pub fn list_archive_conversations( &self, provider: Option<&str>, @@ -956,41 +1249,22 @@ impl Store { rows.collect::>() .map_err(Into::into) } +} - fn archive_content_retention( - &self, - item_id: &str, - ) -> Result> { - let mut statement = self.conn.prepare( - r#" - SELECT content_id, - text_content IS NOT NULL OR binary_content IS NOT NULL, - external_uri IS NOT NULL, - COALESCE(length(CAST(text_content AS BLOB)), 0) - + COALESCE(length(binary_content), 0), - truncated, - content_hash - FROM archive_content_parts - WHERE item_id = ?1 - "#, - )?; - let rows = statement.query_map(params![item_id], |row| { - Ok(( - row.get::<_, String>(0)?, - RetainedContentPart { - quality: ContentRetentionQuality { - materialized: row.get(1)?, - external: row.get(2)?, - stored_bytes: row.get(3)?, - truncated: row.get(4)?, - }, - content_hash: row.get(5)?, - }, - )) - })?; - rows.collect::>() - .map_err(Into::into) - } +/// Byte length standard base64 decodes to, without decoding it. +/// +/// Every encoded payload this store accepts is canonical: the archive types +/// re-encode the bytes they hashed, so the length follows from the encoding +/// alone and the decision to keep or replace a row can be made before paying +/// for the decode. +fn base64_decoded_len(encoded: &str) -> u64 { + let padding = encoded + .bytes() + .rev() + .take_while(|byte| *byte == b'=') + .count() + .min(2) as u64; + (encoded.len() as u64 / 4) * 3 - padding } fn incoming_content_should_replace( @@ -1308,6 +1582,63 @@ mod tests { assert_eq!(summaries[0].conversation_id, conversation.conversation_id); } + /// A file already imported under an earlier reconstruction must be read + /// again, or the archive keeps results the current code would never + /// produce. This is the whole purpose of the import revision, and it only + /// works if the revision takes part in the cached signature. + #[test] + fn entries_imported_under_an_earlier_revision_are_pending() { + let store = Store::in_memory().expect("store"); + let conversation = sample_conversation(); + let entry = ScanFileStateEntry { + cache_key: "/archive/thread.jsonl".to_string(), + cache_signature: "record-signature".to_string(), + }; + store + .store_archive_scan_with_code_changes( + &conversation.source_id, + std::slice::from_ref(&conversation), + std::slice::from_ref(&entry), + &[], + &[], + CoverageStatus::Unavailable, + ) + .expect("store archive scan"); + assert!( + store + .pending_archive_import_entries( + &conversation.source_id, + std::slice::from_ref(&entry) + ) + .expect("unchanged entry") + .is_empty(), + "an unchanged file was re-read" + ); + + // The same file, recorded as an earlier revision had left it. + store + .conn + .execute( + "UPDATE archive_import_state SET cache_signature = ?1 + WHERE source_id = ?2 AND cache_key = ?3", + params![ + statsai_core::hash_text(&format!("archive.v7:{}", entry.cache_signature)), + &conversation.source_id.0, + &entry.cache_key, + ], + ) + .expect("record an earlier revision"); + + assert_eq!( + store + .pending_archive_import_entries(&conversation.source_id, &[entry]) + .expect("earlier revision") + .len(), + 1, + "a file imported under an earlier revision was not re-read" + ); + } + #[test] fn artifact_metadata_changes_make_cached_archive_entry_pending() { let dir = tempfile::tempdir().expect("temp dir"); @@ -1398,6 +1729,89 @@ mod tests { assert_eq!(store.archive_stats().unwrap().conversations, 1); } + /// Re-importing an unchanged conversation must not rewrite its content. + /// + /// The rows and the search index are identical either way, so the only + /// thing a rewrite produces is work — and on a provider whose whole archive + /// re-imports whenever any part of it changes, that work is the import. + #[test] + fn unchanged_content_is_not_rewritten_on_reimport() { + let store = Store::in_memory().expect("store"); + let conversation = sample_conversation(); + let first = store + .upsert_archive_conversations(std::slice::from_ref(&conversation)) + .expect("first upsert"); + assert_eq!(first.content_parts, 2); + assert_eq!(first.binary_bytes, 4); + + let repeated = store + .upsert_archive_conversations(std::slice::from_ref(&conversation)) + .expect("identical re-import"); + + assert_eq!(repeated.content_parts, 0, "identical parts were rewritten"); + assert_eq!(repeated.binary_bytes, 0); + // Skipping the write must not lose content or searchability. + let restored = store + .archive_conversation(&conversation.conversation_id) + .expect("read") + .expect("conversation"); + assert_eq!(restored, conversation); + assert_eq!(store.search_archive("searchable", 10).unwrap().len(), 1); + let stats = store.archive_stats().expect("stats"); + assert_eq!(stats.binary_parts, 1); + assert_eq!(stats.text_parts, 1); + } + + /// A better reconstruction can correct how content is described without + /// changing a byte of it. Skipping the write because the bytes match would + /// leave the archive holding the description the old parser produced. + #[test] + fn corrected_metadata_is_stored_even_when_the_content_is_unchanged() { + let store = Store::in_memory().expect("store"); + let conversation = sample_conversation(); + store + .upsert_archive_conversations(std::slice::from_ref(&conversation)) + .expect("first upsert"); + + let mut corrected = conversation.clone(); + let text_part = &mut corrected.items[0].parts[0]; + text_part.kind = ArchiveContentKind::Json; + let binary_part = &mut corrected.items[0].parts[1]; + binary_part.mime_type = Some("image/webp".to_string()); + binary_part.name = Some("corrected.webp".to_string()); + let result = store + .upsert_archive_conversations(std::slice::from_ref(&corrected)) + .expect("metadata correction"); + + assert_eq!(result.content_parts, 2, "corrections were skipped"); + let restored = store + .archive_conversation(&conversation.conversation_id) + .expect("read") + .expect("conversation"); + assert_eq!(restored.items[0].parts, corrected.items[0].parts); + // The content itself is unchanged, so it stays searchable. + assert_eq!(store.search_archive("searchable", 10).unwrap().len(), 1); + } + + #[test] + fn base64_decoded_len_matches_the_decoded_payload() { + for bytes in [ + [0u8].as_slice(), + [0, 1].as_slice(), + [0, 1, 2].as_slice(), + [0, 1, 2, 255].as_slice(), + [7; 61].as_slice(), + ] { + let encoded = BASE64.encode(bytes); + assert_eq!( + base64_decoded_len(&encoded), + bytes.len() as u64, + "length of {encoded}" + ); + } + assert_eq!(base64_decoded_len(""), 0); + } + #[test] fn reduced_item_rescan_preserves_richer_existing_parts() { let store = Store::in_memory().expect("store"); diff --git a/crates/statsai-store/src/lib.rs b/crates/statsai-store/src/lib.rs index 838f742..e91467b 100644 --- a/crates/statsai-store/src/lib.rs +++ b/crates/statsai-store/src/lib.rs @@ -458,6 +458,25 @@ pub struct Store { conn: Connection, } +/// Restores the store's commit durability when dropped. +/// +/// Held for the length of a bulk import. Restoring on drop rather than at the +/// end of the import keeps a failed import from leaving a long-lived process +/// writing everything else at reduced durability. +pub struct BulkImportDurability<'a> { + store: &'a Store, + restore_to: i64, +} + +impl Drop for BulkImportDurability<'_> { + fn drop(&mut self) { + let _ = self + .store + .conn + .execute_batch(&format!("PRAGMA synchronous = {}", self.restore_to)); + } +} + impl Store { /// Opens a store and applies migrations. /// @@ -478,10 +497,41 @@ impl Store { conn.busy_timeout(SQLITE_BUSY_TIMEOUT)?; let store = Self { conn }; store.migrate()?; + store.configure_connection()?; store.conn.execute_batch("PRAGMA optimize=0x10002;")?; Ok(store) } + /// Applies the per-connection settings the write paths rely on. + /// + /// Commit durability is deliberately left alone: this connection also + /// writes state that exists nowhere else — verifications a person entered, + /// subscriptions, account assignments, privacy identity — and none of that + /// can be collected again from local files. Relaxing durability is scoped + /// to the imports that can, in [`Store::relax_durability_for_bulk_import`]. + /// + /// The page cache is raised because these archives are far larger than the + /// 2MB default, which turns index maintenance into a stream of single-page + /// reads. + fn configure_connection(&self) -> Result<()> { + self.conn.execute_batch( + "PRAGMA cache_size = -65536; + PRAGMA temp_store = MEMORY; + CREATE TEMP TABLE IF NOT EXISTS incoming_records ( + source_record_id TEXT, + item_id TEXT NOT NULL + ); + CREATE INDEX IF NOT EXISTS incoming_records_source_idx + ON incoming_records (source_record_id); + CREATE INDEX IF NOT EXISTS incoming_records_item_idx + ON incoming_records (item_id);", + )?; + // Batched writes issue one statement per batch size, and the archive + // paths alternate between a handful of them. + self.conn.set_prepared_statement_cache_capacity(64); + Ok(()) + } + /// Opens an independent connection to the same file-backed store. /// /// This is useful for background readers that must inspect cache state @@ -521,10 +571,39 @@ impl Store { }; store.conn.busy_timeout(SQLITE_BUSY_TIMEOUT)?; store.migrate()?; + store.configure_connection()?; store.conn.execute_batch("PRAGMA optimize=0x10002;")?; Ok(store) } + /// Relaxes commit durability until the returned guard is dropped. + /// + /// In WAL mode `synchronous = NORMAL` cannot corrupt the database; it only + /// means a power loss may cost the most recently committed transactions. + /// That is an acceptable trade for importing a provider's archive, because + /// each file's rows and the cache entry recording it commit together, so a + /// lost commit is collected again on the next run rather than going + /// silently missing. + /// + /// It is not an acceptable trade for the rest of the store, which holds + /// state that no local file can reproduce, so the relaxation is scoped to + /// the import rather than applied to the connection for good. + /// + /// # Errors + /// + /// Returns an error if SQLite rejects the durability change. + pub fn relax_durability_for_bulk_import(&self) -> Result> { + let restore_to = self + .conn + .query_row("PRAGMA synchronous", [], |row| row.get::<_, i64>(0)) + .context("read commit durability")?; + self.conn.execute_batch("PRAGMA synchronous = NORMAL")?; + Ok(BulkImportDurability { + store: self, + restore_to, + }) + } + fn with_immediate_transaction(&self, operation: impl FnOnce() -> Result) -> Result { if !self.conn.is_autocommit() { return operation(); @@ -5689,6 +5768,40 @@ mod tests { }; use std::path::Path; + /// Importing an archive may trade durability for speed because a lost + /// commit is simply collected again. The rest of the store holds work that + /// no local file can reproduce, so the trade must not outlive the import. + #[test] + fn relaxed_durability_is_scoped_to_the_bulk_import() { + let dir = tempfile::tempdir().expect("temp dir"); + let store = Store::open(&dir.path().join("statsai.sqlite")).expect("store"); + let durability = || { + store + .conn + .query_row("PRAGMA synchronous", [], |row| row.get::<_, i64>(0)) + .expect("read durability") + }; + + let opened_with = durability(); + assert_ne!( + opened_with, 0, + "a store must never open with durability disabled" + ); + + { + let _relaxed = store + .relax_durability_for_bulk_import() + .expect("relax durability"); + assert_eq!(durability(), 1, "the import did not get NORMAL durability"); + } + + assert_eq!( + durability(), + opened_with, + "durability stayed relaxed after the import" + ); + } + #[test] #[cfg(unix)] fn open_restricts_store_directory_and_database_permissions() { diff --git a/crates/statsai/src/main.rs b/crates/statsai/src/main.rs index 9b54210..790fe2e 100644 --- a/crates/statsai/src/main.rs +++ b/crates/statsai/src/main.rs @@ -1,4 +1,4 @@ -use anyhow::{bail, Context, Result}; +use anyhow::{bail, ensure, Context, Result}; use chrono::{DateTime, Datelike, Duration, NaiveDate, Utc}; use clap::{Args, Parser, Subcommand}; use serde::{Deserialize, Serialize}; @@ -6,8 +6,8 @@ use serde_json::{json, Value}; #[cfg(test)] use statsai_adapters::VerifiedSubscriptionState; use statsai_adapters::{ - adapter_for_provider, default_adapters, ProviderAdapter, ScanCandidateFile, ScanDiagnostics, - ScanOptions, VerifiedSourceObservation, + adapter_for_provider, default_adapters, ArchiveScan, ProviderAdapter, ScanCandidateFile, + ScanDiagnostics, ScanOptions, VerifiedSourceObservation, }; #[cfg(test)] use statsai_adapters::{SourceIdentityInference, VerifiedSourceState}; @@ -6089,18 +6089,71 @@ fn collect_conversations( ) -> Result<()> { let canonical_provider_filter = canonical_conversation_provider_filter(provider_filter)?; let configured_sources = store.list_sources()?; - let mut sources_collected = 0u64; - let mut total_conversations = 0u64; - let mut total_items = 0u64; - let mut total_parts = 0u64; - let mut total_binary_bytes = 0u64; - let mut total_missing = 0u64; + // Reduced durability covers the imports and nothing after them. An import + // rebuilds from the provider's own files, so a commit lost to a power cut + // costs a re-collect. The code-change refresh that follows carries metrics + // forward for commits too old for Git to be rescanned for, and writes them + // across several transactions, so it keeps the store's normal durability. + // Anything added below belongs after this block unless it is as reproducible + // as an import. + let totals = { + let _durability = store.relax_durability_for_bulk_import()?; + collect_archive_sources( + store, + canonical_provider_filter, + &configured_sources, + no_cache, + verbose, + )? + }; + println!( + "archive collection: sources={} conversations={} items={} parts={} binary_bytes={} missing={}", + totals.sources, + totals.conversations, + totals.items, + totals.parts, + totals.binary_bytes, + totals.missing, + ); + let code_changes = store.refresh_code_changes(device_id)?; + println!( + "code changes: trace_edits={} repositories={} commits={} trace_matched={} metrics={} trace_coverage={:?} git_coverage={:?}", + code_changes.trace_edits, + code_changes.repositories, + code_changes.commits, + code_changes.matches, + code_changes.metrics, + code_changes.trace_coverage, + code_changes.git_coverage, + ); + Ok(()) +} + +/// What one `conversation collect` run imported. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +struct ArchiveCollectionTotals { + sources: u64, + conversations: u64, + items: u64, + parts: u64, + binary_bytes: u64, + missing: u64, +} +/// Imports every archive source the filter admits. +fn collect_archive_sources( + store: &Store, + canonical_provider_filter: Option<&str>, + configured_sources: &[SourceLocation], + no_cache: bool, + verbose: bool, +) -> Result { + let mut totals = ArchiveCollectionTotals::default(); for adapter in default_adapters() { if canonical_provider_filter.is_some_and(|provider| provider != adapter.provider()) { continue; } - for source in scan_sources_for_adapter(adapter.as_ref(), &configured_sources) { + for source in scan_sources_for_adapter(adapter.as_ref(), configured_sources) { let candidates = adapter.archive_scan_candidates(&source)?; let entries = scan_file_state_entries(&candidates); // An empty candidate list from an unreachable root — an unmounted @@ -6145,12 +6198,12 @@ fn collect_conversations( &pending, verbose, )?; - sources_collected += 1; - total_conversations += collected.conversations; - total_items += collected.items; - total_parts += collected.parts; - total_binary_bytes += collected.binary_bytes; - total_missing += collected.missing; + totals.sources += 1; + totals.conversations += collected.conversations; + totals.items += collected.items; + totals.parts += collected.parts; + totals.binary_bytes += collected.binary_bytes; + totals.missing += collected.missing; if verbose { println!( "{} {}: files={} conversations={} items={} parts={} binary_bytes={} missing={} invalid_records={}", @@ -6167,27 +6220,7 @@ fn collect_conversations( } } } - println!( - "archive collection: sources={} conversations={} items={} parts={} binary_bytes={} missing={}", - sources_collected, - total_conversations, - total_items, - total_parts, - total_binary_bytes, - total_missing, - ); - let code_changes = store.refresh_code_changes(device_id)?; - println!( - "code changes: trace_edits={} repositories={} commits={} trace_matched={} metrics={} trace_coverage={:?} git_coverage={:?}", - code_changes.trace_edits, - code_changes.repositories, - code_changes.commits, - code_changes.matches, - code_changes.metrics, - code_changes.trace_coverage, - code_changes.git_coverage, - ); - Ok(()) + Ok(totals) } #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] @@ -6213,65 +6246,237 @@ fn collect_archive_source_entries( .iter() .map(|candidate| (candidate.cache_key.as_str(), candidate)) .collect::>(); + // Sized once: grouping and the reconstruction budget both ask for this, and + // a source can hold thousands of files. + let source_bytes = pending + .iter() + .map(|entry| { + candidates_by_key + .get(entry.cache_key.as_str()) + .and_then(|candidate| std::fs::metadata(&candidate.path).ok()) + .map_or(0, |metadata| metadata.len()) + }) + .collect::>(); let mut collected = ArchiveSourceCollection::default(); - for (index, entry) in pending.iter().enumerate() { - let candidate = candidates_by_key.get(entry.cache_key.as_str()).copied(); - let candidate_bytes = candidate - .and_then(|candidate| std::fs::metadata(&candidate.path).ok()) - .map_or(0, |metadata| metadata.len()); - let report_candidate = - verbose && (index == 0 || (index + 1) % 25 == 0 || candidate_bytes >= 16 * 1024 * 1024); - if report_candidate { - println!( - "{} {}: collecting file {}/{} ({} bytes) {}", - adapter.provider(), - preview_path_label(source), - index + 1, - pending.len(), - candidate_bytes, - candidate - .and_then(|candidate| candidate.path.file_name()) - .and_then(|name| name.to_str()) - .unwrap_or(entry.cache_key.as_str()), - ); - } - + let mut index = 0; + while index < pending.len() { + let group = archive_collection_group(&source_bytes, index); + let group_entries = &pending[index..index + group]; + + // Reading and reconstructing a transcript is independent per file and + // is what the wall clock is mostly spent on, so the group is parsed on + // several threads. The results are written in file order on this + // thread: two files can describe the same conversation, and the record + // that wins must not depend on which thread finished first. let collect_started = Instant::now(); - let selected = HashSet::from([entry.cache_key.clone()]); - let scan = adapter.collect_archive(source, Some(&selected))?; + let scans = parse_archive_group( + adapter, + source, + group_entries, + &source_bytes[index..index + group], + ); let collect_elapsed = collect_started.elapsed(); + // Reconstruction stops early when it is already holding enough content, + // so the run advances by what came back rather than what was offered. + // Advancing by anything else would step over a file that was never + // stored, and a run that reconstructed nothing has to say so rather + // than skip the file or ask for it forever. + let group = scans.len(); + ensure!( + group > 0, + "reconstructed no archive files from {} at file {}", + preview_path_label(source), + index + 1, + ); + let store_started = Instant::now(); - let write = store.store_archive_scan_with_code_changes( - &source.source_id, - &scan.conversations, - std::slice::from_ref(entry), - &scan.artifact_dependencies, - &scan.trace_edits, - scan.trace_coverage, - )?; + // One transaction per file: a file's rows and the cache entry that + // records it are committed together, and a file that fails stops the + // run only after every file before it has been stored. + for (entry, scan) in group_entries.iter().zip(scans) { + let scan = scan?; + let write = store.store_archive_scan_with_code_changes( + &source.source_id, + &scan.conversations, + std::slice::from_ref(entry), + &scan.artifact_dependencies, + &scan.trace_edits, + scan.trace_coverage, + )?; + collected.files += scan.diagnostics.files_scanned; + collected.conversations += write.conversations; + collected.items += write.items; + collected.parts += write.content_parts; + collected.binary_bytes += write.binary_bytes; + collected.missing += scan.diagnostics.missing_content; + collected.invalid_records += scan.diagnostics.invalid_records; + } let store_elapsed = store_started.elapsed(); - collected.files += scan.diagnostics.files_scanned; - collected.conversations += write.conversations; - collected.items += write.items; - collected.parts += write.content_parts; - collected.binary_bytes += write.binary_bytes; - collected.missing += scan.diagnostics.missing_content; - collected.invalid_records += scan.diagnostics.invalid_records; - if report_candidate { + if verbose { println!( - "{} {}: completed file {}/{} collect={:.1}s store={:.1}s", + "{} {}: collected files {}-{}/{} collect={:.1}s store={:.1}s", adapter.provider(), preview_path_label(source), index + 1, + index + group, pending.len(), collect_elapsed.as_secs_f64(), store_elapsed.as_secs_f64(), ); } + index += group; } Ok(collected) } +/// Files parsed together before their results are written. +const ARCHIVE_COLLECTION_GROUP_FILES: usize = 16; +/// Source bytes a group stops growing at. +/// +/// A group is held in memory in full, and one transcript can be tens of +/// megabytes that expand further once reconstructed, so the bound is on the +/// input size rather than the file count alone. +const ARCHIVE_COLLECTION_GROUP_BYTES: u64 = 64 * 1024 * 1024; + +/// Number of files to take for the group starting at `index`, always at least +/// one so that a single oversized file still makes progress. +fn archive_collection_group(source_bytes: &[u64], index: usize) -> usize { + let mut group = 0; + let mut bytes = 0u64; + while index + group < source_bytes.len() && group < ARCHIVE_COLLECTION_GROUP_FILES { + bytes = bytes.saturating_add(source_bytes[index + group]); + group += 1; + if bytes >= ARCHIVE_COLLECTION_GROUP_BYTES { + break; + } + } + group.max(1) +} + +/// Reconstructed content a group holds before its results are stored. +/// +/// Source size does not predict this: a one-line transcript can name a local +/// artifact of tens of megabytes, which is materialized and carried as base64. +/// Workers therefore stop taking new files once the group has this much +/// outstanding, and the files they did not reach are simply the next group. +const ARCHIVE_COLLECTION_RETAINED_BYTES: usize = 192 * 1024 * 1024; +/// Files reconstructed at once. +/// +/// A file's real cost is only known once it has been read: a one-line +/// transcript may name artifacts many times its own size, so no estimate taken +/// beforehand can bound it. What can be bounded is how many files are being +/// read at once, which is what keeps the worst case a small multiple of the +/// largest single file rather than a multiple of the core count. +const ARCHIVE_COLLECTION_IN_FLIGHT: usize = 4; +/// Least a file is charged against the budget while it is being reconstructed. +/// +/// Set so that the budget alone admits [`ARCHIVE_COLLECTION_IN_FLIGHT`] files +/// of unknown size, and fewer once one of them is known to be large. +const ARCHIVE_COLLECTION_FILE_RESERVE: usize = + ARCHIVE_COLLECTION_RETAINED_BYTES / ARCHIVE_COLLECTION_IN_FLIGHT; + +/// How much of the budget a group has outstanding, and which file is next. +/// +/// Both move together under one lock: a worker that decided to take a file +/// having seen the budget empty, while its peers were deciding the same thing, +/// is how a bound on retained content stops being one. +struct ArchiveGroupClaim { + next: usize, + retained: usize, +} + +/// Reconstructs a leading run of `entries`, in parallel, preserving their order. +/// +/// Returns one result per file reconstructed, which may be fewer than were +/// offered; the caller advances by however many came back. A file that cannot +/// be read is reported in its own position rather than aborting the run, so the +/// caller can still store every file that precedes it. Collection is resumable +/// precisely because a file that was stored stays stored when a later one +/// fails. +fn parse_archive_group( + adapter: &dyn ProviderAdapter, + source: &SourceLocation, + entries: &[ScanFileStateEntry], + source_bytes: &[u64], +) -> Vec> { + let collect_one = |entry: &ScanFileStateEntry| { + let selected = HashSet::from([entry.cache_key.clone()]); + adapter.collect_archive(source, Some(&selected)) + }; + if entries.len() == 1 { + return vec![collect_one(&entries[0])]; + } + let workers = std::thread::available_parallelism() + .map_or(ARCHIVE_COLLECTION_IN_FLIGHT, std::num::NonZeroUsize::get) + .min(entries.len()) + .min(ARCHIVE_COLLECTION_IN_FLIGHT); + let claim = std::sync::Mutex::new(ArchiveGroupClaim { + next: 0, + retained: 0, + }); + let results = (0..entries.len()) + .map(|_| std::sync::Mutex::new(None)) + .collect::>(); + std::thread::scope(|scope| { + for _ in 0..workers { + scope.spawn(|| loop { + let (index, reserved) = { + let mut claim = claim.lock().expect("archive group claim"); + let index = claim.next; + if index >= entries.len() { + return; + } + let reserved = usize::try_from(source_bytes.get(index).copied().unwrap_or(0)) + .unwrap_or(usize::MAX) + .max(ARCHIVE_COLLECTION_FILE_RESERVE); + // The first file is always taken, so one file larger than + // the whole budget still makes progress. + if index > 0 + && claim.retained.saturating_add(reserved) + > ARCHIVE_COLLECTION_RETAINED_BYTES + { + return; + } + claim.next = index + 1; + claim.retained = claim.retained.saturating_add(reserved); + (index, reserved) + }; + let scan = collect_one(&entries[index]); + // What the file actually costs replaces what it was charged. + let actual = scan.as_ref().map_or(0, archive_scan_retained_bytes); + { + let mut claim = claim.lock().expect("archive group claim"); + claim.retained = claim + .retained + .saturating_sub(reserved) + .saturating_add(actual); + } + *results[index].lock().expect("archive scan slot") = Some(scan); + }); + } + }); + // Indices are handed out in order and every one handed out is reconstructed, + // so the results form a leading run. Stopping at the first gap keeps that + // true whatever the workers did. + results + .into_iter() + .map_while(|slot| slot.into_inner().expect("archive scan slot")) + .collect() +} + +/// Reconstructed content a scan is holding in memory. +fn archive_scan_retained_bytes(scan: &ArchiveScan) -> usize { + scan.conversations + .iter() + .flat_map(|conversation| &conversation.items) + .flat_map(|item| &item.parts) + .map(|part| { + part.text.as_ref().map_or(0, String::len) + + part.data_base64.as_ref().map_or(0, String::len) + }) + .sum() +} + fn canonical_conversation_provider_filter(provider: Option<&str>) -> Result> { provider .map(|provider| { @@ -8973,6 +9178,197 @@ mod tests { assert_eq!(pending, vec![entries[1].clone()]); } + /// Files are reconstructed on several threads but must be handed back in + /// the order they were listed: two files can describe the same + /// conversation, and which record wins must not depend on scheduling. + #[test] + fn archive_group_parsing_preserves_file_order() { + struct OrderedArchiveAdapter; + + impl ProviderAdapter for OrderedArchiveAdapter { + fn id(&self) -> &'static str { + "ordered-archive-test" + } + fn version(&self) -> &'static str { + "0" + } + fn provider(&self) -> &'static str { + "archive_test" + } + fn discover(&self) -> Vec { + Vec::new() + } + fn scan_candidates(&self, _source: &SourceLocation) -> Result> { + Ok(Vec::new()) + } + fn scan( + &self, + _source: &SourceLocation, + _options: &ScanOptions, + ) -> Result { + Ok(statsai_adapters::AdapterScan::default()) + } + fn collect_archive( + &self, + _source: &SourceLocation, + selected_cache_keys: Option<&HashSet>, + ) -> Result { + let selected = selected_cache_keys + .and_then(|keys| keys.iter().next()) + .context("selected archive cache key")?; + // The earlier files are the slow ones, so a run that returned + // results as they arrived would reorder them. + let index: u64 = selected.parse().context("cache key index")?; + std::thread::sleep(std::time::Duration::from_millis(40 - index * 4)); + let mut scan = statsai_adapters::ArchiveScan::default(); + scan.diagnostics.files_scanned = index; + Ok(scan) + } + } + + let source = SourceLocation::local_adapter( + "archive_test", + "ordered-archive-test", + "0", + Path::new("/tmp/archive-order-test"), + LocationOrigin::Configured, + ); + let entries = (0..10) + .map(|index| ScanFileStateEntry { + cache_key: index.to_string(), + cache_signature: format!("signature-{index}"), + }) + .collect::>(); + + let scans = parse_archive_group( + &OrderedArchiveAdapter, + &source, + &entries, + &vec![0; entries.len()], + ); + + let order = scans + .into_iter() + .map(|scan| scan.expect("collected archive").diagnostics.files_scanned) + .collect::>(); + assert_eq!(order, (0..10).collect::>()); + } + + /// A transcript's size on disk says nothing about how much content it + /// materializes, so reconstruction stops taking new files once it is + /// holding enough. The files it did not reach must come back for the + /// caller to collect next, never be reported as done. + #[test] + fn archive_group_parsing_stops_before_holding_too_much_content() { + struct HeavyArchiveAdapter; + + impl ProviderAdapter for HeavyArchiveAdapter { + fn id(&self) -> &'static str { + "heavy-archive-test" + } + fn version(&self) -> &'static str { + "0" + } + fn provider(&self) -> &'static str { + "archive_test" + } + fn discover(&self) -> Vec { + Vec::new() + } + fn scan_candidates(&self, _source: &SourceLocation) -> Result> { + Ok(Vec::new()) + } + fn scan( + &self, + _source: &SourceLocation, + _options: &ScanOptions, + ) -> Result { + Ok(statsai_adapters::AdapterScan::default()) + } + fn collect_archive( + &self, + _source: &SourceLocation, + _selected_cache_keys: Option<&HashSet>, + ) -> Result { + // A tiny record naming an artifact that materializes far + // larger than the file it came from. + let mut conversation = ArchiveConversation { + schema_version: statsai_core::ARCHIVE_CONVERSATION_SCHEMA_VERSION.to_string(), + conversation_id: "conv_heavy".to_string(), + provider: "archive_test".to_string(), + source_id: statsai_core::SourceId("heavy".to_string()), + native_conversation_id: "heavy".to_string(), + title: None, + project: None, + started_at: None, + updated_at: None, + completeness: statsai_core::ArchiveCompleteness::Complete, + missing_content_count: 0, + missing_content_scope_id: None, + discarded_source_record_ids: Vec::new(), + superseded_conversation_ids: Vec::new(), + items: Vec::new(), + }; + let item_id = "item_heavy".to_string(); + conversation.items.push(statsai_core::ArchiveItem { + item_id: item_id.clone(), + native_item_id: None, + source_record_id: None, + ordinal: 0, + kind: statsai_core::ArchiveItemKind::Message, + role: None, + created_at: None, + model: None, + tool_name: None, + tool_call_id: None, + status: None, + usage: None, + parts_authoritative: true, + parts: vec![statsai_core::ArchiveContentPart::text( + statsai_core::archive_content_id(&item_id, 0), + 0, + ArchiveContentKind::Text, + "x".repeat(ARCHIVE_COLLECTION_RETAINED_BYTES / 4), + )], + }); + // Held long enough that every worker has tried to claim before + // any capacity is released. + std::thread::sleep(std::time::Duration::from_millis(250)); + let mut scan = statsai_adapters::ArchiveScan::default(); + scan.conversations.push(conversation); + Ok(scan) + } + } + + let source = SourceLocation::local_adapter( + "archive_test", + "heavy-archive-test", + "0", + Path::new("/tmp/archive-heavy-test"), + LocationOrigin::Configured, + ); + let entries = (0..ARCHIVE_COLLECTION_GROUP_FILES) + .map(|index| ScanFileStateEntry { + cache_key: index.to_string(), + cache_signature: format!("signature-{index}"), + }) + .collect::>(); + // A quarter of the budget each, so only a few may be outstanding. + let source_bytes = + vec![ARCHIVE_COLLECTION_RETAINED_BYTES as u64 / 4; ARCHIVE_COLLECTION_GROUP_FILES]; + + let scans = parse_archive_group(&HeavyArchiveAdapter, &source, &entries, &source_bytes); + + // Every worker reaches the budget before any of them finishes, which is + // exactly when a check that does not reserve lets all of them through. + assert!(!scans.is_empty(), "no file was reconstructed"); + assert!( + scans.len() <= 4, + "budget did not gate concurrent claims: {} files were taken", + scans.len() + ); + } + #[test] fn provider_aliases_match_canonical_provider() { assert!(provider_matches("claude_code", "claude"));