mirror of
https://github.com/thegeneralist01/archivr
synced 2026-10-09 12:55:00 +02:00
fix: recover and preserve summary attempts
This commit is contained in:
parent
492b1a3168
commit
e5b9fc8ae0
4 changed files with 231 additions and 93 deletions
|
|
@ -590,70 +590,70 @@ async fn request_entry_summary_handler(
|
|||
include_images: body.include_images,
|
||||
};
|
||||
|
||||
let conn = database::open_or_initialize(&mounted.archive_path)?;
|
||||
let entry_id = database::entry_id_for_uid(&conn, &entry_uid)?
|
||||
.ok_or(ApiError::not_found("entry not found"))?;
|
||||
|
||||
// 2. Extract the same content the summarizer will feed the model, so the
|
||||
// digest below is the identical cache key summarize_entry will compute.
|
||||
let input = summarizer::build_summary_input(&archive_paths, &entry_uid, summary_options)
|
||||
.map_err(|e| {
|
||||
if summarizer::is_unsupported_summary_content_error(&e) {
|
||||
ApiError::bad_request(summarizer::UNSUPPORTED_SUMMARY_CONTENT_MESSAGE)
|
||||
} else {
|
||||
ApiError::bad_request(&format!("{e:#}"))
|
||||
}
|
||||
})?;
|
||||
|
||||
// 3. Cache hit: identical entry + provider + model + prompt + input.
|
||||
if !body.force {
|
||||
if let Some(existing) = database::find_entry_summary(
|
||||
&conn,
|
||||
entry_id,
|
||||
provider.kind(),
|
||||
provider.model(),
|
||||
summarizer::PROMPT_VERSION,
|
||||
&input.input_sha256,
|
||||
)? {
|
||||
if existing.status == "completed" {
|
||||
return Ok((
|
||||
StatusCode::OK,
|
||||
serde_json::to_value(&existing)
|
||||
.map(Json)
|
||||
.map_err(|e| ApiError::internal(&e.to_string()))?,
|
||||
));
|
||||
// 2. Preflight extraction and SQLite cache/attempt work are synchronous
|
||||
// core operations, so keep them off the Axum runtime. This also means the
|
||||
// input claimed here is passed directly to the provider worker below.
|
||||
let preflight_paths = archive_paths.clone();
|
||||
let preflight_uid = entry_uid.clone();
|
||||
let provider_kind = provider.kind().to_string();
|
||||
let provider_model = provider.model().map(str::to_string);
|
||||
let force = body.force;
|
||||
enum PreflightOutcome {
|
||||
Cached(database::EntrySummaryRecord),
|
||||
Pending { input: summarizer::SummaryInput, summary_uid: String },
|
||||
}
|
||||
let outcome = tokio::task::spawn_blocking(move || -> anyhow::Result<PreflightOutcome> {
|
||||
let input = summarizer::build_summary_input(&preflight_paths, &preflight_uid, summary_options)?;
|
||||
let conn = database::open_or_initialize(&preflight_paths.archive_path)?;
|
||||
let entry_id = database::entry_id_for_uid(&conn, &preflight_uid)?
|
||||
.ok_or_else(|| anyhow::anyhow!("entry not found"))?;
|
||||
if !force {
|
||||
if let Some(existing) = database::find_entry_summary(
|
||||
&conn, entry_id, &provider_kind, provider_model.as_deref(),
|
||||
summarizer::PROMPT_VERSION, &input.input_sha256,
|
||||
)? {
|
||||
if existing.status == "completed" {
|
||||
return Ok(PreflightOutcome::Cached(existing));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 4. Claim the row up front so the 202 can name it and the client can poll
|
||||
// immediately, before the worker thread has done anything.
|
||||
let summary_uid = database::upsert_pending_entry_summary(
|
||||
&conn,
|
||||
entry_id,
|
||||
provider.kind(),
|
||||
provider.model(),
|
||||
summarizer::PROMPT_VERSION,
|
||||
&input.input_sha256,
|
||||
)?;
|
||||
drop(conn);
|
||||
let summary_uid = database::upsert_pending_entry_summary(
|
||||
&conn, entry_id, &provider_kind, provider_model.as_deref(),
|
||||
summarizer::PROMPT_VERSION, &input.input_sha256,
|
||||
)?;
|
||||
Ok(PreflightOutcome::Pending { input, summary_uid })
|
||||
})
|
||||
.await
|
||||
.map_err(|e| ApiError::internal(&format!("summary preflight task failed: {e}")))?
|
||||
.map_err(|e| {
|
||||
if summarizer::is_unsupported_summary_content_error(&e) {
|
||||
ApiError::bad_request(summarizer::UNSUPPORTED_SUMMARY_CONTENT_MESSAGE)
|
||||
} else if format!("{e:#}") == "entry not found" {
|
||||
ApiError::not_found("entry not found")
|
||||
} else {
|
||||
ApiError::bad_request(&format!("{e:#}"))
|
||||
}
|
||||
})?;
|
||||
let (input, summary_uid) = match outcome {
|
||||
PreflightOutcome::Cached(existing) => return Ok((
|
||||
StatusCode::OK,
|
||||
serde_json::to_value(&existing).map(Json)
|
||||
.map_err(|e| ApiError::internal(&e.to_string()))?,
|
||||
)),
|
||||
PreflightOutcome::Pending { input, summary_uid } => (input, summary_uid),
|
||||
};
|
||||
|
||||
let archive_path = mounted.archive_path.clone();
|
||||
let entry_uid_bg = entry_uid.clone();
|
||||
let summary_uid_bg = summary_uid.clone();
|
||||
let summary_options_bg = summary_options;
|
||||
tokio::task::spawn_blocking(move || {
|
||||
// summarize_entry owns the pending → running → completed/failed
|
||||
// transitions for this same row (the cache key is identical, so the
|
||||
// upsert above and the one inside it resolve to one row). We only have
|
||||
// to catch the case where it fails before it can record anything.
|
||||
// The input and attempt were claimed during blocking preflight, so no
|
||||
// row is reset and no archive content is read twice.
|
||||
if let Err(e) =
|
||||
summarizer::summarize_entry(
|
||||
summarizer::summarize_prebuilt_entry(
|
||||
&archive_paths,
|
||||
&entry_uid_bg,
|
||||
summary_options_bg,
|
||||
input,
|
||||
&summary_uid_bg,
|
||||
provider.as_ref(),
|
||||
summarizer::PROMPT_VERSION,
|
||||
)
|
||||
{
|
||||
eprintln!("warn: summary {summary_uid_bg}: {e:#}");
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue