1
Fork 0
mirror of https://github.com/thegeneralist01/archivr synced 2026-10-09 21:03:17 +02:00

feat(core): add entry_summaries schema + summarizer trait/providers

Per-entry LLM summaries as a regenerable child record, not a column on
archived_entries and not an on-disk artifact: an entry may carry several
summaries (one per provider/model/prompt version), any of which can be
discarded and recomputed. Generation is manual-only — nothing in capture.rs
calls into this module.

- database.rs: entry_summaries table + index, EntrySummaryRecord, and
  upsert/update/find/latest helpers mirroring the capture_jobs style.
  provider_model is stored as '' rather than NULL because SQLite treats
  NULLs as distinct inside a UNIQUE index, which would stop the CLI
  providers (no model) from ever deduping on the cache key.
- summarizer.rs: SummaryProvider trait with four implementations —
  Anthropic Messages API, OpenAI-compatible chat completions, `claude -p`
  and `codex exec -`. Configuration comes from env vars only (never TOML),
  matching how yt-dlp / single-file / tweet-scraper are resolved, which
  also keeps API keys out of anything the archive persists.
- archive.rs: EntryDetail gains latest_summary, populated by one extra
  LIMIT 1 query in get_entry_detail. EntrySummaryView aliases the DB row
  rather than duplicating it.

Implementation notes:
- No tokio in core. CLI timeouts are enforced structurally: stdout is
  drained on its own thread and handed back over a channel so the calling
  thread can recv_timeout and kill an overrunning child; stdin is written
  on a third thread so a 48 KB prompt cannot deadlock against a child
  waiting for us to read.
- HTML is reduced with regex rather than a parser: html5ever is not in the
  tree, and a model tolerates imperfect whitespace. Paired tags are spelled
  out per tag because Rust's regex engine has no backreferences by design.
- reqwest is declared with only the `blocking` feature here, so bodies are
  serialized via .body(value.to_string()) instead of widening the
  workspace dependency for .json().
- input_sha256 holds a SHA3-256 digest via hash::hash_bytes, the tree's one
  hashing primitive; the content is truncated to 48 KB *before* hashing so
  the cache key describes exactly the bytes the model saw.

Tests: no mockito/wiremock in dev-deps, and adding a mock HTTP server for
one JSON shape is a poor trade, so the two halves that can actually break
are tested directly — request-body builders and response parsers — leaving
only reqwest's own transport uncovered. Plus schema idempotency, cache-key
dedupe, cascade-on-delete, provider_from_env happy/missing-var paths, HTML
and tweet extraction, output normalization, and the CLI runner's stdin
round-trip, timeout kill, and nonzero-exit paths.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
TheGeneralist 2026-08-22 15:41:46 +02:00
parent bf5f95397f
commit fd3f6903e2
Signed by: thegeneralist01
SSH key fingerprint: SHA256:pp9qddbCNmVNoSjevdvQvM5z0DHN7LTa8qBMbcMq/R4
4 changed files with 1519 additions and 0 deletions

View file

@ -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<String>,
pub artifacts: Vec<EntryArtifactSummary>,
/// 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<EntrySummaryView>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
@ -343,12 +354,15 @@ pub fn get_entry_detail(
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
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,
}))
}

View file

@ -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<String>,
pub prompt_version: String,
pub input_sha256: String,
pub status: String,
pub summary_text: Option<String>,
pub error_text: Option<String>,
pub created_at: String,
pub updated_at: String,
pub completed_at: Option<String>,
}
#[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<usize> {
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<EntrySummaryRecord> {
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<String>>(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<Option<i64>> {
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<String> {
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<Option<EntrySummaryRecord>> {
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<Option<EntrySummaryRecord>> {
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<Option<EntrySummaryRecord>> {
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);
}
}

View file

@ -4,3 +4,4 @@ pub mod database;
pub mod downloader;
pub mod hash;
pub mod twitter;
pub mod summarizer;

File diff suppressed because it is too large Load diff