mirror of
https://github.com/thegeneralist01/archivr
synced 2026-10-09 21:03:17 +02:00
Merge branch 'feat/llm-summaries' into integration-all-three
# Conflicts: # crates/archivr-server/static/index.html
This commit is contained in:
commit
6a24b66a04
12 changed files with 1958 additions and 52 deletions
|
|
@ -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,
|
||||
}))
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,3 +4,4 @@ pub mod database;
|
|||
pub mod downloader;
|
||||
pub mod hash;
|
||||
pub mod twitter;
|
||||
pub mod summarizer;
|
||||
|
|
|
|||
1104
crates/archivr-core/src/summarizer.rs
Normal file
1104
crates/archivr-core/src/summarizer.rs
Normal file
File diff suppressed because it is too large
Load diff
|
|
@ -29,7 +29,7 @@ use std::{
|
|||
time::{Duration, Instant},
|
||||
};
|
||||
|
||||
use archivr_core::{archive, capture, database, downloader};
|
||||
use archivr_core::{archive, capture, database, downloader, summarizer};
|
||||
use axum::{
|
||||
Json, Router,
|
||||
extract::{ConnectInfo, DefaultBodyLimit, Multipart, Path, Query, Request, State},
|
||||
|
|
@ -268,6 +268,10 @@ pub fn app_with_state(state: AppState) -> Router {
|
|||
"/api/archives/:archive_id/entries/:entry_uid/rearchive",
|
||||
post(rearchive_handler),
|
||||
)
|
||||
.route(
|
||||
"/api/archives/:archive_id/entries/:entry_uid/summary",
|
||||
get(entry_summary_handler).post(request_entry_summary_handler),
|
||||
)
|
||||
.route(
|
||||
"/api/archives/:archive_id/entries/:entry_uid/favicon",
|
||||
get(serve_entry_favicon),
|
||||
|
|
@ -514,6 +518,148 @@ async fn entry_detail(
|
|||
Ok(Json(detail))
|
||||
}
|
||||
|
||||
#[derive(Debug, serde::Deserialize)]
|
||||
struct SummaryRequestBody {
|
||||
provider: String,
|
||||
#[serde(default)]
|
||||
force: bool,
|
||||
}
|
||||
|
||||
/// `GET /api/archives/:id/entries/:uid/summary`
|
||||
///
|
||||
/// Read-only, gated exactly like entry detail: a guest may read a summary only
|
||||
/// for an entry whose content they could already read. Returns
|
||||
/// `{ entry_uid, summary }` with a null summary when none has been requested,
|
||||
/// which is also what the frontend polls while a job is running.
|
||||
async fn entry_summary_handler(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
Path((archive_id, entry_uid)): Path<(String, String)>,
|
||||
) -> Result<Json<serde_json::Value>, ApiError> {
|
||||
let mounted = mounted_archive(&state, &archive_id)?;
|
||||
let conn = database::open_or_initialize(&mounted.archive_path)?;
|
||||
if matches!(auth_user, AuthUser::Guest)
|
||||
&& !database::is_entry_publicly_accessible(&conn, &entry_uid)?
|
||||
{
|
||||
return Err(ApiError::unauthorized("login required"));
|
||||
}
|
||||
let entry_id = database::entry_id_for_uid(&conn, &entry_uid)?
|
||||
.ok_or(ApiError::not_found("entry not found"))?;
|
||||
let summary = database::latest_entry_summary(&conn, entry_id)?;
|
||||
Ok(Json(
|
||||
serde_json::json!({ "entry_uid": entry_uid, "summary": summary }),
|
||||
))
|
||||
}
|
||||
|
||||
/// `POST /api/archives/:id/entries/:uid/summary`
|
||||
///
|
||||
/// Manual-only summary generation. Body: `{ "provider": "...", "force": false }`.
|
||||
///
|
||||
/// Provider configuration and content extraction are both resolved *before*
|
||||
/// spawning, so a missing env var or an unsummarizable artifact comes back as a
|
||||
/// synchronous 400 naming the exact problem rather than as a background job the
|
||||
/// caller has to poll only to learn about a config typo.
|
||||
///
|
||||
/// Returns 200 with the existing row when an identical cache key already
|
||||
/// completed and `force` is false; otherwise 202 with a pending row.
|
||||
async fn request_entry_summary_handler(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
Path((archive_id, entry_uid)): Path<(String, String)>,
|
||||
Json(body): Json<SummaryRequestBody>,
|
||||
) -> Result<(StatusCode, Json<serde_json::Value>), ApiError> {
|
||||
auth_user.require_role(ROLE_USER)?;
|
||||
let mounted = mounted_archive(&state, &archive_id)?;
|
||||
let archive_paths =
|
||||
archive::read_archive_paths(&mounted.archive_path).map_err(ApiError::from)?;
|
||||
|
||||
// 1. Provider config from the environment. The error text carries the exact
|
||||
// variable name, which is the whole point of returning it as a 400.
|
||||
let provider_cfg = summarizer::provider_from_env(&body.provider)
|
||||
.map_err(|e| ApiError::bad_request(&format!("{e:#}")))?;
|
||||
let provider = summarizer::provider_from_config(provider_cfg);
|
||||
|
||||
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)
|
||||
.map_err(|e| 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()))?,
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 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 archive_path = mounted.archive_path.clone();
|
||||
let entry_uid_bg = entry_uid.clone();
|
||||
let summary_uid_bg = summary_uid.clone();
|
||||
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.
|
||||
if let Err(e) =
|
||||
summarizer::summarize_entry(&archive_paths, &entry_uid_bg, provider.as_ref(), summarizer::PROMPT_VERSION)
|
||||
{
|
||||
eprintln!("warn: summary {summary_uid_bg}: {e:#}");
|
||||
if let Ok(conn) = database::open_or_initialize(&archive_path) {
|
||||
if let Ok(Some(row)) = database::get_entry_summary_by_uid(&conn, &summary_uid_bg) {
|
||||
if row.status != "failed" {
|
||||
database::update_entry_summary_status(
|
||||
&conn,
|
||||
&summary_uid_bg,
|
||||
"failed",
|
||||
None,
|
||||
Some(&format!("{e:#}")),
|
||||
)
|
||||
.ok();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
Ok((
|
||||
StatusCode::ACCEPTED,
|
||||
Json(serde_json::json!({
|
||||
"summary_uid": summary_uid,
|
||||
"status": "pending",
|
||||
"entry_uid": entry_uid,
|
||||
})),
|
||||
))
|
||||
}
|
||||
|
||||
async fn list_runs(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
|
|
|
|||
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
|
|
@ -6,8 +6,8 @@
|
|||
<title>Archivr</title>
|
||||
<link rel="icon" type="image/svg+xml" href="/favicon.svg">
|
||||
<link rel="icon" type="image/x-icon" href="/favicon.ico">
|
||||
<script type="module" crossorigin src="/assets/index-CthL4K20.js"></script>
|
||||
<link rel="stylesheet" crossorigin href="/assets/index-DicH9MNh.css">
|
||||
<script type="module" crossorigin src="/assets/index-DdftptO_.js"></script>
|
||||
<link rel="stylesheet" crossorigin href="/assets/index-e69aZgiQ.css">
|
||||
</head>
|
||||
<body>
|
||||
<div id="root"></div>
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue