From af2a3ceb4962694d549f9da468136a34133eb675 Mon Sep 17 00:00:00 2001 From: Tyler Hallada Date: Wed, 2 Sep 2026 20:44:28 +0000 Subject: [PATCH] Bisect provider-rejected batches in triage and deep assessment DeepSeek's content filter rejects a whole request (400 "Content Exists Risk") when one article trips it, which cost the other articles in the batch their assessment and retried them every run. New curate/batch.rs runs both stages through a bisecting runner: a rejected batch is split until the offending article is isolated, that article is retried once on the editor provider when it is a different one, and a still-rejected article is recorded as a provider_rejected assessment row so it is not retried for assessment_reuse_days. Cache reuse accepts rows from either configured model. Each stage logs reused/requested/rejected counts, the curation: line shows rejections when non-zero, explain prints them, and the llm_assess span reports the deep-set size. Implemented by a Claude agent from an orchestrator brief; verified fmt/clippy(-W dead_code)/test green (354 lib tests). Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01A1rCLQeKBgnBo3oTgHuTMe --- Cargo.lock | 2 +- README.md | 13 + docs/plans/2026-08-15-implementation-notes.md | 10 + docs/runbooks/curation-v2-migration.md | 7 +- src/curate/admit.rs | 45 + src/curate/assess.rs | 292 +++++- src/curate/batch.rs | 839 ++++++++++++++++++ src/curate/embedding.rs | 4 +- src/curate/llm.rs | 2 +- src/curate/mod.rs | 16 +- src/curate/telemetry.rs | 56 ++ src/curate/triage.rs | 446 ++++++++-- src/pipeline.rs | 59 +- src/report.rs | 47 +- 14 files changed, 1711 insertions(+), 127 deletions(-) create mode 100644 src/curate/batch.rs diff --git a/Cargo.lock b/Cargo.lock index c592979..f0291f1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -815,7 +815,7 @@ dependencies = [ [[package]] name = "daily-epub" -version = "0.1.0" +version = "0.2.0" dependencies = [ "ammonia", "anyhow", diff --git a/README.md b/README.md index 2bca22f..a0a9610 100644 --- a/README.md +++ b/README.md @@ -685,5 +685,18 @@ From spec §7, plus what implementation turned up: beginning/middle/end sample, separates editorial quality from reader fit, and records descriptive facets. Utility is normalized over the deep set; embedding leader clusters cap near-duplicates before the 60-item editor shortlist. +- **DeepSeek's content filter rejects whole batches.** A `400 Content Exists + Risk` refuses the entire request when any one article in it trips the input + filter, without saying which. Triage and deep assessment therefore bisect a + rejected batch (halves, then quarters, down to single articles) so the other + articles keep their assessment. A single article the bulk provider still + refuses is retried once on the editor provider with the same prompt when that + is a different provider; if that also fails, an `article_assessments` row with + `kind = 'provider_rejected'` and a NULL score is written so the article is not + sent again for `assessment_reuse_days`, and it is ranked on its other signals. + See them with `explain` (`triage: rejected by provider — deepseek: …`), the + `triage:` / `assess:` log lines (`… 3 rejected (2 recovered on gemini)`), the + ` · N rejected` suffix on the `curation:` line, or + `sqlite3 /var/lib/daily-epub/daily-epub.db "select stage, count(*) from article_assessments where kind = 'provider_rejected' group by stage"`. - **One reader, one issue per day.** There is no multi-user support and no weekly/retrospective edition (spec §6). diff --git a/docs/plans/2026-08-15-implementation-notes.md b/docs/plans/2026-08-15-implementation-notes.md index 6f72995..691b67e 100644 --- a/docs/plans/2026-08-15-implementation-notes.md +++ b/docs/plans/2026-08-15-implementation-notes.md @@ -211,3 +211,13 @@ implementer needs that are easy to get wrong: concurrency (`triage`, `assess`) is the bulk provider's `max_concurrent_requests`; the summaries fan out at the summary provider's. `Models { bulk, editor, summaries }` in the colophon and Behind the paper stay model ids taken from the built clients. +- **Provider rejections** (2026-09-02): `curate::batch` bisects a triage or deep batch that + fails with a non-transient error (`Api`, `Refusal`, `EmptyResponse`, or a response that + parses to zero items) down to single articles; a rejected single is retried once on the + editor client when it is another provider. An article both refuse gets an + `article_assessments` row with `kind = 'provider_rejected'`, `score`/`fit` NULL, + `rationale = ": "`, the bulk `model` and the stage's + `prompt_version`; the cache loaders skip it (leaving the assessment absent) while it is within + `assessment_reuse_days`, `--rescore` ignores it, and `admit::hygiene` never treats it as a low + score. Because a recovered article's row carries the editor's model, reuse accepts rows whose + `model` is either configured model (`triage::reusable_models`). diff --git a/docs/runbooks/curation-v2-migration.md b/docs/runbooks/curation-v2-migration.md index 2e5de4c..87e5c18 100644 --- a/docs/runbooks/curation-v2-migration.md +++ b/docs/runbooks/curation-v2-migration.md @@ -257,7 +257,12 @@ sudo systemd-run --quiet --wait --pty --collect --uid=daily-epub --gid=daily-epu ``` produces the same date's paper with Gemini as the editor. Triage and deep assessments are cached -for three days, so the second run costs only the editor, summaries and the Brief. Compare the two +per article for three days, so the second run re-requests only articles that have no row yet: +ones ingested since the first run, and ones whose earlier batch was rejected by the provider's +content filter (those are bisected, retried once on the editor provider, and then recorded as +`provider_rejected` so they are not retried daily). The `triage:` and `assess:` log lines say how +many were reused versus requested. The bulk-side cost of the second run is therefore small; the +editor, summaries and the Brief are the real spend. Compare the two lineups, the `why` lines and the Brief side by side, and the `providers:` cost line. To switch for good, set `editor = "gemini"` in `[llm]` (and `summary_model` stays `editor`, so summaries move with it). The same trick works for the bulk role: `DAILY_EPUB_LLM__BULK=gemini`. diff --git a/src/curate/admit.rs b/src/curate/admit.rs index 436edc6..b3995d0 100644 --- a/src/curate/admit.rs +++ b/src/curate/admit.rs @@ -494,4 +494,49 @@ mod tests { .expect("thin row"); assert_eq!(reason, "recently_rejected"); } + + #[tokio::test] + async fn provider_rejected_rows_do_not_mark_an_article_recently_rejected() { + let dir = tempfile::tempdir().expect("tempdir"); + let db = Db::open_and_migrate(&dir.path().join("hygiene.db")) + .await + .expect("db"); + sqlx::query( + "INSERT INTO articles (id, canonical_url, title, first_seen) VALUES + (1, 'https://example.com/1', 'Refused', '2026-09-02T00:00:00Z')", + ) + .execute(db.pool()) + .await + .expect("article"); + sqlx::query( + "INSERT INTO article_assessments + (article_id, stage, model, prompt_version, score, fit, kind, rationale, assessed_at) + VALUES (1, 'triage', 'model', 1, NULL, NULL, 'provider_rejected', + 'deepseek: Content Exists Risk', '2026-09-02T04:00:00Z'), + (1, 'deep', 'model', 1, NULL, NULL, 'provider_rejected', + 'deepseek: Content Exists Risk', '2026-09-02T04:00:00Z')", + ) + .execute(db.pool()) + .await + .expect("rejection rows"); + let run_id = db + .start_run( + "2026-09-02".parse().expect("date"), + "2026-09-02T05:30:00Z".parse().expect("timestamp"), + ) + .await + .expect("run"); + let eligible = hygiene( + &db, + run_id, + vec![article(1, "Refused", 500)], + "2026-09-02".parse().expect("date"), + &CurationConfig::default(), + "2026-09-02T05:30:00Z".parse().expect("timestamp"), + ) + .await + .expect("hygiene"); + assert_eq!(eligible.len(), 1, "a NULL score is not a low score"); + assert_eq!(eligible[0].article.id, 1); + } } diff --git a/src/curate/assess.rs b/src/curate/assess.rs index 0f2eefe..f920ab5 100644 --- a/src/curate/assess.rs +++ b/src/curate/assess.rs @@ -3,12 +3,13 @@ use std::collections::{HashMap, HashSet}; use std::fmt::Write as _; -use futures::{StreamExt, stream}; use jiff::Timestamp; use serde_json::Value; use sqlx::Row as _; +use super::batch::{Assessed, BatchRunner, Scored, StageSummary, run_batches}; use super::llm::{LlmClient, strip_code_fence}; +use super::triage::{PROVIDER_REJECTED, reusable_models, write_rejection}; use super::{prompt_text, truncate_words}; use crate::db::{Db, fmt_ts, parse_ts}; use crate::types::{ArticleId, Candidate, Deep, Facets}; @@ -316,10 +317,24 @@ fn as_bool(value: &Value) -> Option { }) } +impl Assessed for DeepItem { + fn article_id(&self) -> ArticleId { + self.id + } +} + +/// Assess the admitted set on `llm` (cache only when `None`), bisecting +/// rejected batches and retrying rejected singles on `fallback` when it is +/// another provider (see [`super::batch`]). +/// +/// Cached rows written by either configured model are reused; a fresh +/// `provider_rejected` row skips the article and leaves its deep assessment +/// absent, so ranking falls back to the present signals (§12.3). #[allow(clippy::too_many_arguments)] pub async fn run( db: &Db, llm: Option<&LlmClient>, + fallback: Option<&LlmClient>, model: &str, candidates: &mut [Candidate], batch_size: usize, @@ -330,22 +345,31 @@ pub async fn run( assessed_at: Timestamp, temperature: f32, sections: &[String], -) -> anyhow::Result { +) -> anyhow::Result { let positions = candidates .iter() .enumerate() .filter(|(_, candidate)| candidate.stage == "admitted") .map(|(index, candidate)| (candidate.article.id, index)) .collect::>(); + let mut summary = StageSummary { + stage: "assess", + pool: positions.len(), + ..StageSummary::default() + }; + let mut known_rejected = HashSet::new(); if !rescore && !positions.is_empty() { let since = assessed_at - jiff::Span::new().hours(assessment_reuse_days.max(0) * 24); + let models = reusable_models(model, fallback); let rows = sqlx::query( - "SELECT article_id, score, fit, kind, facets_json, rationale, category, + "SELECT article_id, model, score, fit, kind, facets_json, rationale, category, paywalled_guess, assessed_at FROM article_assessments - WHERE stage = 'deep' AND model = ? AND prompt_version = ? AND assessed_at >= ?", + WHERE stage = 'deep' AND model IN (?, ?) AND prompt_version = ? + AND assessed_at >= ?", ) - .bind(model) + .bind(models[0]) + .bind(models[1]) .bind(DEEP_PROMPT_VERSION) .bind(fmt_ts(since)) .fetch_all(db.pool()) @@ -355,6 +379,10 @@ pub async fn run( let Some(index) = positions.get(&id).copied() else { continue; }; + if row.get::, _>("kind").as_deref() == Some(PROVIDER_REJECTED) { + known_rejected.insert(id); + continue; + } let (Some(quality), Some(fit)) = ( row.get::, _>("score"), row.get::, _>("fit"), @@ -376,52 +404,47 @@ pub async fn run( .unwrap_or_default(), paywalled_guess: row.get::("paywalled_guess") != 0, facets, - model: model.to_string(), + model: row.get::("model"), prompt_version: DEEP_PROMPT_VERSION, assessed_at: parse_ts( "article_assessments.assessed_at", &row.get::("assessed_at"), )?, }); + summary.reused += 1; } } + summary.known_rejected = known_rejected.len(); let pending = candidates .iter() - .filter(|candidate| candidate.stage == "admitted" && candidate.assessment.deep.is_none()) + .filter(|candidate| { + candidate.stage == "admitted" + && candidate.assessment.deep.is_none() + && !known_rejected.contains(&candidate.article.id) + }) .collect::>(); if let Some(llm) = llm { - let prompts = pending + summary.requested = pending.len(); + let batches = pending .chunks(batch_size.max(1)) - .map(|batch| { - let allowed = batch - .iter() - .map(|candidate| candidate.article.id) - .collect::>(); - (allowed, build_batch_prompt(batch, sections)) - }) + .map(<[&Candidate]>::to_vec) .collect::>(); - let results = stream::iter(prompts) - .map(|(allowed, prompt)| async move { - if let Err(error) = llm.meter.check_budget() { - tracing::warn!(%error, "bulk budget tripped; skipping deep batch"); - return Vec::new(); - } - match llm.complete(&prompt, temperature, true).await { - Ok(raw) => parse_deep_response(&raw, sections) - .into_iter() - .filter(|item| allowed.contains(&item.id)) - .collect(), - Err(error) => { - tracing::warn!(%error, "deep batch failed; its articles remain unassessed"); - Vec::new() - } - } - }) - .buffer_unordered(max_concurrent_requests.max(1)) - .collect::>>() - .await; - for item in results.into_iter().flatten() { + summary.batches = batches.len(); + let build_prompt = |batch: &[&Candidate]| build_batch_prompt(batch, sections); + let parse = |raw: &str| parse_deep_response(raw, sections); + let runner = BatchRunner { + llm, + fallback, + temperature, + build_prompt: &build_prompt, + parse: &parse, + }; + summary.fallback_provider = runner.fallback_provider(); + let outcome = run_batches(&runner, batches, max_concurrent_requests).await; + summary.rejected = outcome.rejected; + summary.recovered = outcome.recovered; + for Scored { item, model } in outcome.items { let Some(index) = positions.get(&item.id).copied() else { continue; }; @@ -432,7 +455,7 @@ pub async fn run( rationale: item.rationale, paywalled_guess: item.paywalled_guess, facets: item.facets, - model: model.to_string(), + model, prompt_version: DEEP_PROMPT_VERSION, assessed_at, }; @@ -450,7 +473,7 @@ pub async fn run( paywalled_guess = excluded.paywalled_guess, assessed_at = excluded.assessed_at", ) .bind(item.id) - .bind(model) + .bind(&deep.model) .bind(DEEP_PROMPT_VERSION) .bind(profile_version) .bind(deep.quality) @@ -464,6 +487,19 @@ pub async fn run( .execute(db.pool()) .await?; candidates[index].assessment.deep = Some(deep); + summary.applied += 1; + } + for rejection in &outcome.rejections { + write_rejection( + db, + "deep", + rejection, + model, + DEEP_PROMPT_VERSION, + profile_version, + assessed_at, + ) + .await?; } } for candidate in candidates @@ -472,17 +508,16 @@ pub async fn run( { candidate.stage = "assessed".into(); } - Ok(candidates - .iter() - .filter(|candidate| candidate.assessment.deep.is_some()) - .count()) + tracing::info!("{}", summary.info_line()); + Ok(summary) } #[cfg(test)] mod tests { use super::*; use crate::config::{CurationConfig, ProviderConfig}; - use crate::curate::llm::{MockBackend, PriceTable, UsageMeter}; + use crate::curate::batch::tests::{FilterBackend, deep_answer_for}; + use crate::curate::llm::{ChatBackend, MockBackend, PriceTable, UsageMeter}; use crate::curate::prefilter::tests::{article, with_social}; use crate::curate::signals::{Neighbour, TopInterest}; use crate::types::{TokenUsage, Triage}; @@ -555,6 +590,7 @@ mod tests { run( db, llm, + None, "deepseek-v4-flash", candidates, batch_size, @@ -568,6 +604,7 @@ mod tests { ) .await .expect("deep assessment never aborts the run") + .assessed() } #[test] @@ -951,6 +988,7 @@ mod tests { let assessed = run( &db, Some(&llm), + None, "deepseek-v4-flash", &mut candidates, 1, @@ -963,7 +1001,8 @@ mod tests { §ions(), ) .await - .expect("assessment"); + .expect("assessment") + .assessed(); assert_eq!(assessed, 1, "only the first batch ran"); assert_eq!(backend.calls(), 1); assert!(llm.meter.budget_exceeded()); @@ -1032,6 +1071,7 @@ mod tests { run( &db, None, + None, "other-model", &mut other_model, 8, @@ -1054,6 +1094,7 @@ mod tests { run( &db, None, + None, "deepseek-v4-flash", &mut stale_prompt, 8, @@ -1079,6 +1120,7 @@ mod tests { run( &db, None, + None, "deepseek-v4-flash", &mut old, 8, @@ -1117,4 +1159,166 @@ mod tests { .expect("row"); assert_eq!(stored, 4.0); } + + fn named_client(provider: &str, model: &str, backend: Arc) -> LlmClient { + LlmClient::with_backend_options( + provider, + model, + "SYSTEM".into(), + None, + UsageMeter::with_prices(PriceTable::from(&ProviderConfig::deepseek()), 10.0), + backend, + ) + } + + fn titled(n: i64) -> Vec { + (1..=n) + .map(|id| { + let mut c = candidate(id, 900); + c.article.title = format!("Piece {id}"); + c + }) + .collect() + } + + async fn assess_with( + db: &Db, + llm: &LlmClient, + fallback: Option<&LlmClient>, + candidates: &mut [Candidate], + rescore: bool, + at: Timestamp, + ) -> StageSummary { + run( + db, + Some(llm), + fallback, + "deepseek-v4-flash", + candidates, + 4, + 4, + 3, + rescore, + Some(1), + at, + 0.3, + §ions(), + ) + .await + .expect("deep assessment never aborts the run") + } + + #[tokio::test] + async fn rejected_deep_batches_are_bisected_and_rejections_persisted() { + let (_dir, db) = db_with_articles(&[1, 2, 3, 4]).await; + let backend = FilterBackend::deep(&["Piece 3"]); + let llm = named_client("deepseek", "deepseek-v4-flash", backend.clone()); + let mut candidates = titled(4); + let summary = assess_with(&db, &llm, None, &mut candidates, false, timestamp()).await; + assert_eq!(backend.calls(), 5); + assert_eq!(summary.assessed(), 3); + assert_eq!((summary.rejected, summary.recovered), (1, 0)); + assert_eq!( + summary.info_line(), + "assess: 4 in pool · 0 reused · 4 requested in 1 batches · 1 rejected" + ); + assert!(candidates[2].assessment.deep.is_none()); + assert_eq!( + candidates[2].stage, "admitted", + "kept for present-signal ranking" + ); + assert!(candidates[0].assessment.deep.is_some()); + let row = sqlx::query( + "SELECT model, score, fit, kind, rationale, facets_json FROM article_assessments + WHERE article_id = 3 AND stage = 'deep'", + ) + .fetch_one(db.pool()) + .await + .expect("rejection row"); + assert_eq!(row.get::("model"), "deepseek-v4-flash"); + assert_eq!(row.get::, _>("score"), None); + assert_eq!(row.get::, _>("fit"), None); + assert_eq!( + row.get::, _>("kind").as_deref(), + Some(PROVIDER_REJECTED) + ); + assert!( + row.get::, _>("rationale") + .is_some_and(|why| why.starts_with("deepseek: 400 Bad Request")) + ); + assert_eq!(row.get::, _>("facets_json"), None); + + // Honoured next run, ignored under --rescore. + let mut cached = titled(4); + let summary = assess_with(&db, &llm, None, &mut cached, false, timestamp()).await; + assert_eq!(backend.calls(), 5); + assert_eq!( + (summary.reused, summary.known_rejected, summary.requested), + (3, 1, 0) + ); + assert_eq!(summary.rejected_total(), 1); + assert!(cached[2].assessment.deep.is_none()); + let mut rescored = titled(4); + let summary = assess_with(&db, &llm, None, &mut rescored, true, timestamp()).await; + assert_eq!(summary.requested, 4); + assert_eq!(backend.calls(), 10); + + // Expired: retried (and rejected again). + let later = timestamp() + jiff::Span::new().hours(4 * 24); + let mut expired = titled(4); + let summary = assess_with(&db, &llm, None, &mut expired, false, later).await; + assert_eq!(summary.known_rejected, 0); + assert_eq!(summary.requested, 4); + } + + #[tokio::test] + async fn deep_fallback_rows_carry_the_editor_model_and_are_reused() { + let (_dir, db) = db_with_articles(&[1, 2]).await; + let backend = FilterBackend::deep(&["Piece 2"]); + let llm = named_client("deepseek", "deepseek-v4-flash", backend.clone()); + let editor_backend = Arc::new(MockBackend::new()); + editor_backend.push(deep_answer_for(&[2]), TokenUsage::default()); + let editor = named_client("anthropic", "claude-opus-5", editor_backend.clone()); + let mut candidates = titled(2); + let summary = assess_with( + &db, + &llm, + Some(&editor), + &mut candidates, + false, + timestamp(), + ) + .await; + assert_eq!(editor_backend.calls(), 1); + assert_eq!( + editor_backend.prompts()[0].user, + backend.prompts_for_single(2), + "the same single-article prompt" + ); + assert_eq!((summary.rejected, summary.recovered), (1, 1)); + assert_eq!(summary.assessed(), 2); + let deep = candidates[1].assessment.deep.as_ref().expect("recovered"); + assert_eq!(deep.model, "claude-opus-5"); + assert_eq!(candidates[1].stage, "assessed"); + let rejected: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM article_assessments WHERE kind = 'provider_rejected'", + ) + .fetch_one(db.pool()) + .await + .expect("count"); + assert_eq!(rejected, 0); + + let mut cached = titled(2); + let summary = assess_with(&db, &llm, Some(&editor), &mut cached, false, timestamp()).await; + assert_eq!(summary.reused, 2, "the editor's row is reusable"); + assert_eq!( + cached[1] + .assessment + .deep + .as_ref() + .map(|deep| deep.model.as_str()), + Some("claude-opus-5") + ); + assert_eq!(backend.calls(), 3); + } } diff --git a/src/curate/batch.rs b/src/curate/batch.rs new file mode 100644 index 0000000..720b737 --- /dev/null +++ b/src/curate/batch.rs @@ -0,0 +1,839 @@ +//! Bisecting batch runner shared by triage and deep assessment. +//! +//! DeepSeek's input content filter rejects a whole request when any one +//! article in the batch trips it (`400 Content Exists Risk`) and does not say +//! which. Sending the batch once and giving up cost every other article its +//! assessment and retried them all the next day. Instead, a batch that fails +//! with a non-transient error is split in half and both halves are sent +//! again, down to single articles. A single article that still fails is +//! "rejected": it is retried once on the editor client when that runs on a +//! different provider, and otherwise reported so the caller can persist a +//! `provider_rejected` row and stop asking (§10, §12.1, §17). +//! +//! Transient errors are retried inside [`LlmClient::complete`] and, once +//! exhausted, leave the batch unassessed as before; a budget trip stops the +//! stage. Concurrency is the caller's `max_concurrent_requests` across the +//! original batches; the bisection inside one batch runs sequentially. + +use std::collections::HashSet; + +use futures::{StreamExt, stream}; + +use super::llm::{LlmClient, LlmError}; +use crate::types::{ArticleId, Candidate}; + +/// Longest rejection message kept in `article_assessments.rationale`. +pub const REJECTION_MESSAGE_CHARS: usize = 200; + +/// A parsed per-article item the runner can attribute to a candidate. +pub trait Assessed { + fn article_id(&self) -> ArticleId; +} + +/// One stage's prompt builder, parser and clients. +pub struct BatchRunner<'a, T> { + /// The bulk client every batch goes to first. + pub llm: &'a LlmClient, + /// The editor client, tried once per rejected article when it runs on a + /// different provider than `llm` and its meter is not tripped. + pub fallback: Option<&'a LlmClient>, + pub temperature: f32, + pub build_prompt: &'a (dyn Fn(&[&Candidate]) -> String + Sync), + pub parse: &'a (dyn Fn(&str) -> Vec + Sync), +} + +impl BatchRunner<'_, T> { + /// The editor client when it is a real alternative to the bulk one. + pub fn usable_fallback(&self) -> Option<&LlmClient> { + self.fallback + .filter(|fallback| fallback.provider() != self.llm.provider()) + .filter(|fallback| fallback.meter.check_budget().is_ok()) + } + + /// The name of the provider rejected articles are retried on, if any. + pub fn fallback_provider(&self) -> Option { + self.fallback + .filter(|fallback| fallback.provider() != self.llm.provider()) + .map(|fallback| fallback.provider().to_string()) + } +} + +/// An item together with the model that produced it (bulk or editor). +#[derive(Debug, Clone, PartialEq)] +pub struct Scored { + pub item: T, + pub model: String, +} + +/// A single article both the bulk provider and the fallback refused. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Rejection { + pub id: ArticleId, + /// The bulk provider's name. + pub provider: String, + /// The bulk provider's message, at most [`REJECTION_MESSAGE_CHARS`] long. + pub message: String, +} + +impl Rejection { + /// `: `, the `rationale` of a `provider_rejected` row. + pub fn rationale(&self) -> String { + format!("{}: {}", self.provider, self.message) + } +} + +/// What a set of batches produced. +#[derive(Debug, Clone, PartialEq)] +pub struct BatchOutcome { + pub items: Vec>, + /// Articles nobody would assess; the caller persists these. + pub rejections: Vec, + /// Single articles the bulk provider rejected, recovered or not. + pub rejected: usize, + /// Of `rejected`, those the fallback client assessed. + pub recovered: usize, + /// Requests made on the bulk client. + pub requests: usize, + /// True when the bulk budget tripped and work was left undone. + pub budget_stopped: bool, +} + +impl Default for BatchOutcome { + fn default() -> Self { + Self { + items: Vec::new(), + rejections: Vec::new(), + rejected: 0, + recovered: 0, + requests: 0, + budget_stopped: false, + } + } +} + +impl BatchOutcome { + fn absorb(&mut self, other: Self) { + self.items.extend(other.items); + self.rejections.extend(other.rejections); + self.rejected += other.rejected; + self.recovered += other.recovered; + self.requests += other.requests; + self.budget_stopped |= other.budget_stopped; + } +} + +/// The per-stage totals behind the `triage:` / `assess:` log line and the +/// run report's `*_reused` / `*_rejected` counts. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct StageSummary { + pub stage: &'static str, + /// Articles the stage was responsible for this run. + pub pool: usize, + /// Served from cached rows within `assessment_reuse_days`. + pub reused: usize, + /// Skipped because a `provider_rejected` row is still fresh. + pub known_rejected: usize, + /// Sent to the bulk provider. + pub requested: usize, + /// Original batches before any bisection. + pub batches: usize, + /// Assessments written this run (bulk and editor). + pub applied: usize, + /// Single articles the bulk provider rejected this run. + pub rejected: usize, + /// Of `rejected`, those the editor client assessed instead. + pub recovered: usize, + pub fallback_provider: Option, +} + +impl StageSummary { + /// Articles carrying this stage's assessment after the run. + pub fn assessed(&self) -> usize { + self.reused + self.applied + } + + /// Articles without an assessment because a provider rejected them: + /// this run's unrecovered rejections plus the cached ones. + pub fn rejected_total(&self) -> usize { + self.known_rejected + self.rejected.saturating_sub(self.recovered) + } + + /// `triage: 398 in pool · 210 reused · 188 requested in 8 batches · 3 rejected (2 recovered on gemini)` + pub fn info_line(&self) -> String { + let mut line = format!( + "{}: {} in pool · {} reused · {} requested in {} batches · {} rejected", + self.stage, self.pool, self.reused, self.requested, self.batches, self.rejected + ); + if self.recovered > 0 { + line.push_str(&format!( + " ({} recovered on {})", + self.recovered, + self.fallback_provider.as_deref().unwrap_or("editor") + )); + } + if self.known_rejected > 0 { + line.push_str(&format!(" · {} known rejected", self.known_rejected)); + } + line + } +} + +/// Errors the provider will keep returning for the same input. +fn is_rejection(error: &LlmError) -> bool { + matches!( + error, + LlmError::Api { .. } | LlmError::Refusal { .. } | LlmError::EmptyResponse { .. } + ) +} + +/// The provider's own words, without the ` request failed:` prefix. +fn rejection_message(error: &LlmError) -> String { + match error { + LlmError::Api { message, .. } => message.clone(), + LlmError::Refusal { .. } => "returned a refusal".into(), + LlmError::EmptyResponse { .. } => "returned an empty completion".into(), + other => other.to_string(), + } +} + +fn truncate_chars(text: &str, max_chars: usize) -> String { + let collapsed = text.split_whitespace().collect::>().join(" "); + if collapsed.chars().count() <= max_chars { + return collapsed; + } + let mut cut = collapsed + .chars() + .take(max_chars.saturating_sub(1)) + .collect::(); + cut.push('…'); + cut +} + +enum Failure { + Rejected(String), + NoItems, +} + +struct Work<'c> { + batch: Vec<&'c Candidate>, + /// A batch that parses to zero items is split once; its halves are not. + zero_split: bool, +} + +/// Run every batch on the bulk client, `max_concurrent` in flight. +pub async fn run_batches( + runner: &BatchRunner<'_, T>, + batches: Vec>, + max_concurrent: usize, +) -> BatchOutcome { + let outcomes = stream::iter(batches) + .map(|batch| run_one(runner, batch)) + .buffer_unordered(max_concurrent.max(1)) + .collect::>() + .await; + let mut merged = BatchOutcome::default(); + for outcome in outcomes { + merged.absorb(outcome); + } + merged +} + +/// One original batch, bisected as far as it needs to be. +async fn run_one( + runner: &BatchRunner<'_, T>, + batch: Vec<&Candidate>, +) -> BatchOutcome { + let original = batch.len(); + let provider = runner.llm.provider(); + let mut out = BatchOutcome::default(); + let mut warned = false; + // Explicit LIFO stack instead of recursion: the left half is pushed last + // so it runs first, which keeps the request order predictable. + let mut stack = vec![Work { + batch, + zero_split: true, + }]; + while let Some(work) = stack.pop() { + if let Err(error) = runner.llm.meter.check_budget() { + tracing::warn!(%error, "bulk budget tripped; skipping the rest of the batch"); + out.budget_stopped = true; + return out; + } + let size = work.batch.len(); + let prompt = (runner.build_prompt)(&work.batch); + let allowed = work + .batch + .iter() + .map(|candidate| candidate.article.id) + .collect::>(); + out.requests += 1; + let failure = match runner.llm.complete(&prompt, runner.temperature, true).await { + Ok(raw) => { + let items = (runner.parse)(&raw) + .into_iter() + .filter(|item| allowed.contains(&item.article_id())) + .collect::>(); + if items.is_empty() { + Failure::NoItems + } else { + out.items.extend(items.into_iter().map(|item| Scored { + item, + model: runner.llm.model.clone(), + })); + continue; + } + } + Err(LlmError::BudgetExceeded { spent, limit }) => { + tracing::warn!( + spent, + limit, + "bulk budget tripped; skipping the rest of the batch" + ); + out.budget_stopped = true; + return out; + } + Err(error) if is_rejection(&error) => Failure::Rejected(rejection_message(&error)), + Err(error) => { + tracing::warn!(%error, size, "batch failed; its articles remain unassessed"); + continue; + } + }; + if size == 1 { + let candidate = work.batch[0]; + let message = match failure { + Failure::Rejected(message) => message, + Failure::NoItems => "returned no assessment for the article".into(), + }; + resolve_single(runner, candidate, &prompt, message, &mut out).await; + continue; + } + match failure { + Failure::Rejected(message) => { + if !warned { + tracing::warn!( + provider, + message = %truncate_chars(&message, REJECTION_MESSAGE_CHARS), + "batch of {original} rejected by {provider}; bisecting" + ); + warned = true; + } + push_halves(&mut stack, work.batch, work.zero_split); + } + Failure::NoItems if work.zero_split => { + tracing::warn!(provider, size, "batch parsed to zero items; bisecting once"); + push_halves(&mut stack, work.batch, false); + } + Failure::NoItems => { + tracing::warn!( + provider, + size, + "half batch parsed to zero items again; its articles remain unassessed" + ); + } + } + } + out +} + +fn push_halves<'c>(stack: &mut Vec>, batch: Vec<&'c Candidate>, zero_split: bool) { + let mut left = batch; + let right = left.split_off(left.len() / 2); + stack.push(Work { + batch: right, + zero_split, + }); + stack.push(Work { + batch: left, + zero_split, + }); +} + +/// A single article the bulk provider would not assess: try the editor +/// client once, else report it as rejected. +async fn resolve_single( + runner: &BatchRunner<'_, T>, + candidate: &Candidate, + prompt: &str, + message: String, + out: &mut BatchOutcome, +) { + let id = candidate.article.id; + let provider = runner.llm.provider(); + out.rejected += 1; + if let Some(fallback) = runner.usable_fallback() { + let fallback_provider = fallback.provider(); + match fallback.complete(prompt, runner.temperature, true).await { + Ok(raw) => { + if let Some(item) = (runner.parse)(&raw) + .into_iter() + .find(|item| item.article_id() == id) + { + tracing::info!( + article_id = id, + provider = fallback_provider, + "article rejected by {provider}; assessed on {fallback_provider} instead" + ); + out.items.push(Scored { + item, + model: fallback.model.clone(), + }); + out.recovered += 1; + return; + } + tracing::warn!( + article_id = id, + provider = fallback_provider, + "fallback returned no assessment for the rejected article" + ); + } + Err(error) => { + tracing::warn!( + article_id = id, + provider = fallback_provider, + %error, + "fallback failed for the rejected article" + ); + } + } + } + let message = truncate_chars(&message, REJECTION_MESSAGE_CHARS); + tracing::warn!( + article_id = id, + provider, + %message, + "article rejected by the provider; recorded so it is not retried" + ); + out.rejections.push(Rejection { + id, + provider: provider.to_string(), + message, + }); +} + +#[cfg(test)] +pub(crate) mod tests { + use super::*; + use crate::config::ProviderConfig; + use crate::curate::llm::{ + ChatBackend, ChatCompletion, ChatRequest, MockBackend, PriceTable, UsageMeter, + }; + use crate::curate::prefilter::tests::article; + use crate::curate::triage::{TriageItem, build_batch_prompt, parse_triage_response}; + use crate::http::RetryPolicy; + use crate::types::TokenUsage; + use std::future::Future; + use std::pin::Pin; + use std::sync::{Arc, Mutex}; + use std::time::Duration; + + /// A backend that answers for every id it finds in the prompt unless the + /// prompt mentions a forbidden title, in which case it rejects the whole + /// request the way DeepSeek's content filter does. + #[derive(Debug)] + pub(crate) struct FilterBackend { + forbidden: Vec, + /// Prompts answered with `{"articles": []}` (by request index). + empty_on: Vec, + answer: fn(&[i64]) -> String, + seen: Mutex>, + } + + impl FilterBackend { + /// Answers in the triage shape. + pub(crate) fn new(forbidden: &[&str], empty_on: &[usize]) -> Arc { + Self::with_answer(forbidden, empty_on, answer_for) + } + + /// Answers in the deep-assessment shape. + pub(crate) fn deep(forbidden: &[&str]) -> Arc { + Self::with_answer(forbidden, &[], deep_answer_for) + } + + fn with_answer( + forbidden: &[&str], + empty_on: &[usize], + answer: fn(&[i64]) -> String, + ) -> Arc { + Arc::new(Self { + forbidden: forbidden.iter().map(|s| s.to_string()).collect(), + empty_on: empty_on.to_vec(), + answer, + seen: Mutex::new(Vec::new()), + }) + } + + pub(crate) fn calls(&self) -> usize { + self.seen.lock().map(|seen| seen.len()).unwrap_or(0) + } + + /// The prompt of the first request that carried only `id`. + pub(crate) fn prompts_for_single(&self, id: i64) -> String { + self.seen + .lock() + .ok() + .and_then(|seen| { + seen.iter() + .find(|prompt| ids_in(prompt) == vec![id]) + .cloned() + }) + .unwrap_or_default() + } + + /// The ids in each request, in order. + pub(crate) fn requests(&self) -> Vec> { + self.seen + .lock() + .map(|seen| seen.iter().map(|prompt| ids_in(prompt)).collect()) + .unwrap_or_default() + } + } + + fn ids_in(prompt: &str) -> Vec { + prompt + .lines() + .filter_map(|line| line.strip_prefix("--- id: ")) + .filter_map(|id| id.trim().parse().ok()) + .collect() + } + + /// A triage answer for `ids`. + pub(crate) fn answer_for(ids: &[i64]) -> String { + let items = ids + .iter() + .map(|id| format!(r#"{{"id":{id},"interest":6,"kind":"essay","why":"fine"}}"#)) + .collect::>() + .join(","); + format!(r#"{{"articles":[{items}]}}"#) + } + + /// A deep-assessment answer for `ids`. + pub(crate) fn deep_answer_for(ids: &[i64]) -> String { + let items = ids + .iter() + .map(|id| { + format!( + r#"{{"id":{id},"quality":7,"fit":6,"category":"Top Stories","rationale":"fine","facets":{{"format":"analysis_essay"}}}}"# + ) + }) + .collect::>() + .join(","); + format!(r#"{{"articles":[{items}]}}"#) + } + + impl ChatBackend for FilterBackend { + fn complete<'a>( + &'a self, + req: ChatRequest, + ) -> Pin> + Send + 'a>> { + Box::pin(async move { + let index = self.calls(); + if let Ok(mut seen) = self.seen.lock() { + seen.push(req.user.clone()); + } + if self.forbidden.iter().any(|title| req.user.contains(title)) { + return Err(LlmError::api( + "deepseek", + "400 Bad Request: {\"error\":{\"message\":\"Content Exists Risk\"}}", + )); + } + let content = if self.empty_on.contains(&index) { + r#"{"articles":[]}"#.to_string() + } else { + (self.answer)(&ids_in(&req.user)) + }; + Ok(ChatCompletion { + content, + usage: TokenUsage::default(), + }) + }) + } + } + + fn client(provider: &str, model: &str, backend: Arc) -> LlmClient { + LlmClient::with_backend_options( + provider, + model, + "SYSTEM".into(), + None, + UsageMeter::with_prices(PriceTable::from(&ProviderConfig::deepseek()), 10.0), + backend, + ) + } + + fn candidates(n: i64) -> Vec { + (1..=n) + .map(|id| Candidate::new(article(id, &format!("Title number {id}"), 800), false)) + .collect() + } + + async fn run( + llm: &LlmClient, + fallback: Option<&LlmClient>, + candidates: &[Candidate], + batch_size: usize, + ) -> BatchOutcome { + let runner = BatchRunner { + llm, + fallback, + temperature: 0.3, + build_prompt: &build_batch_prompt, + parse: &parse_triage_response, + }; + let all = candidates.iter().collect::>(); + let batches = all.chunks(batch_size).map(<[&Candidate]>::to_vec).collect(); + run_batches(&runner, batches, 4).await + } + + fn sorted_ids(outcome: &BatchOutcome) -> Vec { + let mut ids = outcome + .items + .iter() + .map(|scored| scored.item.id) + .collect::>(); + ids.sort_unstable(); + ids + } + + #[tokio::test] + async fn one_bad_article_in_four_costs_five_requests_and_one_rejection() { + let backend = FilterBackend::new(&["Title number 3"], &[]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let candidates = candidates(4); + let outcome = run(&llm, None, &candidates, 4).await; + + // [1,2,3,4] is rejected → split. [1,2] succeeds. [3,4] is rejected → + // split. [3] alone is rejected → recorded. [4] succeeds. 1 + 2 + 2 = 5. + assert_eq!( + backend.requests(), + vec![vec![1, 2, 3, 4], vec![1, 2], vec![3, 4], vec![3], vec![4]] + ); + assert_eq!(outcome.requests, 5); + assert_eq!(sorted_ids(&outcome), vec![1, 2, 4]); + assert!(outcome.items.iter().all(|s| s.model == "deepseek-v4-flash")); + assert_eq!(outcome.rejected, 1); + assert_eq!(outcome.recovered, 0); + assert_eq!(outcome.rejections.len(), 1); + let rejection = &outcome.rejections[0]; + assert_eq!(rejection.id, 3); + assert_eq!(rejection.provider, "deepseek"); + assert!(rejection.message.contains("Content Exists Risk")); + assert!( + rejection + .rationale() + .starts_with("deepseek: 400 Bad Request"), + "{}", + rejection.rationale() + ); + assert!(!outcome.budget_stopped); + } + + #[tokio::test] + async fn every_article_rejected_gives_every_article_a_rejection_in_2n_minus_1_requests() { + let titles = (1..=8) + .map(|id| format!("Title number {id}")) + .collect::>(); + let backend = + FilterBackend::new(&titles.iter().map(String::as_str).collect::>(), &[]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let candidates = candidates(8); + let outcome = run(&llm, None, &candidates, 8).await; + assert_eq!(backend.calls(), 2 * 8 - 1); + assert_eq!(outcome.requests, 15); + assert!(outcome.items.is_empty()); + let mut rejected = outcome + .rejections + .iter() + .map(|rejection| rejection.id) + .collect::>(); + rejected.sort_unstable(); + assert_eq!(rejected, (1..=8).collect::>()); + assert_eq!(outcome.rejected, 8); + } + + #[tokio::test] + async fn a_rejected_article_is_recovered_on_an_editor_of_another_provider() { + let backend = FilterBackend::new(&["Title number 2"], &[]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let editor_backend = Arc::new(MockBackend::new()); + editor_backend.push(answer_for(&[2]), TokenUsage::default()); + let editor = client("gemini", "gemini-3.8-flash", editor_backend.clone()); + let candidates = candidates(2); + let outcome = run(&llm, Some(&editor), &candidates, 2).await; + assert_eq!(backend.requests(), vec![vec![1, 2], vec![1], vec![2]]); + assert_eq!(editor_backend.calls(), 1); + let prompt = editor_backend.prompts()[0].clone(); + assert_eq!(prompt.system.as_str(), "SYSTEM", "same system prompt"); + assert!(prompt.user.contains("--- id: 2\n"), "same article prompt"); + assert_eq!(sorted_ids(&outcome), vec![1, 2]); + let recovered = outcome + .items + .iter() + .find(|scored| scored.item.id == 2) + .expect("recovered"); + assert_eq!(recovered.model, "gemini-3.8-flash"); + assert!(outcome.rejections.is_empty(), "no rejection row"); + assert_eq!((outcome.rejected, outcome.recovered), (1, 1)); + + // The same provider name is no alternative: no fallback attempt. + let backend = FilterBackend::new(&["Title number 2"], &[]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let same_backend = Arc::new(MockBackend::new()); + same_backend.push(answer_for(&[2]), TokenUsage::default()); + let same = client("deepseek", "deepseek-v4-flash", same_backend.clone()); + let outcome = run(&llm, Some(&same), &candidates, 2).await; + assert_eq!(same_backend.calls(), 0); + assert_eq!(outcome.rejections.len(), 1); + assert_eq!(outcome.recovered, 0); + + // A tripped editor meter is skipped too, and a failing editor still + // yields a rejection with the bulk provider's message. + let backend = FilterBackend::new(&["Title number 2"], &[]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let tripped_backend = Arc::new(MockBackend::new()); + let tripped = client("gemini", "gemini-3.8-flash", tripped_backend.clone()); + tripped.meter.preload_cost(100.0); + let outcome = run(&llm, Some(&tripped), &candidates, 2).await; + assert_eq!(tripped_backend.calls(), 0); + assert_eq!(outcome.rejections.len(), 1); + + let backend = FilterBackend::new(&["Title number 2"], &[]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let failing_backend = Arc::new(MockBackend::new()); + failing_backend.push_llm_error(LlmError::refusal("gemini")); + let failing = client("gemini", "gemini-3.8-flash", failing_backend.clone()); + let outcome = run(&llm, Some(&failing), &candidates, 2).await; + assert_eq!(failing_backend.calls(), 1); + assert_eq!(outcome.rejections.len(), 1); + assert_eq!(outcome.rejections[0].provider, "deepseek"); + assert!( + outcome.rejections[0] + .message + .contains("Content Exists Risk") + ); + } + + #[tokio::test] + async fn zero_items_bisect_once_and_a_single_zero_is_a_rejection() { + // Batch of 4 parses to zero items → its two halves are sent; both fine. + let backend = FilterBackend::new(&[], &[0]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let candidates = candidates(4); + let outcome = run(&llm, None, &candidates, 4).await; + assert_eq!( + backend.requests(), + vec![vec![1, 2, 3, 4], vec![1, 2], vec![3, 4]] + ); + assert_eq!(sorted_ids(&outcome), vec![1, 2, 3, 4]); + assert!(outcome.rejections.is_empty()); + + // A half that parses to zero items again is not split further: its + // articles simply stay unassessed for this run. + let backend = FilterBackend::new(&[], &[0, 1]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let outcome = run(&llm, None, &candidates, 4).await; + assert_eq!(backend.calls(), 3); + assert_eq!(sorted_ids(&outcome), vec![3, 4]); + assert!(outcome.rejections.is_empty()); + + // A single article that parses to zero items is a rejection. + let backend = FilterBackend::new(&[], &[0]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let outcome = run(&llm, None, &candidates[..1], 1).await; + assert_eq!(backend.calls(), 1); + assert!(outcome.items.is_empty()); + assert_eq!(outcome.rejections.len(), 1); + assert_eq!( + outcome.rejections[0].rationale(), + "deepseek: returned no assessment for the article" + ); + } + + #[tokio::test] + async fn transient_and_budget_failures_are_not_bisected() { + let backend = Arc::new(MockBackend::new()); + for _ in 0..3 { + backend.push_llm_error(LlmError::Transient { + provider: "deepseek".into(), + message: "503".into(), + }); + } + let llm = + client("deepseek", "deepseek-v4-flash", backend.clone()).with_retry(RetryPolicy { + max_attempts: 3, + base_delay: Duration::from_millis(1), + max_delay: Duration::from_millis(2), + }); + let candidates = candidates(4); + // Transient errors are retried inside the client and, once exhausted, + // leave the batch unassessed without bisecting it. + let outcome = run(&llm, None, &candidates, 4).await; + assert_eq!(backend.calls(), 3, "three attempts, one batch"); + assert!(outcome.items.is_empty()); + assert!( + outcome.rejections.is_empty(), + "transient is not a rejection" + ); + assert_eq!(outcome.rejected, 0); + + let backend = Arc::new(MockBackend::new()); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + llm.meter.preload_cost(100.0); + let outcome = run(&llm, None, &candidates, 2).await; + assert_eq!(backend.calls(), 0); + assert!(outcome.budget_stopped); + } + + #[test] + fn summary_line_matches_the_documented_shape() { + let summary = StageSummary { + stage: "triage", + pool: 398, + reused: 210, + known_rejected: 0, + requested: 188, + batches: 8, + applied: 187, + rejected: 3, + recovered: 2, + fallback_provider: Some("gemini".into()), + }; + assert_eq!( + summary.info_line(), + "triage: 398 in pool · 210 reused · 188 requested in 8 batches · 3 rejected (2 recovered on gemini)" + ); + assert_eq!(summary.assessed(), 397); + assert_eq!(summary.rejected_total(), 1); + let quiet = StageSummary { + stage: "assess", + pool: 120, + reused: 20, + known_rejected: 2, + requested: 98, + batches: 13, + applied: 98, + ..StageSummary::default() + }; + assert_eq!( + quiet.info_line(), + "assess: 120 in pool · 20 reused · 98 requested in 13 batches · 0 rejected · 2 known rejected" + ); + assert_eq!(quiet.rejected_total(), 2); + } + + #[test] + fn rejection_messages_are_short_and_single_line() { + let long = "x".repeat(500); + let cut = truncate_chars(&long, REJECTION_MESSAGE_CHARS); + assert_eq!(cut.chars().count(), REJECTION_MESSAGE_CHARS); + assert!(cut.ends_with('…')); + assert_eq!(truncate_chars("a\n b", 10), "a b"); + assert_eq!( + rejection_message(&LlmError::refusal("anthropic")), + "returned a refusal" + ); + assert_eq!( + rejection_message(&LlmError::api("deepseek", "400: nope")), + "400: nope" + ); + } +} diff --git a/src/curate/embedding.rs b/src/curate/embedding.rs index 710cfa4..48d8e5c 100644 --- a/src/curate/embedding.rs +++ b/src/curate/embedding.rs @@ -416,8 +416,8 @@ pub fn decode_blob(bytes: &[u8], dimension: usize) -> Result, Embedding }); } let mut vector = Vec::with_capacity(dimension); - for chunk in bytes.chunks_exact(4) { - let value = f32::from_le_bytes([chunk[0], chunk[1], chunk[2], chunk[3]]); + for chunk in bytes.as_chunks::<4>().0 { + let value = f32::from_le_bytes(*chunk); if !value.is_finite() { return Err(EmbeddingError::NonFinite); } diff --git a/src/curate/llm.rs b/src/curate/llm.rs index 0682ad2..af78e36 100644 --- a/src/curate/llm.rs +++ b/src/curate/llm.rs @@ -627,7 +627,7 @@ impl LlmClient { } #[cfg(test)] - fn with_retry(mut self, retry: RetryPolicy) -> Self { + pub(crate) fn with_retry(mut self, retry: RetryPolicy) -> Self { self.retry = retry; self } diff --git a/src/curate/mod.rs b/src/curate/mod.rs index cab708e..5ae78ca 100644 --- a/src/curate/mod.rs +++ b/src/curate/mod.rs @@ -12,6 +12,7 @@ pub mod admit; pub mod assess; +pub mod batch; pub mod editor; pub mod editorial; pub mod embedding; @@ -50,23 +51,30 @@ impl Curator { /// A no-op under `--skip-llm`: like triage, nothing is read or written and /// utility falls back to the present signals (§12.3, §17). When the bulk /// provider is down or its budget trips, cached rows are still reused and - /// the failed batches simply stay unassessed. + /// the failed batches simply stay unassessed. A batch the provider's + /// content filter rejects is bisected, and a rejected single article is + /// retried on the editor when that is another provider (`curate::batch`). pub async fn assess( &self, candidates: &mut [Candidate], rescore: bool, profile_version: Option, assessed_at: Timestamp, - ) -> anyhow::Result { + ) -> anyhow::Result { let Some(bulk) = self.llms.bulk.as_ref() else { tracing::info!("--skip-llm: deep assessment skipped"); - return Ok(0); + return Ok(batch::StageSummary::default()); }; - let span = tracing::info_span!("llm_assess", candidates = candidates.len()); + let admitted = candidates + .iter() + .filter(|candidate| candidate.stage == "admitted") + .count(); + let span = tracing::info_span!("llm_assess", candidates = admitted); let _guard = span.enter(); assess::run( &self.db, Some(bulk), + self.llms.editor.as_ref(), &bulk.model, candidates, self.config.llm.deep_batch_size, diff --git a/src/curate/telemetry.rs b/src/curate/telemetry.rs index 171e14d..5a576bd 100644 --- a/src/curate/telemetry.rs +++ b/src/curate/telemetry.rs @@ -429,6 +429,12 @@ pub async fn render_explain(db: &Db, row: &ExplainRow) -> Result, _>("rationale") .unwrap_or_default(); + if assessment.get::, _>("kind").as_deref() + == Some(super::triage::PROVIDER_REJECTED) + { + let _ = writeln!(out, " {stage}: rejected by provider — {rationale}"); + continue; + } if stage == "deep" { let _ = writeln!( out, @@ -1007,6 +1013,56 @@ mod tests { assert_eq!(parsed["exploration"], true); } + #[tokio::test] + async fn explain_names_provider_rejections_instead_of_scores() { + let (_dir, db) = db_with_articles(&[1]).await; + let run_id = db.start_run(date(), Timestamp::now()).await.unwrap(); + write( + &db, + &CandidateRun { + run_id, + article_id: 1, + stage: "admitted", + excluded_reason: None, + admitted_by: Some("[\"interest\"]"), + signals_json: "{}", + utility: None, + rank_utility: None, + cluster_id: None, + cluster_rank: None, + editor_why: None, + }, + ) + .await + .unwrap(); + sqlx::query( + "INSERT INTO article_assessments + (article_id, stage, model, prompt_version, score, fit, kind, rationale, assessed_at) + VALUES (1, 'triage', 'deepseek-v4-flash', 1, NULL, NULL, 'provider_rejected', + 'deepseek: 400 Bad Request: Content Exists Risk', '2026-09-02T04:00:00Z'), + (1, 'deep', 'deepseek-v4-flash', 1, NULL, NULL, 'provider_rejected', + 'deepseek: returned a refusal', '2026-09-02T04:00:00Z')", + ) + .execute(db.pool()) + .await + .unwrap(); + let text = explain(&db, date(), None, &ExplainTarget::Article(1)) + .await + .unwrap(); + assert!( + text.contains( + " triage: rejected by provider — deepseek: 400 Bad Request: Content Exists Risk\n" + ), + "{text}" + ); + assert!( + text.contains(" deep: rejected by provider — deepseek: returned a refusal\n"), + "{text}" + ); + assert!(!text.contains("interest —"), "{text}"); + assert!(!text.contains("quality —"), "{text}"); + } + #[tokio::test] async fn rows_are_upserted_with_every_column_replaced() { let (_dir, db) = db_with_articles(&[1]).await; diff --git a/src/curate/triage.rs b/src/curate/triage.rs index c0b03c0..34e10f0 100644 --- a/src/curate/triage.rs +++ b/src/curate/triage.rs @@ -3,17 +3,20 @@ use std::collections::{HashMap, HashSet}; use std::fmt::Write as _; -use futures::{StreamExt, stream}; use jiff::Timestamp; use serde_json::Value; use sqlx::Row as _; +use super::batch::{Assessed, BatchRunner, Rejection, Scored, StageSummary, run_batches}; use super::llm::{LlmClient, strip_code_fence}; use super::{prompt_text, truncate_words}; use crate::db::{Db, fmt_ts, parse_ts}; use crate::types::{ArticleId, Candidate, Triage}; pub const TRIAGE_PROMPT_VERSION: i64 = 1; +/// The `kind` of an `article_assessments` row recording that the provider +/// refused the article; `score` (and `fit`) are NULL and `rationale` says why. +pub const PROVIDER_REJECTED: &str = "provider_rejected"; pub const TRIAGE_INSTRUCTIONS: &str = r#"TASK: first-pass triage of today's candidate articles for The Daily EPUB. You see only each article's opening. Decide how much THIS reader (profile in your @@ -337,10 +340,36 @@ fn compare_signal( .then_with(|| left_id.cmp(&right_id)) } +impl Assessed for TriageItem { + fn article_id(&self) -> ArticleId { + self.id + } +} + +/// The models whose cached rows a stage may reuse: the bulk model and, when +/// an editor on another provider can recover rejected articles, its model. +pub fn reusable_models<'a>(model: &'a str, fallback: Option<&'a LlmClient>) -> [&'a str; 2] { + [ + model, + fallback + .map(|client| client.model.as_str()) + .unwrap_or(model), + ] +} + +/// Triage every pool article without a fresh cached assessment on the bulk +/// client, bisecting rejected batches and retrying rejected singles on +/// `fallback` when it is another provider (see [`super::batch`]). +/// +/// Cached rows count as reused whether the bulk or the editor model wrote +/// them; a fresh `provider_rejected` row skips the article and leaves its +/// triage absent. Articles with a reusable deep row are reused as well: they +/// need no triage to be admitted. #[allow(clippy::too_many_arguments)] pub async fn run( db: &Db, llm: &LlmClient, + fallback: Option<&LlmClient>, candidates: &mut [Candidate], pool: &HashSet, batch_size: usize, @@ -350,36 +379,46 @@ pub async fn run( profile_version: Option, assessed_at: Timestamp, temperature: f32, -) -> anyhow::Result { +) -> anyhow::Result { + let positions = candidates + .iter() + .enumerate() + .map(|(index, candidate)| (candidate.article.id, index)) + .collect::>(); let mut reusable_deep = HashSet::new(); + let mut known_rejected = HashSet::new(); if !rescore { let since = assessed_at - jiff::Span::new().hours(assessment_reuse_days.max(0) * 24); + let models = reusable_models(&llm.model, fallback); let rows = sqlx::query( - "SELECT article_id, stage, score, kind, rationale, assessed_at + "SELECT article_id, stage, model, score, kind, rationale, assessed_at FROM article_assessments - WHERE model = ? AND assessed_at >= ? + WHERE model IN (?, ?) AND assessed_at >= ? AND ((stage = 'triage' AND prompt_version = ?) OR (stage = 'deep' AND prompt_version = ?))", ) - .bind(&llm.model) + .bind(models[0]) + .bind(models[1]) .bind(fmt_ts(since)) .bind(TRIAGE_PROMPT_VERSION) .bind(super::assess::DEEP_PROMPT_VERSION) .fetch_all(db.pool()) .await?; - let pool_ids = pool; - let positions = candidates - .iter() - .enumerate() - .map(|(index, candidate)| (candidate.article.id, index)) - .collect::>(); for row in rows { let id = row.get::("article_id"); - if !pool_ids.contains(&id) { + if !pool.contains(&id) { continue; } + let rejected = + row.get::, _>("kind").as_deref() == Some(PROVIDER_REJECTED); if row.get::("stage") == "deep" { - reusable_deep.insert(id); + if !rejected { + reusable_deep.insert(id); + } + continue; + } + if rejected { + known_rejected.insert(id); continue; } let Some(score) = row.get::, _>("score") else { @@ -398,7 +437,7 @@ pub async fn run( why: row .get::, _>("rationale") .unwrap_or_default(), - model: llm.model.clone(), + model: row.get::("model"), prompt_version: TRIAGE_PROMPT_VERSION, assessed_at: timestamp, }); @@ -406,52 +445,51 @@ pub async fn run( } } - let pending = candidates + let in_pool = candidates .iter() .filter(|candidate| { - pool.contains(&candidate.article.id) - && candidate.excluded_reason.is_none() - && candidate.assessment.triage.is_none() - && !reusable_deep.contains(&candidate.article.id) + pool.contains(&candidate.article.id) && candidate.excluded_reason.is_none() }) .collect::>(); - let prompts = pending - .chunks(batch_size.max(1)) - .map(|batch| { - let allowed = batch - .iter() - .map(|candidate| candidate.article.id) - .collect::>(); - (allowed, build_batch_prompt(batch)) - }) - .collect::>(); - let results = stream::iter(prompts) - .map(|(allowed, prompt)| async move { - if let Err(error) = llm.meter.check_budget() { - tracing::warn!(%error, "bulk budget tripped; skipping triage batch"); - return Vec::new(); - } - match llm.complete(&prompt, temperature, true).await { - Ok(raw) => parse_triage_response(&raw) - .into_iter() - .filter(|item| allowed.contains(&item.id)) - .collect(), - Err(error) => { - tracing::warn!(%error, "triage batch failed; its articles remain untriaged"); - Vec::new() - } - } - }) - .buffer_unordered(max_concurrent_requests.max(1)) - .collect::>>() - .await; - let positions = candidates + let pending = in_pool .iter() - .enumerate() - .map(|(index, candidate)| (candidate.article.id, index)) - .collect::>(); - let mut applied = 0; - for item in results.into_iter().flatten() { + .copied() + .filter(|candidate| { + candidate.assessment.triage.is_none() + && !reusable_deep.contains(&candidate.article.id) + && !known_rejected.contains(&candidate.article.id) + }) + .collect::>(); + let known_rejected = in_pool + .iter() + .filter(|candidate| known_rejected.contains(&candidate.article.id)) + .count(); + let mut summary = StageSummary { + stage: "triage", + pool: in_pool.len(), + reused: in_pool.len() - pending.len() - known_rejected, + known_rejected, + requested: pending.len(), + ..StageSummary::default() + }; + let batches = pending + .chunks(batch_size.max(1)) + .map(<[&Candidate]>::to_vec) + .collect::>(); + summary.batches = batches.len(); + let runner = BatchRunner { + llm, + fallback, + temperature, + build_prompt: &build_batch_prompt, + parse: &parse_triage_response, + }; + summary.fallback_provider = runner.fallback_provider(); + let outcome = run_batches(&runner, batches, max_concurrent_requests).await; + summary.rejected = outcome.rejected; + summary.recovered = outcome.recovered; + + for Scored { item, model } in outcome.items { let Some(index) = positions.get(&item.id).copied() else { continue; }; @@ -459,7 +497,7 @@ pub async fn run( interest: item.interest, kind: item.kind, why: item.why, - model: llm.model.clone(), + model, prompt_version: TRIAGE_PROMPT_VERSION, assessed_at, }; @@ -486,7 +524,19 @@ pub async fn run( .execute(db.pool()) .await?; candidates[index].assessment.triage = Some(triage); - applied += 1; + summary.applied += 1; + } + for rejection in &outcome.rejections { + write_rejection( + db, + "triage", + rejection, + &llm.model, + TRIAGE_PROMPT_VERSION, + profile_version, + assessed_at, + ) + .await?; } for candidate in candidates .iter_mut() @@ -494,19 +544,294 @@ pub async fn run( { candidate.stage = "triaged".into(); } - Ok(applied) + tracing::info!("{}", summary.info_line()); + Ok(summary) +} + +/// Persist a `provider_rejected` row so the article is not sent again while +/// the row is fresh (`score` and `fit` NULL; the rationale names the provider). +pub async fn write_rejection( + db: &Db, + stage: &str, + rejection: &Rejection, + model: &str, + prompt_version: i64, + profile_version: Option, + assessed_at: Timestamp, +) -> anyhow::Result<()> { + sqlx::query( + "INSERT INTO article_assessments + (article_id, stage, model, prompt_version, profile_version, score, fit, kind, + facets_json, rationale, category, paywalled_guess, assessed_at) + VALUES (?, ?, ?, ?, ?, NULL, NULL, ?, NULL, ?, NULL, 0, ?) + ON CONFLICT(article_id, stage) DO UPDATE SET + model = excluded.model, prompt_version = excluded.prompt_version, + profile_version = excluded.profile_version, score = NULL, fit = NULL, + kind = excluded.kind, facets_json = NULL, rationale = excluded.rationale, + category = NULL, paywalled_guess = 0, assessed_at = excluded.assessed_at", + ) + .bind(rejection.id) + .bind(stage) + .bind(model) + .bind(prompt_version) + .bind(profile_version) + .bind(PROVIDER_REJECTED) + .bind(rejection.rationale()) + .bind(fmt_ts(assessed_at)) + .execute(db.pool()) + .await?; + Ok(()) } #[cfg(test)] mod tests { use super::*; use crate::config::ProviderConfig; - use crate::curate::llm::{MockBackend, PriceTable, UsageMeter}; + use crate::curate::batch::tests::{FilterBackend, answer_for}; + use crate::curate::llm::{ChatBackend, MockBackend, PriceTable, UsageMeter}; use crate::curate::prefilter::tests::article; use crate::curate::signals::{Neighbour, TopInterest}; use crate::types::TokenUsage; use std::sync::Arc; + fn client(provider: &str, model: &str, backend: Arc) -> LlmClient { + LlmClient::with_backend_options( + provider, + model, + "SYSTEM".into(), + None, + UsageMeter::with_prices(PriceTable::from(&ProviderConfig::deepseek()), 10.0), + backend, + ) + } + + fn at() -> Timestamp { + "2026-09-02T05:30:00Z".parse().expect("timestamp") + } + + async fn db_with_articles(ids: &[i64]) -> (tempfile::TempDir, Db) { + let dir = tempfile::tempdir().expect("tempdir"); + let db = Db::open_and_migrate(&dir.path().join("triage.db")) + .await + .expect("db"); + for id in ids { + sqlx::query( + "INSERT INTO articles (id, canonical_url, title, first_seen) + VALUES (?, ?, 'A', '2026-09-02T00:00:00Z')", + ) + .bind(id) + .bind(format!("https://example.com/{id}")) + .execute(db.pool()) + .await + .expect("article"); + } + (dir, db) + } + + fn pool_of(n: i64) -> (Vec, HashSet) { + let candidates = (1..=n) + .map(|id| Candidate::new(article(id, &format!("Piece {id}"), 800), false)) + .collect::>(); + (candidates, (1..=n).collect()) + } + + /// `run` with the defaults these tests share: batches of 4, a 3-day cache. + async fn triage( + db: &Db, + llm: &LlmClient, + fallback: Option<&LlmClient>, + candidates: &mut [Candidate], + pool: &HashSet, + rescore: bool, + assessed_at: Timestamp, + ) -> StageSummary { + run( + db, + llm, + fallback, + candidates, + pool, + 4, + 4, + 3, + rescore, + Some(1), + assessed_at, + 0.3, + ) + .await + .expect("triage never aborts the run") + } + + async fn rejection_rows(db: &Db) -> Vec<(i64, String, String)> { + sqlx::query( + "SELECT article_id, model, rationale FROM article_assessments + WHERE stage = 'triage' AND kind = 'provider_rejected' AND score IS NULL + ORDER BY article_id", + ) + .fetch_all(db.pool()) + .await + .expect("rows") + .iter() + .map(|row| { + ( + row.get::("article_id"), + row.get::("model"), + row.get::, _>("rationale") + .unwrap_or_default(), + ) + }) + .collect() + } + + #[tokio::test] + async fn rejections_are_persisted_honoured_and_expire() { + let (_dir, db) = db_with_articles(&[1, 2, 3, 4]).await; + let backend = FilterBackend::new(&["Piece 3"], &[]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let (mut candidates, pool) = pool_of(4); + let summary = triage(&db, &llm, None, &mut candidates, &pool, false, at()).await; + assert_eq!(backend.calls(), 5, "1 + 2 + 2 for one bad article in four"); + assert_eq!( + (summary.pool, summary.requested, summary.batches), + (4, 4, 1) + ); + assert_eq!((summary.rejected, summary.recovered), (1, 0)); + assert_eq!(summary.applied, 3); + assert_eq!(summary.rejected_total(), 1); + assert_eq!( + summary.info_line(), + "triage: 4 in pool · 0 reused · 4 requested in 1 batches · 1 rejected" + ); + assert!(candidates[2].assessment.triage.is_none()); + assert_ne!(candidates[2].stage, "triaged"); + assert!( + candidates + .iter() + .filter(|candidate| candidate.article.id != 3) + .all(|candidate| candidate.assessment.triage.is_some()) + ); + let rows = rejection_rows(&db).await; + assert_eq!(rows.len(), 1); + assert_eq!((rows[0].0, rows[0].1.as_str()), (3, "deepseek-v4-flash")); + assert!( + rows[0].2.starts_with("deepseek: 400 Bad Request"), + "{}", + rows[0].2 + ); + + // Next run inside the window: nothing is asked for 3, and 1, 2, 4 are reused. + let (mut cached, _) = pool_of(4); + let summary = triage(&db, &llm, None, &mut cached, &pool, false, at()).await; + assert_eq!(backend.calls(), 5, "no request at all"); + assert_eq!((summary.reused, summary.known_rejected), (3, 1)); + assert_eq!(summary.requested, 0); + assert_eq!(summary.rejected_total(), 1); + assert!(cached[2].assessment.triage.is_none()); + assert!( + summary + .info_line() + .ends_with("0 rejected · 1 known rejected") + ); + + // `--rescore` ignores rejection rows like any other cached row. + let (mut rescored, _) = pool_of(4); + let summary = triage(&db, &llm, None, &mut rescored, &pool, true, at()).await; + assert_eq!(summary.requested, 4); + assert_eq!(summary.known_rejected, 0); + assert_eq!(backend.calls(), 10); + assert_eq!(rejection_rows(&db).await.len(), 1); + + // An expired rejection row is retried. + sqlx::query( + "UPDATE article_assessments SET assessed_at = '2026-08-01T00:00:00Z' + WHERE article_id = 3", + ) + .execute(db.pool()) + .await + .expect("age the row"); + let (mut expired, _) = pool_of(4); + let summary = triage(&db, &llm, None, &mut expired, &pool, false, at()).await; + assert_eq!(summary.requested, 1, "only the expired rejection"); + assert_eq!(backend.calls(), 11); + assert_eq!(backend.requests()[10], vec![3]); + assert_eq!(rejection_rows(&db).await.len(), 1, "rejected again"); + } + + #[tokio::test] + async fn fallback_assessments_carry_the_editor_model_and_are_reused() { + let (_dir, db) = db_with_articles(&[1, 2]).await; + let backend = FilterBackend::new(&["Piece 2"], &[]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let editor_backend = Arc::new(MockBackend::new()); + editor_backend.push(answer_for(&[2]), TokenUsage::default()); + let editor = client("gemini", "gemini-3.8-flash", editor_backend.clone()); + let (mut candidates, pool) = pool_of(2); + let summary = triage( + &db, + &llm, + Some(&editor), + &mut candidates, + &pool, + false, + at(), + ) + .await; + assert_eq!(editor_backend.calls(), 1); + assert_eq!((summary.rejected, summary.recovered), (1, 1)); + assert_eq!(summary.rejected_total(), 0); + assert_eq!( + summary.info_line(), + "triage: 2 in pool · 0 reused · 2 requested in 1 batches · 1 rejected (1 recovered on gemini)" + ); + let recovered = candidates[1].assessment.triage.as_ref().expect("recovered"); + assert_eq!(recovered.model, "gemini-3.8-flash"); + assert_eq!(candidates[1].stage, "triaged"); + assert!(rejection_rows(&db).await.is_empty(), "no rejection row"); + let stored: String = sqlx::query_scalar( + "SELECT model FROM article_assessments WHERE article_id = 2 AND stage = 'triage'", + ) + .fetch_one(db.pool()) + .await + .expect("row"); + assert_eq!(stored, "gemini-3.8-flash"); + + // The editor's row is reusable while the editor is configured… + let (mut cached, _) = pool_of(2); + let summary = triage(&db, &llm, Some(&editor), &mut cached, &pool, false, at()).await; + assert_eq!(summary.reused, 2); + assert_eq!(backend.calls(), 3); + assert_eq!( + cached[1] + .assessment + .triage + .as_ref() + .map(|triage| triage.model.as_str()), + Some("gemini-3.8-flash") + ); + + // …and not otherwise: without an editor only bulk-model rows count. + let (mut alone, _) = pool_of(2); + let summary = triage(&db, &llm, None, &mut alone, &pool, false, at()).await; + assert_eq!((summary.reused, summary.requested), (1, 1)); + } + + #[tokio::test] + async fn an_editor_on_the_bulk_provider_is_not_a_fallback() { + let (_dir, db) = db_with_articles(&[1, 2]).await; + let backend = FilterBackend::new(&["Piece 2"], &[]); + let llm = client("deepseek", "deepseek-v4-flash", backend.clone()); + let same_backend = Arc::new(MockBackend::new()); + same_backend.push(answer_for(&[2]), TokenUsage::default()); + let same = client("deepseek", "deepseek-v4-flash", same_backend.clone()); + let (mut candidates, pool) = pool_of(2); + let summary = triage(&db, &llm, Some(&same), &mut candidates, &pool, false, at()).await; + assert_eq!(same_backend.calls(), 0); + assert_eq!(summary.fallback_provider, None); + assert_eq!((summary.rejected, summary.recovered), (1, 0)); + assert_eq!(rejection_rows(&db).await.len(), 1); + } + const TRIAGE_FIXTURE: &str = include_str!(concat!( env!("CARGO_MANIFEST_DIR"), "/tests/fixtures/deepseek_triage_batch.json" @@ -597,6 +922,7 @@ mod tests { run( &db, &llm, + None, &mut first, &pool, 25, @@ -615,6 +941,7 @@ mod tests { run( &db, &llm, + None, &mut cached, &pool, 25, @@ -645,6 +972,7 @@ mod tests { run( &db, &llm, + None, &mut rescored, &pool, 25, diff --git a/src/pipeline.rs b/src/pipeline.rs index 2dd4f95..b934939 100644 --- a/src/pipeline.rs +++ b/src/pipeline.rs @@ -469,9 +469,10 @@ async fn run_stages( .flatten() .and_then(|value| value.parse().ok()); if let Some(bulk) = curator.llms.bulk.as_ref() { - if let Err(error) = triage::run( + match triage::run( db, bulk, + curator.llms.editor.as_ref(), &mut personalized, &triage_pool, config.llm.triage_batch_size, @@ -484,9 +485,13 @@ async fn run_stages( ) .await { - report.warn(format!( + Ok(summary) => { + report.counts.triage_reused = summary.reused as i64; + report.counts.triage_rejected = summary.rejected_total() as i64; + } + Err(error) => report.warn(format!( "triage degraded; admission continues without it: {error:#}" - )); + )), } } else { tracing::info!("--skip-llm or no bulk provider: triage skipped"); @@ -516,7 +521,7 @@ async fn run_stages( // --- Stage 9: deep assessment (§12.1) --- let stage = Timestamp::now(); - if let Err(error) = curator + match curator .assess( &mut personalized, ctx.rescore, @@ -525,9 +530,13 @@ async fn run_stages( ) .await { - report.warn(format!( + Ok(summary) => { + report.counts.deep_reused = summary.reused as i64; + report.counts.deep_rejected = summary.rejected_total() as i64; + } + Err(error) => report.warn(format!( "deep assessment degraded; ranking continues on present signals: {error:#}" - )); + )), } report.counts.assessed = personalized .iter() @@ -1666,16 +1675,28 @@ mod tests { assert_eq!(thin, "{}"); // Admission replaces the old prefilter and carries retriever telemetry. - // The bulk provider is "down": the client exists but every call fails, so - // the deep set is ranked on present signals and the editor falls back - // to utility order (§17). + // The bulk provider is "down": the client exists but the one deep batch + // fails transiently through every retry (a content-filter rejection + // would be bisected instead), so the deep set is ranked on present + // signals and the editor falls back to utility order (§17). let bulk_backend = Arc::new(ChatMockBackend::new()); + for _ in 0..3 { + bulk_backend.push_llm_error(crate::curate::llm::LlmError::Transient { + provider: "deepseek".into(), + message: "503".into(), + }); + } let bulk = LlmClient::with_backend( &h.config.providers["deepseek"].model, "SYSTEM".into(), UsageMeter::for_provider(&h.config.providers["deepseek"]), bulk_backend.clone(), - ); + ) + .with_retry(crate::http::RetryPolicy { + max_attempts: 3, + base_delay: std::time::Duration::from_millis(1), + max_delay: std::time::Duration::from_millis(2), + }); let curator = Curator::new( h.config.clone(), h.db.clone(), @@ -1699,9 +1720,21 @@ mod tests { let assessed = curator .assess(&mut features, false, None, now()) .await - .unwrap(); + .unwrap() + .assessed(); assert_eq!(assessed, 0, "every deep batch failed"); - assert_eq!(bulk_backend.calls(), 1, "one batch was attempted"); + assert_eq!( + bulk_backend.calls(), + 3, + "one batch was attempted, three times; nothing was bisected" + ); + let rejected: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM article_assessments WHERE kind = 'provider_rejected'", + ) + .fetch_one(h.db.pool()) + .await + .unwrap(); + assert_eq!(rejected, 0, "a transient failure is not a rejection"); assert!(features.iter().all(|c| c.assessment.deep.is_none())); let embeddings = features .iter() @@ -1743,7 +1776,7 @@ mod tests { let lineup = curator.select(candidates, run_date()).await.unwrap(); assert_eq!( bulk_backend.calls(), - 2, + 4, "the editor tried the bulk fallback" ); let selected = lineup diff --git a/src/report.rs b/src/report.rs index 17b5482..e2827b4 100644 --- a/src/report.rs +++ b/src/report.rs @@ -84,6 +84,13 @@ pub struct StageCounts { pub verdicts_in_prompt: i64, /// Articles with a reusable or newly produced triage assessment. pub triaged: i64, + /// Triage assessments served from cached rows. + #[serde(default)] + pub triage_reused: i64, + /// Pool articles without a triage because a provider rejected them + /// (this run's unrecovered rejections plus fresh `provider_rejected` rows). + #[serde(default)] + pub triage_rejected: i64, /// Articles admitted to close reading. pub admitted: i64, /// First admitting retriever counts. @@ -92,6 +99,13 @@ pub struct StageCounts { pub exploration_selected: i64, /// Deep assessments, whether reused or newly produced. pub assessed: i64, + /// Deep assessments served from cached rows. + #[serde(default)] + pub deep_reused: i64, + /// Admitted articles without a deep assessment because a provider + /// rejected them. + #[serde(default)] + pub deep_rejected: i64, /// Candidates shown to the editor after diversification. pub shortlisted: i64, /// Leader clusters formed over the deep set. @@ -287,10 +301,24 @@ impl RunReport { } else { providers }; + let rejected = |count: i64| { + if count > 0 { + format!(" · {count} rejected") + } else { + String::new() + } + }; [ format!( - "curation: {} considered → {} eligible → {} triaged → {} assessed → {} shortlisted → {} selected", - c.articles, c.eligible, c.triaged, c.assessed, c.shortlisted, c.selected + "curation: {} considered → {} eligible → {} triaged{} → {} assessed{} → {} shortlisted → {} selected", + c.articles, + c.eligible, + c.triaged, + rejected(c.triage_rejected), + c.assessed, + rejected(c.deep_rejected), + c.shortlisted, + c.selected ), format!( "admission: triage {} · interest {} · knn {} · exploration {} · blend {} · auto {}", @@ -435,6 +463,21 @@ mod tests { ); assert_eq!(RunReport::format_duration(48), "48s"); assert_eq!(RunReport::format_duration(3600), "60m00s"); + + // Provider rejections show up only when there were any. + r.counts.triaged = 395; + r.counts.triage_rejected = 3; + r.counts.deep_rejected = 0; + assert_eq!( + r.info_block()[0], + "curation: 412 considered → 398 eligible → 395 triaged · 3 rejected → 120 assessed → 60 shortlisted → 17 selected" + ); + r.counts.assessed = 119; + r.counts.deep_rejected = 1; + assert_eq!( + r.info_block()[0], + "curation: 412 considered → 398 eligible → 395 triaged · 3 rejected → 119 assessed · 1 rejected → 60 shortlisted → 17 selected" + ); } #[test]