diff --git a/crates/archivr-core/src/archive.rs b/crates/archivr-core/src/archive.rs index e84d4a2..ec467ae 100644 --- a/crates/archivr-core/src/archive.rs +++ b/crates/archivr-core/src/archive.rs @@ -36,6 +36,14 @@ pub struct EntrySummary { pub cacheable_bytes: i64, } +/// One stored LLM summary, as exposed over the API. +/// +/// Aliased rather than redefined: the DB row is already the exact shape the +/// frontend needs, and a second near-identical struct would only add a mapping +/// step to keep in sync. The `View` name exists because `EntrySummary` in this +/// module is the *entry listing* row, an unrelated thing. +pub use crate::database::EntrySummaryRecord as EntrySummaryView; + #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] pub struct EntryDetail { pub summary: EntrySummary, @@ -43,6 +51,9 @@ pub struct EntryDetail { pub source_metadata_json: String, pub display_metadata_json: Option, pub artifacts: Vec, + /// Most recently updated summary for this entry, if any has ever been + /// requested. Always `None` on a fresh capture — summarization is manual. + pub latest_summary: Option, } #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] @@ -343,12 +354,15 @@ pub fn get_entry_detail( })? .collect::>>()?; + let latest_summary = database::latest_entry_summary(conn, entry_id)?; + Ok(Some(EntryDetail { summary, structured_root_relpath, source_metadata_json, display_metadata_json, artifacts, + latest_summary, })) } diff --git a/crates/archivr-core/src/database.rs b/crates/archivr-core/src/database.rs index 95b4f26..34b1d2e 100644 --- a/crates/archivr-core/src/database.rs +++ b/crates/archivr-core/src/database.rs @@ -110,6 +110,28 @@ pub struct CaptureJobRecord { pub updated_at: String, } +/// One row of `entry_summaries` — a regenerable LLM summary of an entry. +/// +/// `provider_model` is stored as `''` (not NULL) when a provider has no explicit +/// model, because SQLite treats NULLs as distinct inside a UNIQUE index and a +/// NULL model would defeat the `(entry_id, provider_kind, provider_model, +/// prompt_version, input_sha256)` dedupe key. Readers map `''` back to `None`. +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] +pub struct EntrySummaryRecord { + pub summary_uid: String, + pub entry_uid: String, + pub provider_kind: String, + pub provider_model: Option, + pub prompt_version: String, + pub input_sha256: String, + pub status: String, + pub summary_text: Option, + pub error_text: Option, + pub created_at: String, + pub updated_at: String, + pub completed_at: Option, +} + #[derive(Debug, Clone, serde::Serialize)] pub struct UserSummary { pub user_uid: String, @@ -322,6 +344,24 @@ pub fn initialize_schema(conn: &Connection) -> Result<()> { created_at TEXT NOT NULL, updated_at TEXT NOT NULL ); + CREATE TABLE IF NOT EXISTS entry_summaries ( + 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, + 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, + UNIQUE(entry_id, provider_kind, provider_model, prompt_version, input_sha256) + ); + CREATE INDEX IF NOT EXISTS idx_entry_summaries_entry_updated + ON entry_summaries(entry_id, updated_at DESC); 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); @@ -1354,6 +1394,186 @@ pub fn fail_stalled_capture_jobs(conn: &Connection) -> Result { Ok(n) } +// ── entry_summaries ──────────────────────────────────────────────────────── +// +// Summaries are a regenerable child record of an entry, never a column on +// `archived_entries`: an entry can carry several (one per provider/model/prompt +// version), and any of them can be thrown away and recomputed. Generation is +// manual-only — nothing in `capture.rs` writes here. + +/// `SELECT` list shared by every `entry_summaries` read, so all readers build +/// an identical `EntrySummaryRecord` from the same column ordering. +const ENTRY_SUMMARY_COLS: &str = "SELECT s.summary_uid, e.entry_uid, s.provider_kind, s.provider_model, + s.prompt_version, s.input_sha256, s.status, s.summary_text, + s.error_text, s.created_at, s.updated_at, s.completed_at + FROM entry_summaries s + JOIN archived_entries e ON e.id = s.entry_id"; + +fn map_entry_summary(row: &rusqlite::Row<'_>) -> rusqlite::Result { + Ok(EntrySummaryRecord { + summary_uid: row.get(0)?, + entry_uid: row.get(1)?, + provider_kind: row.get(2)?, + // '' is the stored stand-in for "this provider has no model"; see the + // doc comment on EntrySummaryRecord for why it is not NULL. + provider_model: row + .get::<_, Option>(3)? + .filter(|m| !m.is_empty()), + prompt_version: row.get(4)?, + input_sha256: row.get(5)?, + status: row.get(6)?, + summary_text: row.get(7)?, + error_text: row.get(8)?, + created_at: row.get(9)?, + updated_at: row.get(10)?, + completed_at: row.get(11)?, + }) +} + +/// Resolves `entry_uid` to its integer primary key. `Ok(None)` if absent. +pub fn entry_id_for_uid(conn: &Connection, entry_uid: &str) -> Result> { + conn.query_row( + "SELECT id FROM archived_entries WHERE entry_uid = ?1", + [entry_uid], + |row| row.get(0), + ) + .optional() + .map_err(Into::into) +} + +/// Creates (or resets to `pending`) the summary row for one cache key. +/// +/// The `ON CONFLICT` arm is what makes "Regenerate" work: re-requesting the same +/// (entry, provider, model, prompt, input) reuses the existing row rather than +/// violating the UNIQUE index, clearing any previous text/error so the UI does +/// not show a stale result next to a running job. Returns the row's `summary_uid`. +pub fn upsert_pending_entry_summary( + conn: &Connection, + entry_id: i64, + provider_kind: &str, + provider_model: Option<&str>, + prompt_version: &str, + input_sha256: &str, +) -> Result { + let summary_uid = format!("sum_{}", &Uuid::new_v4().simple().to_string()[..10]); + let now = now_timestamp(); + let model = provider_model.unwrap_or(""); + conn.execute( + "INSERT INTO entry_summaries + (summary_uid, entry_id, provider_kind, provider_model, prompt_version, + input_sha256, status, summary_text, error_text, created_at, updated_at, completed_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, 'pending', NULL, NULL, ?7, ?7, NULL) + ON CONFLICT(entry_id, provider_kind, provider_model, prompt_version, input_sha256) + DO UPDATE SET status = 'pending', summary_text = NULL, error_text = NULL, + completed_at = NULL, updated_at = ?7", + rusqlite::params![ + summary_uid, + entry_id, + provider_kind, + model, + prompt_version, + input_sha256, + now + ], + )?; + // On the conflict path the pre-existing row keeps its original summary_uid, + // so read it back rather than returning the one we just generated. + let stored: String = conn.query_row( + "SELECT summary_uid FROM entry_summaries + WHERE entry_id = ?1 AND provider_kind = ?2 AND provider_model = ?3 + AND prompt_version = ?4 AND input_sha256 = ?5", + rusqlite::params![entry_id, provider_kind, model, prompt_version, input_sha256], + |row| row.get(0), + )?; + Ok(stored) +} + +/// Moves a summary row through `running` → `completed` / `failed`. +/// `completed_at` is stamped only on the terminal `completed` transition. +pub fn update_entry_summary_status( + conn: &Connection, + summary_uid: &str, + status: &str, + summary_text: Option<&str>, + error_text: Option<&str>, +) -> Result<()> { + let now = now_timestamp(); + let completed_at = if status == "completed" { + Some(now.clone()) + } else { + None + }; + conn.execute( + "UPDATE entry_summaries + SET status = ?1, summary_text = ?2, error_text = ?3, + completed_at = ?4, updated_at = ?5 + WHERE summary_uid = ?6", + rusqlite::params![status, summary_text, error_text, completed_at, now, summary_uid], + )?; + Ok(()) +} + +/// Returns one summary by its public uid. +pub fn get_entry_summary_by_uid( + conn: &Connection, + summary_uid: &str, +) -> Result> { + conn.query_row( + &format!("{ENTRY_SUMMARY_COLS} WHERE s.summary_uid = ?1"), + [summary_uid], + map_entry_summary, + ) + .optional() + .map_err(Into::into) +} + +/// Looks up the row for one exact cache key — used to short-circuit a POST when +/// a completed summary for identical input already exists and `force` is false. +pub fn find_entry_summary( + conn: &Connection, + entry_id: i64, + provider_kind: &str, + provider_model: Option<&str>, + prompt_version: &str, + input_sha256: &str, +) -> Result> { + conn.query_row( + &format!( + "{ENTRY_SUMMARY_COLS} + WHERE s.entry_id = ?1 AND s.provider_kind = ?2 AND s.provider_model = ?3 + AND s.prompt_version = ?4 AND s.input_sha256 = ?5" + ), + rusqlite::params![ + entry_id, + provider_kind, + provider_model.unwrap_or(""), + prompt_version, + input_sha256 + ], + map_entry_summary, + ) + .optional() + .map_err(Into::into) +} + +/// Most recently touched summary for an entry, whatever its status. +/// Backs `EntryDetail.latest_summary` and the GET summary route. +pub fn latest_entry_summary( + conn: &Connection, + entry_id: i64, +) -> Result> { + conn.query_row( + &format!( + "{ENTRY_SUMMARY_COLS} WHERE s.entry_id = ?1 + ORDER BY 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, @@ -4176,4 +4396,184 @@ mod tests { assert_eq!(child_count, 1, "only the child item must be counted"); } + + // ── entry_summaries ──────────────────────────────────────────────────── + + #[test] + fn initialize_schema_is_idempotent_for_entry_summaries() { + // initialize_schema runs on *every* open_or_initialize, so re-running it + // against a populated DB must be a no-op, not an error. + let c = conn(); + let entry = create_entry_fixture(&c, "private", None, None); + upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "abc").unwrap(); + + initialize_schema(&c).unwrap(); + initialize_schema(&c).unwrap(); + + let exists: i64 = c + .query_row( + "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='entry_summaries'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(exists, 1); + assert!(latest_entry_summary(&c, entry.id).unwrap().is_some()); + } + + #[test] + fn entry_summary_lifecycle_pending_running_completed() { + let c = conn(); + let entry = create_entry_fixture(&c, "private", None, None); + let uid = + upsert_pending_entry_summary(&c, entry.id, "anthropic_http", Some("m1"), "v1", "sha1") + .unwrap(); + assert!(uid.starts_with("sum_"), "got {uid}"); + assert_eq!(uid.len(), "sum_".len() + 10); + + let rec = get_entry_summary_by_uid(&c, &uid).unwrap().unwrap(); + assert_eq!(rec.status, "pending"); + assert_eq!(rec.entry_uid, entry.entry_uid); + assert_eq!(rec.provider_model.as_deref(), Some("m1")); + assert!(rec.completed_at.is_none()); + + update_entry_summary_status(&c, &uid, "running", None, None).unwrap(); + assert_eq!( + get_entry_summary_by_uid(&c, &uid).unwrap().unwrap().status, + "running" + ); + + update_entry_summary_status(&c, &uid, "completed", Some("{\"summary\":\"s\"}"), None) + .unwrap(); + let rec = get_entry_summary_by_uid(&c, &uid).unwrap().unwrap(); + assert_eq!(rec.status, "completed"); + assert_eq!(rec.summary_text.as_deref(), Some("{\"summary\":\"s\"}")); + assert!(rec.completed_at.is_some(), "completed rows must be stamped"); + } + + #[test] + fn entry_summary_failure_records_error_and_no_completed_at() { + let c = conn(); + let entry = create_entry_fixture(&c, "private", None, None); + let uid = upsert_pending_entry_summary(&c, entry.id, "codex_cli", None, "v1", "s").unwrap(); + update_entry_summary_status(&c, &uid, "failed", None, Some("boom")).unwrap(); + let rec = get_entry_summary_by_uid(&c, &uid).unwrap().unwrap(); + assert_eq!(rec.status, "failed"); + assert_eq!(rec.error_text.as_deref(), Some("boom")); + assert!(rec.completed_at.is_none()); + } + + #[test] + fn upsert_pending_entry_summary_reuses_the_row_for_an_identical_cache_key() { + 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. + 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"); + + let n: i64 = c + .query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0)) + .unwrap(); + assert_eq!(n, 1); + + let rec = get_entry_summary_by_uid(&c, &first).unwrap().unwrap(); + assert_eq!(rec.status, "pending"); + assert!(rec.summary_text.is_none(), "stale text must be cleared"); + } + + #[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. + 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(); + upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "sha").unwrap(); + let n: i64 = c + .query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0)) + .unwrap(); + assert_eq!(n, 1); + // …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); + } + + #[test] + fn different_cache_keys_produce_separate_summaries() { + 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(); + upsert_pending_entry_summary(&c, entry.id, "codex_cli", None, "v1", "sha").unwrap(); + upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v2", "sha").unwrap(); + upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "other").unwrap(); + let n: i64 = c + .query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0)) + .unwrap(); + assert_eq!(n, 4); + } + + #[test] + fn find_entry_summary_matches_only_the_exact_cache_key() { + let c = conn(); + let entry = create_entry_fixture(&c, "private", None, None); + let uid = + upsert_pending_entry_summary(&c, entry.id, "openai_compatible", Some("gpt"), "v1", "s") + .unwrap(); + let hit = find_entry_summary(&c, entry.id, "openai_compatible", Some("gpt"), "v1", "s") + .unwrap() + .unwrap(); + assert_eq!(hit.summary_uid, uid); + assert!( + find_entry_summary(&c, entry.id, "openai_compatible", Some("gpt"), "v1", "changed") + .unwrap() + .is_none(), + "a changed input digest must miss the cache" + ); + } + + #[test] + fn latest_entry_summary_returns_the_most_recently_updated_row() { + let c = conn(); + let entry = create_entry_fixture(&c, "private", None, None); + let a = upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "a").unwrap(); + let b = upsert_pending_entry_summary(&c, entry.id, "codex_cli", None, "v1", "b").unwrap(); + update_entry_summary_status(&c, &a, "completed", Some("first"), None).unwrap(); + update_entry_summary_status(&c, &b, "completed", Some("second"), None).unwrap(); + // Same-timestamp ties break on id DESC, so the later insert wins either way. + assert_eq!(latest_entry_summary(&c, entry.id).unwrap().unwrap().summary_uid, b); + } + + #[test] + fn latest_entry_summary_is_none_for_an_unsummarized_entry() { + let c = conn(); + let entry = create_entry_fixture(&c, "private", None, None); + assert!(latest_entry_summary(&c, entry.id).unwrap().is_none()); + } + + #[test] + fn deleting_an_entry_cascades_to_its_summaries() { + let c = conn(); + c.pragma_update(None, "foreign_keys", "ON").unwrap(); + let entry = create_entry_fixture(&c, "private", None, None); + upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "a").unwrap(); + assert!(delete_entry(&c, &entry.entry_uid).unwrap()); + let n: i64 = c + .query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0)) + .unwrap(); + assert_eq!(n, 0); + } + + #[test] + fn entry_id_for_uid_resolves_and_misses() { + let c = conn(); + let entry = create_entry_fixture(&c, "private", None, None); + assert_eq!(entry_id_for_uid(&c, &entry.entry_uid).unwrap(), Some(entry.id)); + assert_eq!(entry_id_for_uid(&c, "ent_nope").unwrap(), None); + } } diff --git a/crates/archivr-core/src/lib.rs b/crates/archivr-core/src/lib.rs index cd97ffb..907f9e6 100644 --- a/crates/archivr-core/src/lib.rs +++ b/crates/archivr-core/src/lib.rs @@ -4,3 +4,4 @@ pub mod database; pub mod downloader; pub mod hash; pub mod twitter; +pub mod summarizer; diff --git a/crates/archivr-core/src/summarizer.rs b/crates/archivr-core/src/summarizer.rs new file mode 100644 index 0000000..d29ca7f --- /dev/null +++ b/crates/archivr-core/src/summarizer.rs @@ -0,0 +1,1104 @@ +//! Per-entry LLM summaries. +//! +//! A summary is a *regenerable child record* of an entry (`entry_summaries`), +//! not a column on `archived_entries` and not an artifact on disk: an entry can +//! carry several summaries (one per provider / model / prompt version), any of +//! them can be discarded and recomputed, and none of them is part of the +//! preserved capture. Generation is manual-only — nothing in `capture.rs` calls +//! into this module. +//! +//! Four providers sit behind one [`SummaryProvider`] trait: two HTTP APIs and +//! two local CLIs. They are configured by environment variable, never by TOML, +//! matching how the rest of the tree resolves external tools (`ARCHIVR_YT_DLP`, +//! `ARCHIVR_SINGLE_FILE`, `ARCHIVR_TWEET_SCRAPER`, …) and keeping API keys out +//! of any file the archive would otherwise persist. + +use anyhow::{Context, Result, anyhow, bail}; +use std::{ + env, + io::Write, + path::{Path, PathBuf}, + process::{Command, Stdio}, + sync::mpsc, + thread, + time::Duration, +}; + +use crate::{archive::ArchivePaths, database, hash}; + +/// Bump whenever the prompt text below changes in a way that would produce a +/// materially different summary. It is part of the `entry_summaries` cache key, +/// so a bump makes every stored summary regenerate on next request instead of +/// silently mixing outputs from two different prompts. +pub const PROMPT_VERSION: &str = "v1-2026-08-22"; + +/// Upper bound on characters fed to a model. Archived pages run to hundreds of +/// kilobytes; past this point we are paying for tokens that do not change a +/// five-sentence summary. Truncation happens *before* hashing so the cache key +/// describes exactly what the model saw. +const MAX_INPUT_CHARS: usize = 48_000; + +const DEFAULT_HTTP_TIMEOUT_SECS: u64 = 120; +const DEFAULT_CLI_TIMEOUT_SECS: u64 = 300; + +/// The instruction half of the prompt. JSON output is requested because parsing +/// prose out of a free-form answer is the single most fragile part of an LLM +/// integration; a JSON object survives models that like to add pleasantries. +const SYSTEM_PROMPT: &str = "\ +You summarize archived web content for a personal archive index. + +Reply with a single JSON object and nothing else — no markdown fence, no prose +before or after. The object has exactly these keys: + + \"tldr\": one sentence, at most 25 words. + \"summary\": 4 to 6 sentences of plain English describing what the content + says, its claims, and its conclusion. No preamble like + \"This article discusses\". + \"tags\": an array of at most 5 short lowercase topic tags. + +Write in English regardless of the source language. If the content is too +short or empty to summarize, still return the object and say so in \"summary\"."; + +// ── Request / output types ───────────────────────────────────────────────── + +/// Everything the prompt builder needs about one entry. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct SummaryRequest { + pub entry_uid: String, + pub title: Option, + pub source_kind: String, + pub entity_kind: String, + pub content: String, +} + +/// What a provider produced. `model` is echoed back because HTTP providers may +/// resolve an alias (`claude-3-5-sonnet-latest`) to a dated concrete model, and +/// the concrete one is what we want recorded against the summary. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct SummaryOutput { + pub text: String, + pub model: Option, +} + +/// One way of turning a [`SummaryRequest`] into text. +/// +/// `Send + Sync` so a boxed provider can cross into the server's +/// `spawn_blocking` worker. +pub trait SummaryProvider: Send + Sync { + /// Stable identifier persisted as `entry_summaries.provider_kind`. + fn kind(&self) -> &'static str; + fn model(&self) -> Option<&str>; + fn summarize(&self, request: &SummaryRequest) -> Result; +} + +// ── Configuration ────────────────────────────────────────────────────────── + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct HttpProviderConfig { + /// Full URL, e.g. `https://api.anthropic.com/v1/messages`. + pub endpoint: String, + pub api_key: String, + pub model: String, + pub timeout_secs: u64, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct CliProviderConfig { + /// Resolved via `ARCHIVR_CLAUDE_CLI` / `ARCHIVR_CODEX_CLI`; a bare name is + /// left for the OS to resolve on `PATH`, as elsewhere in the tree. + pub executable: PathBuf, + pub model: Option, + pub timeout_secs: u64, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum ProviderConfig { + AnthropicHttp(HttpProviderConfig), + OpenAiCompatible(HttpProviderConfig), + ClaudeCli(CliProviderConfig), + CodexCli(CliProviderConfig), +} + +pub const PROVIDER_KINDS: [&str; 4] = [ + "anthropic_http", + "openai_compatible", + "claude_cli", + "codex_cli", +]; + +pub fn provider_from_config(cfg: ProviderConfig) -> Box { + match cfg { + ProviderConfig::AnthropicHttp(c) => Box::new(AnthropicHttpProvider(c)), + ProviderConfig::OpenAiCompatible(c) => Box::new(OpenAiCompatibleProvider(c)), + ProviderConfig::ClaudeCli(c) => Box::new(ClaudeCliProvider(c)), + ProviderConfig::CodexCli(c) => Box::new(CodexCliProvider(c)), + } +} + +/// Reads a required env var, failing with the *exact variable name* so the +/// server can hand a caller an actionable 400 rather than "not configured". +fn required_env(name: &str) -> Result { + match env::var(name) { + Ok(v) if !v.trim().is_empty() => Ok(v), + _ => bail!("missing required environment variable: {name}"), + } +} + +fn env_or(name: &str, default: &str) -> String { + env::var(name) + .ok() + .filter(|v| !v.trim().is_empty()) + .unwrap_or_else(|| default.to_string()) +} + +fn optional_env(name: &str) -> Option { + env::var(name).ok().filter(|v| !v.trim().is_empty()) +} + +fn env_timeout(name: &str, default: u64) -> u64 { + env::var(name) + .ok() + .and_then(|v| v.trim().parse::().ok()) + .filter(|v| *v > 0) + .unwrap_or(default) +} + +/// Builds a provider configuration for `kind` purely from the environment. +pub fn provider_from_env(kind: &str) -> Result { + match kind { + "anthropic_http" => Ok(ProviderConfig::AnthropicHttp(HttpProviderConfig { + endpoint: env_or("ARCHIVR_ANTHROPIC_URL", "https://api.anthropic.com/v1/messages"), + api_key: required_env("ARCHIVR_ANTHROPIC_API_KEY")?, + model: env_or("ARCHIVR_ANTHROPIC_MODEL", "claude-3-5-sonnet-latest"), + timeout_secs: env_timeout("ARCHIVR_SUMMARY_HTTP_TIMEOUT", DEFAULT_HTTP_TIMEOUT_SECS), + })), + "openai_compatible" => Ok(ProviderConfig::OpenAiCompatible(HttpProviderConfig { + endpoint: env_or("ARCHIVR_OPENAI_URL", "https://api.openai.com/v1/chat/completions"), + api_key: required_env("ARCHIVR_OPENAI_API_KEY")?, + model: env_or("ARCHIVR_OPENAI_MODEL", "gpt-4o-mini"), + timeout_secs: env_timeout("ARCHIVR_SUMMARY_HTTP_TIMEOUT", DEFAULT_HTTP_TIMEOUT_SECS), + })), + "claude_cli" => Ok(ProviderConfig::ClaudeCli(CliProviderConfig { + executable: PathBuf::from(env_or("ARCHIVR_CLAUDE_CLI", "claude")), + model: optional_env("ARCHIVR_CLAUDE_MODEL"), + timeout_secs: env_timeout("ARCHIVR_SUMMARY_CLI_TIMEOUT", DEFAULT_CLI_TIMEOUT_SECS), + })), + "codex_cli" => Ok(ProviderConfig::CodexCli(CliProviderConfig { + executable: PathBuf::from(env_or("ARCHIVR_CODEX_CLI", "codex")), + model: optional_env("ARCHIVR_CODEX_MODEL"), + timeout_secs: env_timeout("ARCHIVR_SUMMARY_CLI_TIMEOUT", DEFAULT_CLI_TIMEOUT_SECS), + })), + other => bail!( + "unknown summary provider: {other} (expected one of {})", + PROVIDER_KINDS.join(", ") + ), + } +} + +// ── Prompt assembly ──────────────────────────────────────────────────────── + +/// The user half of the prompt: entry metadata as a small header, then content. +pub fn build_user_prompt(request: &SummaryRequest) -> String { + let mut s = String::new(); + if let Some(title) = request.title.as_deref().filter(|t| !t.trim().is_empty()) { + s.push_str(&format!("Title: {title}\n")); + } + s.push_str(&format!( + "Source: {} / {}\n\nContent:\n{}\n", + request.source_kind, request.entity_kind, request.content + )); + s +} + +/// CLIs take a single prompt string on stdin, so the system half is prepended +/// rather than passed as a separate role. +fn build_combined_prompt(request: &SummaryRequest) -> String { + format!("{SYSTEM_PROMPT}\n\n---\n\n{}", build_user_prompt(request)) +} + +// ── HTTP providers ───────────────────────────────────────────────────────── + +fn http_client(timeout_secs: u64) -> Result { + reqwest::blocking::Client::builder() + // reqwest's own timeout covers connect + read, which is all a + // request/response provider needs — no watchdog thread required. + .timeout(Duration::from_secs(timeout_secs)) + .build() + .context("failed to build HTTP client for summary provider") +} + +/// Body builder kept separate from the transport so it can be unit-tested +/// without a network round-trip. +pub fn anthropic_request_body(model: &str, request: &SummaryRequest) -> serde_json::Value { + serde_json::json!({ + "model": model, + "max_tokens": 1024, + "messages": [{ + "role": "user", + "content": build_combined_prompt(request), + }], + }) +} + +pub fn openai_request_body(model: &str, request: &SummaryRequest) -> serde_json::Value { + serde_json::json!({ + "model": model, + "messages": [ + { "role": "system", "content": SYSTEM_PROMPT }, + { "role": "user", "content": build_user_prompt(request) }, + ], + }) +} + +struct AnthropicHttpProvider(HttpProviderConfig); + +impl SummaryProvider for AnthropicHttpProvider { + fn kind(&self) -> &'static str { + "anthropic_http" + } + fn model(&self) -> Option<&str> { + Some(&self.0.model) + } + fn summarize(&self, request: &SummaryRequest) -> Result { + let body = anthropic_request_body(&self.0.model, request); + let resp = http_client(self.0.timeout_secs)? + .post(&self.0.endpoint) + .header("x-api-key", &self.0.api_key) + .header("anthropic-version", "2023-06-01") + .header("content-type", "application/json") + .body(body.to_string()) + .send() + .with_context(|| format!("request to {} failed", self.0.endpoint))?; + let status = resp.status(); + let text = resp.text().unwrap_or_default(); + if !status.is_success() { + bail!("anthropic API returned {status}: {}", truncate_for_error(&text)); + } + parse_anthropic_response(&text) + } +} + +pub fn parse_anthropic_response(body: &str) -> Result { + let json: serde_json::Value = + serde_json::from_str(body).context("anthropic response was not JSON")?; + let text = json["content"][0]["text"] + .as_str() + .ok_or_else(|| anyhow!("anthropic response had no content[0].text"))?; + Ok(SummaryOutput { + text: text.to_string(), + model: json["model"].as_str().map(str::to_string), + }) +} + +struct OpenAiCompatibleProvider(HttpProviderConfig); + +impl SummaryProvider for OpenAiCompatibleProvider { + fn kind(&self) -> &'static str { + "openai_compatible" + } + fn model(&self) -> Option<&str> { + Some(&self.0.model) + } + fn summarize(&self, request: &SummaryRequest) -> Result { + let body = openai_request_body(&self.0.model, request); + let resp = http_client(self.0.timeout_secs)? + .post(&self.0.endpoint) + .header("authorization", format!("Bearer {}", self.0.api_key)) + .header("content-type", "application/json") + .body(body.to_string()) + .send() + .with_context(|| format!("request to {} failed", self.0.endpoint))?; + let status = resp.status(); + let text = resp.text().unwrap_or_default(); + if !status.is_success() { + bail!( + "openai-compatible API returned {status}: {}", + truncate_for_error(&text) + ); + } + parse_openai_response(&text) + } +} + +pub fn parse_openai_response(body: &str) -> Result { + let json: serde_json::Value = + serde_json::from_str(body).context("openai-compatible response was not JSON")?; + let text = json["choices"][0]["message"]["content"] + .as_str() + .ok_or_else(|| anyhow!("response had no choices[0].message.content"))?; + Ok(SummaryOutput { + text: text.to_string(), + model: json["model"].as_str().map(str::to_string), + }) +} + +fn truncate_for_error(s: &str) -> String { + let trimmed = s.trim(); + if trimmed.chars().count() <= 400 { + return trimmed.to_string(); + } + trimmed.chars().take(400).collect::() + "…" +} + +// ── CLI providers ────────────────────────────────────────────────────────── + +/// Runs `executable args…`, writes `prompt` to its stdin, and returns stdout. +/// +/// `archivr-core` deliberately has no async runtime and the tree carries no +/// `wait_timeout` dependency, so the timeout is enforced by structure rather +/// than by a library: stdout is drained on its own thread and handed back over +/// a channel, which leaves the calling thread free to `recv_timeout` and kill +/// the child if it overruns. stdin is written on a third thread because a +/// 48 KB prompt can exceed the pipe buffer, and writing it inline would +/// deadlock against a child that is waiting for us to read its output. +fn run_cli( + executable: &Path, + args: &[&str], + prompt: &str, + timeout_secs: u64, +) -> Result { + let mut child = Command::new(executable) + .args(args) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .with_context(|| format!("failed to spawn {}", executable.display()))?; + + let mut stdin = child + .stdin + .take() + .ok_or_else(|| anyhow!("failed to open stdin for {}", executable.display()))?; + let prompt_owned = prompt.to_string(); + thread::spawn(move || { + let _ = stdin.write_all(prompt_owned.as_bytes()); + // Dropping stdin closes the pipe, which is what tells the CLI the + // prompt is complete. + }); + + let stdout = child + .stdout + .take() + .ok_or_else(|| anyhow!("failed to open stdout for {}", executable.display()))?; + let (tx, rx) = mpsc::channel(); + thread::spawn(move || { + let mut buf = String::new(); + use std::io::Read; + let mut stdout = stdout; + let res = stdout.read_to_string(&mut buf).map(|_| buf); + let _ = tx.send(res); + }); + + let collected = match rx.recv_timeout(Duration::from_secs(timeout_secs)) { + Ok(res) => res.with_context(|| format!("failed to read stdout of {}", executable.display()))?, + Err(_) => { + let _ = child.kill(); + let _ = child.wait(); + bail!( + "{} timed out after {timeout_secs}s", + executable.display() + ); + } + }; + + let status = child + .wait() + .with_context(|| format!("failed to wait for {}", executable.display()))?; + if !status.success() { + let mut stderr = String::new(); + if let Some(mut e) = child.stderr.take() { + use std::io::Read; + let _ = e.read_to_string(&mut stderr); + } + bail!( + "{} exited with {status}: {}", + executable.display(), + truncate_for_error(&stderr) + ); + } + Ok(collected) +} + +struct ClaudeCliProvider(CliProviderConfig); + +impl SummaryProvider for ClaudeCliProvider { + fn kind(&self) -> &'static str { + "claude_cli" + } + fn model(&self) -> Option<&str> { + self.0.model.as_deref() + } + fn summarize(&self, request: &SummaryRequest) -> Result { + // `claude -p --output-format text` is the documented one-shot + // ("print") mode of the Claude Code CLI: it reads the prompt from + // stdin, writes the answer to stdout, and exits. + let mut args: Vec<&str> = vec!["-p", "--output-format", "text"]; + if let Some(model) = self.0.model.as_deref() { + args.push("--model"); + args.push(model); + } + let out = run_cli(&self.0.executable, &args, &build_combined_prompt(request), self.0.timeout_secs)?; + Ok(SummaryOutput { + text: out, + model: self.0.model.clone(), + }) + } +} + +/// Codex invocation lives in its own module because its one-shot interface is +/// the least stable of the four. +/// +/// Primary form is `codex exec -`, which reads the prompt from stdin. Older +/// builds only accept the prompt as a positional argument, so a failure to +/// spawn/parse falls back to `codex exec `. +/// +/// TESTED: neither form was exercised end-to-end — `codex` is not installed on +/// the machine this was written on (`which codex` → not found). The `claude` +/// CLI path *was* smoke-tested. Treat the codex path as best-effort until +/// someone with the binary confirms it. +mod codex { + use super::*; + + pub fn run(cfg: &CliProviderConfig, prompt: &str) -> Result { + let mut args: Vec = vec!["exec".into(), "-".into()]; + if let Some(model) = cfg.model.as_deref() { + args.push("--model".into()); + args.push(model.into()); + } + let arg_refs: Vec<&str> = args.iter().map(String::as_str).collect(); + match run_cli(&cfg.executable, &arg_refs, prompt, cfg.timeout_secs) { + Ok(out) => Ok(out), + Err(primary) => { + // Fallback: prompt as a positional argument, no stdin. + let mut fb: Vec = vec!["exec".into()]; + if let Some(model) = cfg.model.as_deref() { + fb.push("--model".into()); + fb.push(model.into()); + } + fb.push(prompt.into()); + let out = Command::new(&cfg.executable) + .args(&fb) + .output() + .with_context(|| format!("codex `exec -` failed ({primary:#}); positional fallback also failed to spawn"))?; + if !out.status.success() { + bail!( + "codex `exec -` failed ({primary:#}); positional fallback exited with {}: {}", + out.status, + truncate_for_error(&String::from_utf8_lossy(&out.stderr)) + ); + } + Ok(String::from_utf8_lossy(&out.stdout).to_string()) + } + } + } +} + +struct CodexCliProvider(CliProviderConfig); + +impl SummaryProvider for CodexCliProvider { + fn kind(&self) -> &'static str { + "codex_cli" + } + fn model(&self) -> Option<&str> { + self.0.model.as_deref() + } + fn summarize(&self, request: &SummaryRequest) -> Result { + let out = codex::run(&self.0, &build_combined_prompt(request))?; + Ok(SummaryOutput { + text: out, + model: self.0.model.clone(), + }) + } +} + +// ── Content extraction ───────────────────────────────────────────────────── + +/// Strips markup from an archived HTML page. +/// +/// Deliberately regex-based rather than a real parser: `html5ever` is not in +/// the dependency tree, and pulling a full HTML parser in to feed a language +/// model — which tolerates imperfect whitespace and stray angle brackets +/// fine — is not worth the build cost. `\ +

Hello world

"; + let text = strip_html(html); + assert!(text.contains("Hello world")); + assert!(!text.contains("color:red")); + assert!(!text.contains("var x")); + } + + #[test] + fn strip_html_keeps_paragraphs_apart() { + let text = strip_html("

One

Two

"); + // Without block-level break handling these would run together as + // "OneTwo", which reads as a single garbled sentence to the model. + assert!(text.contains("One")); + assert!(text.contains("Two")); + assert!(!text.contains("OneTwo")); + } + + #[test] + fn strip_html_decodes_common_entities() { + assert_eq!(strip_html("

a & b  c

").replace('\u{a0}', " "), "a & b c"); + } + + #[test] + fn extract_tweet_text_handles_flat_and_wrapped_shapes() { + let flat = serde_json::json!({ "full_text": "tweet body" }); + assert_eq!(extract_tweet_text(&flat).unwrap(), "tweet body"); + + let wrapped = serde_json::json!({ "tweet": { "text": "wrapped body" } }); + assert_eq!(extract_tweet_text(&wrapped).unwrap(), "wrapped body"); + + let threaded = serde_json::json!({ + "full_text": "first", + "thread": [{ "full_text": "second" }], + }); + assert_eq!(extract_tweet_text(&threaded).unwrap(), "first\n\nsecond"); + + assert!(extract_tweet_text(&serde_json::json!({ "id": 1 })).is_none()); + } + + // ── Output normalization ─────────────────────────────────────────────── + + #[test] + fn normalize_summary_json_passes_through_clean_json() { + let raw = r#"{"tldr":"t","summary":"s","tags":["a"]}"#; + let v: serde_json::Value = serde_json::from_str(&normalize_summary_json(raw)).unwrap(); + assert_eq!(v["tldr"], "t"); + assert_eq!(v["tags"][0], "a"); + } + + #[test] + fn normalize_summary_json_strips_markdown_fences() { + let raw = "```json\n{\"tldr\":\"t\",\"summary\":\"s\",\"tags\":[]}\n```"; + let v: serde_json::Value = serde_json::from_str(&normalize_summary_json(raw)).unwrap(); + assert_eq!(v["summary"], "s"); + } + + #[test] + fn normalize_summary_json_recovers_json_wrapped_in_prose() { + let raw = "Sure! Here you go:\n{\"tldr\":\"t\",\"summary\":\"s\",\"tags\":[]}\nHope that helps."; + let v: serde_json::Value = serde_json::from_str(&normalize_summary_json(raw)).unwrap(); + assert_eq!(v["tldr"], "t"); + } + + #[test] + fn normalize_summary_json_wraps_unparseable_output_rather_than_losing_it() { + // A model that ignored the format instruction still produced something + // a human can read; discarding it would be worse than a missing tldr. + let v: serde_json::Value = + serde_json::from_str(&normalize_summary_json("just prose")).unwrap(); + assert_eq!(v["summary"], "just prose"); + assert_eq!(v["tldr"], ""); + } + + // ── CLI runner ───────────────────────────────────────────────────────── + + #[test] + fn run_cli_round_trips_stdin_to_stdout() { + // `cat` stands in for a provider CLI: it proves the prompt reaches the + // child's stdin and the child's stdout comes back intact. + let out = run_cli(Path::new("cat"), &[], "prompt text", 30).unwrap(); + assert_eq!(out, "prompt text"); + } + + #[test] + fn run_cli_kills_a_child_that_overruns_its_timeout() { + let err = run_cli(Path::new("sleep"), &["30"], "", 1).unwrap_err().to_string(); + assert!(err.contains("timed out"), "got: {err}"); + } + + #[test] + fn run_cli_reports_a_nonzero_exit() { + let err = run_cli(Path::new("false"), &[], "", 30).unwrap_err().to_string(); + assert!(err.contains("exited with"), "got: {err}"); + } +}