diff --git a/crates/archivr-core/src/database.rs b/crates/archivr-core/src/database.rs index 34b1d2e..eb7372a 100644 --- a/crates/archivr-core/src/database.rs +++ b/crates/archivr-core/src/database.rs @@ -349,7 +349,7 @@ pub fn initialize_schema(conn: &Connection) -> Result<()> { summary_uid TEXT NOT NULL UNIQUE, entry_id INTEGER NOT NULL REFERENCES archived_entries(id) ON DELETE CASCADE, provider_kind TEXT NOT NULL, - provider_model TEXT, + provider_model TEXT NOT NULL DEFAULT '', prompt_version TEXT NOT NULL, input_sha256 TEXT NOT NULL, status TEXT NOT NULL CHECK(status IN ('pending','running','completed','failed')), @@ -357,11 +357,12 @@ pub fn initialize_schema(conn: &Connection) -> Result<()> { error_text TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, - completed_at TEXT, - UNIQUE(entry_id, provider_kind, provider_model, prompt_version, input_sha256) + completed_at TEXT ); CREATE INDEX IF NOT EXISTS idx_entry_summaries_entry_updated ON entry_summaries(entry_id, updated_at DESC); + CREATE INDEX IF NOT EXISTS idx_entry_summaries_cache_lookup + ON entry_summaries(entry_id, provider_kind, provider_model, prompt_version, input_sha256, status); CREATE INDEX IF NOT EXISTS idx_capture_jobs_status ON capture_jobs(status); CREATE INDEX IF NOT EXISTS idx_archive_run_items_run_id ON archive_run_items(run_id); CREATE INDEX IF NOT EXISTS idx_archived_entries_source_identity_id ON archived_entries(source_identity_id); @@ -483,6 +484,52 @@ pub fn initialize_schema(conn: &Connection) -> Result<()> { [], ); + // Summary attempts used to be unique by cache key, which meant forced + // regeneration erased the last completed result. Rebuild that small table + // without the cache-key constraint while retaining all existing rows. + let summary_table_sql: Option = conn + .query_row( + "SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'entry_summaries'", + [], + |row| row.get(0), + ) + .optional()?; + if summary_table_sql + .as_deref() + .is_some_and(|sql| sql.contains("UNIQUE(entry_id, provider_kind, provider_model, prompt_version, input_sha256)")) + { + conn.execute_batch( + "BEGIN; + CREATE TABLE entry_summaries_rebuilt ( + id INTEGER PRIMARY KEY, + summary_uid TEXT NOT NULL UNIQUE, + entry_id INTEGER NOT NULL REFERENCES archived_entries(id) ON DELETE CASCADE, + provider_kind TEXT NOT NULL, + provider_model TEXT NOT NULL DEFAULT '', + prompt_version TEXT NOT NULL, + input_sha256 TEXT NOT NULL, + status TEXT NOT NULL CHECK(status IN ('pending','running','completed','failed')), + summary_text TEXT, + error_text TEXT, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + completed_at TEXT + ); + INSERT INTO entry_summaries_rebuilt + SELECT id, summary_uid, entry_id, provider_kind, COALESCE(provider_model, ''), + prompt_version, input_sha256, status, summary_text, error_text, + created_at, updated_at, completed_at + FROM entry_summaries; + DROP TABLE entry_summaries; + ALTER TABLE entry_summaries_rebuilt RENAME TO entry_summaries; + CREATE INDEX idx_entry_summaries_entry_updated + ON entry_summaries(entry_id, updated_at DESC); + CREATE INDEX idx_entry_summaries_cache_lookup + ON entry_summaries(entry_id, provider_kind, provider_model, prompt_version, input_sha256, status); + COMMIT;", + )?; + } + Ok(()) } @@ -1394,6 +1441,21 @@ pub fn fail_stalled_capture_jobs(conn: &Connection) -> Result { Ok(n) } +/// Marks summary attempts left pending or running by a previous process as +/// failed. The message is deliberately stable and does not expose internals. +pub fn fail_stalled_entry_summaries(conn: &Connection) -> Result { + let now = now_timestamp(); + conn.execute( + "UPDATE entry_summaries + SET status = 'failed', + error_text = 'Summary generation was interrupted by a server restart.', + updated_at = ?1 + WHERE status IN ('pending', 'running')", + [now], + ) + .map_err(Into::into) +} + // ── entry_summaries ──────────────────────────────────────────────────────── // // Summaries are a regenerable child record of an entry, never a column on @@ -1441,12 +1503,10 @@ pub fn entry_id_for_uid(conn: &Connection, entry_uid: &str) -> Result Result> { + conn.query_row( + &format!( + "{ENTRY_SUMMARY_COLS} WHERE s.entry_id = ?1 AND s.status = 'completed' + ORDER BY s.completed_at DESC, s.updated_at DESC, s.id DESC LIMIT 1" + ), + [entry_id], + map_entry_summary, + ) + .optional() + .map_err(Into::into) +} + pub fn create_archive_run( conn: &Connection, created_by_user_id: i64, @@ -4464,33 +4532,72 @@ mod tests { } #[test] - fn upsert_pending_entry_summary_reuses_the_row_for_an_identical_cache_key() { + fn regenerating_a_completed_summary_creates_a_new_attempt_and_preserves_completion() { let c = conn(); let entry = create_entry_fixture(&c, "private", None, None); let first = upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "sha").unwrap(); update_entry_summary_status(&c, &first, "completed", Some("old"), None).unwrap(); - // Regenerating the same key must reset the row in place, not add a second. + // Regeneration must retain the finished result while a new attempt runs. let second = upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "sha").unwrap(); - assert_eq!(first, second, "same cache key must keep the same summary_uid"); + assert_ne!(first, second, "regeneration needs a distinct attempt uid"); let n: i64 = c .query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0)) .unwrap(); - assert_eq!(n, 1); + assert_eq!(n, 2); - let rec = get_entry_summary_by_uid(&c, &first).unwrap().unwrap(); + let completed = get_entry_summary_by_uid(&c, &first).unwrap().unwrap(); + assert_eq!(completed.status, "completed"); + assert_eq!(completed.summary_text.as_deref(), Some("old")); + let rec = get_entry_summary_by_uid(&c, &second).unwrap().unwrap(); assert_eq!(rec.status, "pending"); - assert!(rec.summary_text.is_none(), "stale text must be cleared"); + update_entry_summary_status(&c, &second, "failed", None, Some("boom")).unwrap(); + assert_eq!( + latest_completed_entry_summary(&c, entry.id) + .unwrap() + .unwrap() + .summary_uid, + first, + "a failed regeneration must not replace the previous completed result" + ); } #[test] - fn entry_summaries_with_no_model_still_dedupe() { - // Regression guard: a NULL provider_model would compare as distinct in - // SQLite's UNIQUE index, so the CLI providers (which have no model) - // would accumulate a new row on every regenerate. We store '' instead. + fn fail_stalled_entry_summaries_marks_pending_and_running_with_restart_message() { + let c = conn(); + let entry = create_entry_fixture(&c, "private", None, None); + let pending = upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "a") + .unwrap(); + let running = upsert_pending_entry_summary(&c, entry.id, "codex_cli", None, "v1", "b") + .unwrap(); + let completed = upsert_pending_entry_summary(&c, entry.id, "anthropic_http", None, "v1", "c") + .unwrap(); + update_entry_summary_status(&c, &running, "running", None, None).unwrap(); + update_entry_summary_status(&c, &completed, "completed", Some("done"), None).unwrap(); + + assert_eq!(fail_stalled_entry_summaries(&c).unwrap(), 2); + for uid in [&pending, &running] { + let row = get_entry_summary_by_uid(&c, uid).unwrap().unwrap(); + assert_eq!(row.status, "failed"); + assert_eq!( + row.error_text.as_deref(), + Some("Summary generation was interrupted by a server restart.") + ); + assert!(!row.updated_at.is_empty()); + } + assert_eq!( + get_entry_summary_by_uid(&c, &completed).unwrap().unwrap().status, + "completed" + ); + } + + #[test] + fn entry_summaries_with_no_model_preserve_none_across_attempts() { + // CLI providers have no explicit model, so the stored empty-string + // sentinel must always map back to None even when attempts accumulate. let c = conn(); let entry = create_entry_fixture(&c, "private", None, None); upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "sha").unwrap(); @@ -4498,7 +4605,7 @@ mod tests { let n: i64 = c .query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0)) .unwrap(); - assert_eq!(n, 1); + assert_eq!(n, 2); // …and it reads back as None, not as an empty-string model. let rec = latest_entry_summary(&c, entry.id).unwrap().unwrap(); assert_eq!(rec.provider_model, None); diff --git a/crates/archivr-core/src/summarizer.rs b/crates/archivr-core/src/summarizer.rs index 68ea9f4..b98fcda 100644 --- a/crates/archivr-core/src/summarizer.rs +++ b/crates/archivr-core/src/summarizer.rs @@ -1184,14 +1184,34 @@ pub fn summarize_entry( prompt_version, &input.input_sha256, )?; - database::update_entry_summary_status(&conn, &summary_uid, "running", None, None)?; + summarize_prebuilt_entry( + archive_paths, + input, + &summary_uid, + provider, + ) +} + +/// Runs a previously validated and claimed summary attempt. +/// +/// The server uses this after doing its preflight in a blocking task, avoiding +/// a second filesystem/SQLite extraction and ensuring provider output updates +/// the exact pending row returned to the caller. +pub fn summarize_prebuilt_entry( + archive_paths: &ArchivePaths, + input: SummaryInput, + summary_uid: &str, + provider: &dyn SummaryProvider, +) -> Result { + let conn = database::open_or_initialize(&archive_paths.archive_path)?; + database::update_entry_summary_status(&conn, summary_uid, "running", None, None)?; match provider.summarize(&input.request) { Ok(output) => { let text = normalize_summary_json(&output.text); database::update_entry_summary_status( &conn, - &summary_uid, + summary_uid, "completed", Some(&text), None, @@ -1201,7 +1221,7 @@ pub fn summarize_entry( let msg = format!("{e:#}"); database::update_entry_summary_status( &conn, - &summary_uid, + summary_uid, "failed", None, Some(&msg), @@ -1210,7 +1230,7 @@ pub fn summarize_entry( } } - database::get_entry_summary_by_uid(&conn, &summary_uid)? + database::get_entry_summary_by_uid(&conn, summary_uid)? .ok_or_else(|| anyhow!("summary row disappeared after write")) } diff --git a/crates/archivr-server/src/main.rs b/crates/archivr-server/src/main.rs index b191e9e..72be507 100644 --- a/crates/archivr-server/src/main.rs +++ b/crates/archivr-server/src/main.rs @@ -98,6 +98,17 @@ async fn main() -> Result<()> { ), _ => {} } + match archivr_core::database::fail_stalled_entry_summaries(&conn) { + Ok(n) if n > 0 => eprintln!( + "info: marked {n} stalled summary attempt(s) as failed in '{}'", + archive.id + ), + Err(e) => eprintln!( + "warn: stalled summary cleanup failed for '{}': {e:#}", + archive.id + ), + _ => {} + } } } diff --git a/crates/archivr-server/src/routes.rs b/crates/archivr-server/src/routes.rs index 0f220e4..2d3b14c 100644 --- a/crates/archivr-server/src/routes.rs +++ b/crates/archivr-server/src/routes.rs @@ -590,70 +590,70 @@ async fn request_entry_summary_handler( include_images: body.include_images, }; - let conn = database::open_or_initialize(&mounted.archive_path)?; - let entry_id = database::entry_id_for_uid(&conn, &entry_uid)? - .ok_or(ApiError::not_found("entry not found"))?; - - // 2. Extract the same content the summarizer will feed the model, so the - // digest below is the identical cache key summarize_entry will compute. - let input = summarizer::build_summary_input(&archive_paths, &entry_uid, summary_options) - .map_err(|e| { - if summarizer::is_unsupported_summary_content_error(&e) { - ApiError::bad_request(summarizer::UNSUPPORTED_SUMMARY_CONTENT_MESSAGE) - } else { - ApiError::bad_request(&format!("{e:#}")) - } - })?; - - // 3. Cache hit: identical entry + provider + model + prompt + input. - if !body.force { - if let Some(existing) = database::find_entry_summary( - &conn, - entry_id, - provider.kind(), - provider.model(), - summarizer::PROMPT_VERSION, - &input.input_sha256, - )? { - if existing.status == "completed" { - return Ok(( - StatusCode::OK, - serde_json::to_value(&existing) - .map(Json) - .map_err(|e| ApiError::internal(&e.to_string()))?, - )); + // 2. Preflight extraction and SQLite cache/attempt work are synchronous + // core operations, so keep them off the Axum runtime. This also means the + // input claimed here is passed directly to the provider worker below. + let preflight_paths = archive_paths.clone(); + let preflight_uid = entry_uid.clone(); + let provider_kind = provider.kind().to_string(); + let provider_model = provider.model().map(str::to_string); + let force = body.force; + enum PreflightOutcome { + Cached(database::EntrySummaryRecord), + Pending { input: summarizer::SummaryInput, summary_uid: String }, + } + let outcome = tokio::task::spawn_blocking(move || -> anyhow::Result { + let input = summarizer::build_summary_input(&preflight_paths, &preflight_uid, summary_options)?; + let conn = database::open_or_initialize(&preflight_paths.archive_path)?; + let entry_id = database::entry_id_for_uid(&conn, &preflight_uid)? + .ok_or_else(|| anyhow::anyhow!("entry not found"))?; + if !force { + if let Some(existing) = database::find_entry_summary( + &conn, entry_id, &provider_kind, provider_model.as_deref(), + summarizer::PROMPT_VERSION, &input.input_sha256, + )? { + if existing.status == "completed" { + return Ok(PreflightOutcome::Cached(existing)); + } } } - } - - // 4. Claim the row up front so the 202 can name it and the client can poll - // immediately, before the worker thread has done anything. - let summary_uid = database::upsert_pending_entry_summary( - &conn, - entry_id, - provider.kind(), - provider.model(), - summarizer::PROMPT_VERSION, - &input.input_sha256, - )?; - drop(conn); + let summary_uid = database::upsert_pending_entry_summary( + &conn, entry_id, &provider_kind, provider_model.as_deref(), + summarizer::PROMPT_VERSION, &input.input_sha256, + )?; + Ok(PreflightOutcome::Pending { input, summary_uid }) + }) + .await + .map_err(|e| ApiError::internal(&format!("summary preflight task failed: {e}")))? + .map_err(|e| { + if summarizer::is_unsupported_summary_content_error(&e) { + ApiError::bad_request(summarizer::UNSUPPORTED_SUMMARY_CONTENT_MESSAGE) + } else if format!("{e:#}") == "entry not found" { + ApiError::not_found("entry not found") + } else { + ApiError::bad_request(&format!("{e:#}")) + } + })?; + let (input, summary_uid) = match outcome { + PreflightOutcome::Cached(existing) => return Ok(( + StatusCode::OK, + serde_json::to_value(&existing).map(Json) + .map_err(|e| ApiError::internal(&e.to_string()))?, + )), + PreflightOutcome::Pending { input, summary_uid } => (input, summary_uid), + }; let archive_path = mounted.archive_path.clone(); - let entry_uid_bg = entry_uid.clone(); let summary_uid_bg = summary_uid.clone(); - let summary_options_bg = summary_options; tokio::task::spawn_blocking(move || { - // summarize_entry owns the pending → running → completed/failed - // transitions for this same row (the cache key is identical, so the - // upsert above and the one inside it resolve to one row). We only have - // to catch the case where it fails before it can record anything. + // The input and attempt were claimed during blocking preflight, so no + // row is reset and no archive content is read twice. if let Err(e) = - summarizer::summarize_entry( + summarizer::summarize_prebuilt_entry( &archive_paths, - &entry_uid_bg, - summary_options_bg, + input, + &summary_uid_bg, provider.as_ref(), - summarizer::PROMPT_VERSION, ) { eprintln!("warn: summary {summary_uid_bg}: {e:#}");