mirror of
https://github.com/thegeneralist01/archivr
synced 2026-10-09 12:55:00 +02:00
Merge branch 'sol-summary-lifecycle' into integration-all-three
This commit is contained in:
commit
53fab8c76c
4 changed files with 335 additions and 95 deletions
|
|
@ -349,7 +349,7 @@ pub fn initialize_schema(conn: &Connection) -> Result<()> {
|
||||||
summary_uid TEXT NOT NULL UNIQUE,
|
summary_uid TEXT NOT NULL UNIQUE,
|
||||||
entry_id INTEGER NOT NULL REFERENCES archived_entries(id) ON DELETE CASCADE,
|
entry_id INTEGER NOT NULL REFERENCES archived_entries(id) ON DELETE CASCADE,
|
||||||
provider_kind TEXT NOT NULL,
|
provider_kind TEXT NOT NULL,
|
||||||
provider_model TEXT,
|
provider_model TEXT NOT NULL DEFAULT '',
|
||||||
prompt_version TEXT NOT NULL,
|
prompt_version TEXT NOT NULL,
|
||||||
input_sha256 TEXT NOT NULL,
|
input_sha256 TEXT NOT NULL,
|
||||||
status TEXT NOT NULL CHECK(status IN ('pending','running','completed','failed')),
|
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,
|
error_text TEXT,
|
||||||
created_at TEXT NOT NULL,
|
created_at TEXT NOT NULL,
|
||||||
updated_at TEXT NOT NULL,
|
updated_at TEXT NOT NULL,
|
||||||
completed_at TEXT,
|
completed_at TEXT
|
||||||
UNIQUE(entry_id, provider_kind, provider_model, prompt_version, input_sha256)
|
|
||||||
);
|
);
|
||||||
CREATE INDEX IF NOT EXISTS idx_entry_summaries_entry_updated
|
CREATE INDEX IF NOT EXISTS idx_entry_summaries_entry_updated
|
||||||
ON entry_summaries(entry_id, updated_at DESC);
|
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_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_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);
|
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<String> = 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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -1394,6 +1441,21 @@ pub fn fail_stalled_capture_jobs(conn: &Connection) -> Result<usize> {
|
||||||
Ok(n)
|
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<usize> {
|
||||||
|
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 ────────────────────────────────────────────────────────
|
// ── entry_summaries ────────────────────────────────────────────────────────
|
||||||
//
|
//
|
||||||
// Summaries are a regenerable child record of an entry, never a column on
|
// 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<Option<i64
|
||||||
.map_err(Into::into)
|
.map_err(Into::into)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Creates (or resets to `pending`) the summary row for one cache key.
|
/// Creates a fresh pending summary attempt for one cache key.
|
||||||
///
|
///
|
||||||
/// The `ON CONFLICT` arm is what makes "Regenerate" work: re-requesting the same
|
/// Attempts are intentionally not unique by cache key: a forced regeneration
|
||||||
/// (entry, provider, model, prompt, input) reuses the existing row rather than
|
/// must leave an older completed result available while the new attempt runs.
|
||||||
/// 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(
|
pub fn upsert_pending_entry_summary(
|
||||||
conn: &Connection,
|
conn: &Connection,
|
||||||
entry_id: i64,
|
entry_id: i64,
|
||||||
|
|
@ -1462,10 +1522,7 @@ pub fn upsert_pending_entry_summary(
|
||||||
"INSERT INTO entry_summaries
|
"INSERT INTO entry_summaries
|
||||||
(summary_uid, entry_id, provider_kind, provider_model, prompt_version,
|
(summary_uid, entry_id, provider_kind, provider_model, prompt_version,
|
||||||
input_sha256, status, summary_text, error_text, created_at, updated_at, completed_at)
|
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)
|
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![
|
rusqlite::params![
|
||||||
summary_uid,
|
summary_uid,
|
||||||
entry_id,
|
entry_id,
|
||||||
|
|
@ -1476,16 +1533,7 @@ pub fn upsert_pending_entry_summary(
|
||||||
now
|
now
|
||||||
],
|
],
|
||||||
)?;
|
)?;
|
||||||
// On the conflict path the pre-existing row keeps its original summary_uid,
|
Ok(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`.
|
/// Moves a summary row through `running` → `completed` / `failed`.
|
||||||
|
|
@ -1541,7 +1589,9 @@ pub fn find_entry_summary(
|
||||||
&format!(
|
&format!(
|
||||||
"{ENTRY_SUMMARY_COLS}
|
"{ENTRY_SUMMARY_COLS}
|
||||||
WHERE s.entry_id = ?1 AND s.provider_kind = ?2 AND s.provider_model = ?3
|
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"
|
AND s.prompt_version = ?4 AND s.input_sha256 = ?5
|
||||||
|
ORDER BY CASE WHEN s.status = 'completed' THEN 0 ELSE 1 END,
|
||||||
|
s.completed_at DESC, s.updated_at DESC, s.id DESC"
|
||||||
),
|
),
|
||||||
rusqlite::params![
|
rusqlite::params![
|
||||||
entry_id,
|
entry_id,
|
||||||
|
|
@ -1574,6 +1624,24 @@ pub fn latest_entry_summary(
|
||||||
.map_err(Into::into)
|
.map_err(Into::into)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The most recent completed summary, excluding in-flight and failed attempts.
|
||||||
|
/// This is the stable summary shown in entry detail and used by free-text search.
|
||||||
|
pub fn latest_completed_entry_summary(
|
||||||
|
conn: &Connection,
|
||||||
|
entry_id: i64,
|
||||||
|
) -> Result<Option<EntrySummaryRecord>> {
|
||||||
|
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(
|
pub fn create_archive_run(
|
||||||
conn: &Connection,
|
conn: &Connection,
|
||||||
created_by_user_id: i64,
|
created_by_user_id: i64,
|
||||||
|
|
@ -4464,33 +4532,72 @@ mod tests {
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[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 c = conn();
|
||||||
let entry = create_entry_fixture(&c, "private", None, None);
|
let entry = create_entry_fixture(&c, "private", None, None);
|
||||||
let first =
|
let first =
|
||||||
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();
|
||||||
update_entry_summary_status(&c, &first, "completed", Some("old"), None).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 =
|
let second =
|
||||||
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();
|
||||||
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
|
let n: i64 = c
|
||||||
.query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0))
|
.query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0))
|
||||||
.unwrap();
|
.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_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]
|
#[test]
|
||||||
fn entry_summaries_with_no_model_still_dedupe() {
|
fn fail_stalled_entry_summaries_marks_pending_and_running_with_restart_message() {
|
||||||
// Regression guard: a NULL provider_model would compare as distinct in
|
let c = conn();
|
||||||
// SQLite's UNIQUE index, so the CLI providers (which have no model)
|
let entry = create_entry_fixture(&c, "private", None, None);
|
||||||
// would accumulate a new row on every regenerate. We store '' instead.
|
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 c = conn();
|
||||||
let entry = create_entry_fixture(&c, "private", None, None);
|
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();
|
||||||
|
|
@ -4498,7 +4605,7 @@ mod tests {
|
||||||
let n: i64 = c
|
let n: i64 = c
|
||||||
.query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0))
|
.query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0))
|
||||||
.unwrap();
|
.unwrap();
|
||||||
assert_eq!(n, 1);
|
assert_eq!(n, 2);
|
||||||
// …and it reads back as None, not as an empty-string model.
|
// …and it reads back as None, not as an empty-string model.
|
||||||
let rec = latest_entry_summary(&c, entry.id).unwrap().unwrap();
|
let rec = latest_entry_summary(&c, entry.id).unwrap().unwrap();
|
||||||
assert_eq!(rec.provider_model, None);
|
assert_eq!(rec.provider_model, None);
|
||||||
|
|
|
||||||
|
|
@ -1184,14 +1184,34 @@ pub fn summarize_entry(
|
||||||
prompt_version,
|
prompt_version,
|
||||||
&input.input_sha256,
|
&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<database::EntrySummaryRecord> {
|
||||||
|
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) {
|
match provider.summarize(&input.request) {
|
||||||
Ok(output) => {
|
Ok(output) => {
|
||||||
let text = normalize_summary_json(&output.text);
|
let text = normalize_summary_json(&output.text);
|
||||||
database::update_entry_summary_status(
|
database::update_entry_summary_status(
|
||||||
&conn,
|
&conn,
|
||||||
&summary_uid,
|
summary_uid,
|
||||||
"completed",
|
"completed",
|
||||||
Some(&text),
|
Some(&text),
|
||||||
None,
|
None,
|
||||||
|
|
@ -1201,7 +1221,7 @@ pub fn summarize_entry(
|
||||||
let msg = format!("{e:#}");
|
let msg = format!("{e:#}");
|
||||||
database::update_entry_summary_status(
|
database::update_entry_summary_status(
|
||||||
&conn,
|
&conn,
|
||||||
&summary_uid,
|
summary_uid,
|
||||||
"failed",
|
"failed",
|
||||||
None,
|
None,
|
||||||
Some(&msg),
|
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"))
|
.ok_or_else(|| anyhow!("summary row disappeared after write"))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
),
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -513,8 +513,13 @@ async fn entry_detail(
|
||||||
return Err(ApiError::unauthorized("login required"));
|
return Err(ApiError::unauthorized("login required"));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
let detail = archive::get_entry_detail(&conn, &entry_uid)?
|
let mut detail = archive::get_entry_detail(&conn, &entry_uid)?
|
||||||
.ok_or(ApiError::not_found("entry not found"))?;
|
.ok_or(ApiError::not_found("entry not found"))?;
|
||||||
|
if matches!(auth_user, AuthUser::Guest) {
|
||||||
|
let entry_id = database::entry_id_for_uid(&conn, &entry_uid)?
|
||||||
|
.ok_or(ApiError::not_found("entry not found"))?;
|
||||||
|
detail.latest_summary = database::latest_completed_entry_summary(&conn, entry_id)?;
|
||||||
|
}
|
||||||
Ok(Json(detail))
|
Ok(Json(detail))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -547,7 +552,11 @@ async fn entry_summary_handler(
|
||||||
}
|
}
|
||||||
let entry_id = database::entry_id_for_uid(&conn, &entry_uid)?
|
let entry_id = database::entry_id_for_uid(&conn, &entry_uid)?
|
||||||
.ok_or(ApiError::not_found("entry not found"))?;
|
.ok_or(ApiError::not_found("entry not found"))?;
|
||||||
let summary = database::latest_entry_summary(&conn, entry_id)?;
|
let summary = if matches!(auth_user, AuthUser::Guest) {
|
||||||
|
database::latest_completed_entry_summary(&conn, entry_id)?
|
||||||
|
} else {
|
||||||
|
database::latest_entry_summary(&conn, entry_id)?
|
||||||
|
};
|
||||||
Ok(Json(
|
Ok(Json(
|
||||||
serde_json::json!({ "entry_uid": entry_uid, "summary": summary }),
|
serde_json::json!({ "entry_uid": entry_uid, "summary": summary }),
|
||||||
))
|
))
|
||||||
|
|
@ -590,70 +599,72 @@ async fn request_entry_summary_handler(
|
||||||
include_images: body.include_images,
|
include_images: body.include_images,
|
||||||
};
|
};
|
||||||
|
|
||||||
let conn = database::open_or_initialize(&mounted.archive_path)?;
|
// 2. Preflight extraction and SQLite cache/attempt work are synchronous
|
||||||
let entry_id = database::entry_id_for_uid(&conn, &entry_uid)?
|
// core operations, so keep them off the Axum runtime. This also means the
|
||||||
.ok_or(ApiError::not_found("entry not found"))?;
|
// input claimed here is passed directly to the provider worker below.
|
||||||
|
let preflight_paths = archive_paths.clone();
|
||||||
// 2. Extract the same content the summarizer will feed the model, so the
|
let preflight_uid = entry_uid.clone();
|
||||||
// digest below is the identical cache key summarize_entry will compute.
|
let provider_kind = provider.kind().to_string();
|
||||||
let input = summarizer::build_summary_input(&archive_paths, &entry_uid, summary_options)
|
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<PreflightOutcome> {
|
||||||
|
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"))?;
|
||||||
|
// Preserve the route's historical 404 before attempting content
|
||||||
|
// extraction, whose own missing-entry error includes the uid.
|
||||||
|
let input = summarizer::build_summary_input(&preflight_paths, &preflight_uid, summary_options)?;
|
||||||
|
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));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
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| {
|
.map_err(|e| {
|
||||||
if summarizer::is_unsupported_summary_content_error(&e) {
|
if summarizer::is_unsupported_summary_content_error(&e) {
|
||||||
ApiError::bad_request(summarizer::UNSUPPORTED_SUMMARY_CONTENT_MESSAGE)
|
ApiError::bad_request(summarizer::UNSUPPORTED_SUMMARY_CONTENT_MESSAGE)
|
||||||
|
} else if format!("{e:#}") == "entry not found" {
|
||||||
|
ApiError::not_found("entry not found")
|
||||||
} else {
|
} else {
|
||||||
ApiError::bad_request(&format!("{e:#}"))
|
ApiError::bad_request(&format!("{e:#}"))
|
||||||
}
|
}
|
||||||
})?;
|
})?;
|
||||||
|
let (input, summary_uid) = match outcome {
|
||||||
// 3. Cache hit: identical entry + provider + model + prompt + input.
|
PreflightOutcome::Cached(existing) => return Ok((
|
||||||
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,
|
StatusCode::OK,
|
||||||
serde_json::to_value(&existing)
|
serde_json::to_value(&existing).map(Json)
|
||||||
.map(Json)
|
|
||||||
.map_err(|e| ApiError::internal(&e.to_string()))?,
|
.map_err(|e| ApiError::internal(&e.to_string()))?,
|
||||||
));
|
)),
|
||||||
}
|
PreflightOutcome::Pending { input, summary_uid } => (input, summary_uid),
|
||||||
}
|
};
|
||||||
}
|
|
||||||
|
|
||||||
// 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 archive_path = mounted.archive_path.clone();
|
||||||
let entry_uid_bg = entry_uid.clone();
|
|
||||||
let summary_uid_bg = summary_uid.clone();
|
let summary_uid_bg = summary_uid.clone();
|
||||||
let summary_options_bg = summary_options;
|
|
||||||
tokio::task::spawn_blocking(move || {
|
tokio::task::spawn_blocking(move || {
|
||||||
// summarize_entry owns the pending → running → completed/failed
|
// The input and attempt were claimed during blocking preflight, so no
|
||||||
// transitions for this same row (the cache key is identical, so the
|
// row is reset and no archive content is read twice.
|
||||||
// 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) =
|
if let Err(e) =
|
||||||
summarizer::summarize_entry(
|
summarizer::summarize_prebuilt_entry(
|
||||||
&archive_paths,
|
&archive_paths,
|
||||||
&entry_uid_bg,
|
input,
|
||||||
summary_options_bg,
|
&summary_uid_bg,
|
||||||
provider.as_ref(),
|
provider.as_ref(),
|
||||||
summarizer::PROMPT_VERSION,
|
|
||||||
)
|
)
|
||||||
{
|
{
|
||||||
eprintln!("warn: summary {summary_uid_bg}: {e:#}");
|
eprintln!("warn: summary {summary_uid_bg}: {e:#}");
|
||||||
|
|
@ -3162,6 +3173,97 @@ mod tests {
|
||||||
assert!(!error.contains("v1 unsupported"));
|
assert!(!error.contains("v1 unsupported"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn summary_request_for_missing_entry_returns_not_found_before_preflight() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let (registry, _archive_path, auth_path) = make_test_registry(&dir);
|
||||||
|
let session_cookie = make_test_session(&auth_path);
|
||||||
|
let previous_codex_cli = std::env::var_os("ARCHIVR_CODEX_CLI");
|
||||||
|
unsafe { std::env::set_var("ARCHIVR_CODEX_CLI", "/usr/bin/false") };
|
||||||
|
|
||||||
|
let response = app(registry, auth_path)
|
||||||
|
.oneshot(
|
||||||
|
Request::builder()
|
||||||
|
.method("POST")
|
||||||
|
.uri("/api/archives/test/entries/ent_missing/summary")
|
||||||
|
.header("content-type", "application/json")
|
||||||
|
.header("cookie", &session_cookie)
|
||||||
|
.body(json_body(&serde_json::json!({ "provider": "codex_cli" })))
|
||||||
|
.unwrap(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
match previous_codex_cli {
|
||||||
|
Some(value) => unsafe { std::env::set_var("ARCHIVR_CODEX_CLI", value) },
|
||||||
|
None => unsafe { std::env::remove_var("ARCHIVR_CODEX_CLI") },
|
||||||
|
}
|
||||||
|
|
||||||
|
assert_eq!(response.status(), StatusCode::NOT_FOUND);
|
||||||
|
assert_eq!(body_json(response).await["error"], "entry not found");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn public_summary_endpoints_hide_failed_diagnostics_but_authenticated_users_keep_them() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let (registry, archive_path, auth_path) = make_test_registry(&dir);
|
||||||
|
let entry = make_test_entry(&archive_path);
|
||||||
|
let session = make_test_session(&auth_path);
|
||||||
|
let conn = database::open_or_initialize(&archive_path).unwrap();
|
||||||
|
let summary_uid = database::upsert_pending_entry_summary(
|
||||||
|
&conn, entry.id, "codex_cli", None, summarizer::PROMPT_VERSION, "failed-public-test",
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
database::update_entry_summary_status(
|
||||||
|
&conn, &summary_uid, "failed", None, Some("provider secret: raw diagnostic"),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
drop(conn);
|
||||||
|
let collection = api_make_collection(
|
||||||
|
registry.clone(), auth_path.clone(), &session, "Public summaries", "public-summaries", 3, false,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
api_add_to_coll(
|
||||||
|
registry.clone(), auth_path.clone(), &session, &collection, &entry.entry_uid, 3,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
|
||||||
|
let public_summary = app(registry.clone(), auth_path.clone())
|
||||||
|
.oneshot(Request::builder()
|
||||||
|
.uri(format!("/api/archives/test/entries/{}/summary", entry.entry_uid))
|
||||||
|
.body(Body::empty()).unwrap())
|
||||||
|
.await.unwrap();
|
||||||
|
assert_eq!(public_summary.status(), StatusCode::OK);
|
||||||
|
assert!(body_json(public_summary).await["summary"].is_null());
|
||||||
|
|
||||||
|
let public_detail = app(registry.clone(), auth_path.clone())
|
||||||
|
.oneshot(Request::builder()
|
||||||
|
.uri(format!("/api/archives/test/entries/{}", entry.entry_uid))
|
||||||
|
.body(Body::empty()).unwrap())
|
||||||
|
.await.unwrap();
|
||||||
|
assert_eq!(public_detail.status(), StatusCode::OK);
|
||||||
|
assert!(body_json(public_detail).await["latest_summary"].is_null());
|
||||||
|
|
||||||
|
let authenticated_summary = app(registry.clone(), auth_path.clone())
|
||||||
|
.oneshot(Request::builder()
|
||||||
|
.uri(format!("/api/archives/test/entries/{}/summary", entry.entry_uid))
|
||||||
|
.header("cookie", &session)
|
||||||
|
.body(Body::empty()).unwrap())
|
||||||
|
.await.unwrap();
|
||||||
|
let authenticated_summary = body_json(authenticated_summary).await;
|
||||||
|
assert_eq!(authenticated_summary["summary"]["status"], "failed");
|
||||||
|
assert_eq!(authenticated_summary["summary"]["error_text"], "provider secret: raw diagnostic");
|
||||||
|
|
||||||
|
let authenticated_detail = app(registry, auth_path)
|
||||||
|
.oneshot(Request::builder()
|
||||||
|
.uri(format!("/api/archives/test/entries/{}", entry.entry_uid))
|
||||||
|
.header("cookie", &session)
|
||||||
|
.body(Body::empty()).unwrap())
|
||||||
|
.await.unwrap();
|
||||||
|
let authenticated_detail = body_json(authenticated_detail).await;
|
||||||
|
assert_eq!(authenticated_detail["latest_summary"]["status"], "failed");
|
||||||
|
assert_eq!(authenticated_detail["latest_summary"]["error_text"], "provider secret: raw diagnostic");
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn background_summary_failure_stores_safe_copy_only_for_unsupported_content() {
|
fn background_summary_failure_stores_safe_copy_only_for_unsupported_content() {
|
||||||
let dir = tempfile::tempdir().unwrap();
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue