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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01A1rCLQeKBgnBo3oTgHuTMe
This commit is contained in:
2026-09-02 20:44:28 +00:00
co-authored by Claude Fable 5.1
parent 7a842b61b6
commit af2a3ceb49
14 changed files with 1711 additions and 127 deletions
+45
View File
@@ -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);
}
}
+248 -44
View File
@@ -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<bool> {
})
}
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<usize> {
) -> anyhow::Result<StageSummary> {
let positions = candidates
.iter()
.enumerate()
.filter(|(_, candidate)| candidate.stage == "admitted")
.map(|(index, candidate)| (candidate.article.id, index))
.collect::<HashMap<_, _>>();
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::<Option<String>, _>("kind").as_deref() == Some(PROVIDER_REJECTED) {
known_rejected.insert(id);
continue;
}
let (Some(quality), Some(fit)) = (
row.get::<Option<f64>, _>("score"),
row.get::<Option<f64>, _>("fit"),
@@ -376,52 +404,47 @@ pub async fn run(
.unwrap_or_default(),
paywalled_guess: row.get::<i64, _>("paywalled_guess") != 0,
facets,
model: model.to_string(),
model: row.get::<String, _>("model"),
prompt_version: DEEP_PROMPT_VERSION,
assessed_at: parse_ts(
"article_assessments.assessed_at",
&row.get::<String, _>("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::<Vec<_>>();
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::<HashSet<_>>();
(allowed, build_batch_prompt(batch, sections))
})
.map(<[&Candidate]>::to_vec)
.collect::<Vec<_>>();
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::<Vec<Vec<DeepItem>>>()
.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 {
&sections(),
)
.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<dyn ChatBackend>) -> 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<Candidate> {
(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,
&sections(),
)
.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::<String, _>("model"), "deepseek-v4-flash");
assert_eq!(row.get::<Option<f64>, _>("score"), None);
assert_eq!(row.get::<Option<f64>, _>("fit"), None);
assert_eq!(
row.get::<Option<String>, _>("kind").as_deref(),
Some(PROVIDER_REJECTED)
);
assert!(
row.get::<Option<String>, _>("rationale")
.is_some_and(|why| why.starts_with("deepseek: 400 Bad Request"))
);
assert_eq!(row.get::<Option<String>, _>("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);
}
}
+839
View File
@@ -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<T> + Sync),
}
impl<T> 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<String> {
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<T> {
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 {
/// `<provider>: <message>`, 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<T> {
pub items: Vec<Scored<T>>,
/// Articles nobody would assess; the caller persists these.
pub rejections: Vec<Rejection>,
/// 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<T> Default for BatchOutcome<T> {
fn default() -> Self {
Self {
items: Vec::new(),
rejections: Vec::new(),
rejected: 0,
recovered: 0,
requests: 0,
budget_stopped: false,
}
}
}
impl<T> BatchOutcome<T> {
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<String>,
}
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 `<provider> 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::<Vec<_>>().join(" ");
if collapsed.chars().count() <= max_chars {
return collapsed;
}
let mut cut = collapsed
.chars()
.take(max_chars.saturating_sub(1))
.collect::<String>();
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<T: Assessed>(
runner: &BatchRunner<'_, T>,
batches: Vec<Vec<&Candidate>>,
max_concurrent: usize,
) -> BatchOutcome<T> {
let outcomes = stream::iter(batches)
.map(|batch| run_one(runner, batch))
.buffer_unordered(max_concurrent.max(1))
.collect::<Vec<_>>()
.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<T: Assessed>(
runner: &BatchRunner<'_, T>,
batch: Vec<&Candidate>,
) -> BatchOutcome<T> {
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::<HashSet<_>>();
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::<Vec<_>>();
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<Work<'c>>, 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<T: Assessed>(
runner: &BatchRunner<'_, T>,
candidate: &Candidate,
prompt: &str,
message: String,
out: &mut BatchOutcome<T>,
) {
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<String>,
/// Prompts answered with `{"articles": []}` (by request index).
empty_on: Vec<usize>,
answer: fn(&[i64]) -> String,
seen: Mutex<Vec<String>>,
}
impl FilterBackend {
/// Answers in the triage shape.
pub(crate) fn new(forbidden: &[&str], empty_on: &[usize]) -> Arc<Self> {
Self::with_answer(forbidden, empty_on, answer_for)
}
/// Answers in the deep-assessment shape.
pub(crate) fn deep(forbidden: &[&str]) -> Arc<Self> {
Self::with_answer(forbidden, &[], deep_answer_for)
}
fn with_answer(
forbidden: &[&str],
empty_on: &[usize],
answer: fn(&[i64]) -> String,
) -> Arc<Self> {
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<Vec<i64>> {
self.seen
.lock()
.map(|seen| seen.iter().map(|prompt| ids_in(prompt)).collect())
.unwrap_or_default()
}
}
fn ids_in(prompt: &str) -> Vec<i64> {
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::<Vec<_>>()
.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::<Vec<_>>()
.join(",");
format!(r#"{{"articles":[{items}]}}"#)
}
impl ChatBackend for FilterBackend {
fn complete<'a>(
&'a self,
req: ChatRequest,
) -> Pin<Box<dyn Future<Output = Result<ChatCompletion, LlmError>> + 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<dyn ChatBackend>) -> 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<Candidate> {
(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<TriageItem> {
let runner = BatchRunner {
llm,
fallback,
temperature: 0.3,
build_prompt: &build_batch_prompt,
parse: &parse_triage_response,
};
let all = candidates.iter().collect::<Vec<_>>();
let batches = all.chunks(batch_size).map(<[&Candidate]>::to_vec).collect();
run_batches(&runner, batches, 4).await
}
fn sorted_ids(outcome: &BatchOutcome<TriageItem>) -> Vec<i64> {
let mut ids = outcome
.items
.iter()
.map(|scored| scored.item.id)
.collect::<Vec<_>>();
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::<Vec<_>>();
let backend =
FilterBackend::new(&titles.iter().map(String::as_str).collect::<Vec<_>>(), &[]);
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::<Vec<_>>();
rejected.sort_unstable();
assert_eq!(rejected, (1..=8).collect::<Vec<_>>());
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"
);
}
}
+2 -2
View File
@@ -416,8 +416,8 @@ pub fn decode_blob(bytes: &[u8], dimension: usize) -> Result<Vec<f32>, 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);
}
+1 -1
View File
@@ -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
}
+12 -4
View File
@@ -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<i64>,
assessed_at: Timestamp,
) -> anyhow::Result<usize> {
) -> anyhow::Result<batch::StageSummary> {
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,
+56
View File
@@ -429,6 +429,12 @@ pub async fn render_explain(db: &Db, row: &ExplainRow) -> Result<String, sqlx::E
let rationale = assessment
.get::<Option<String>, _>("rationale")
.unwrap_or_default();
if assessment.get::<Option<String>, _>("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;
+387 -59
View File
@@ -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<ArticleId>,
batch_size: usize,
@@ -350,36 +379,46 @@ pub async fn run(
profile_version: Option<i64>,
assessed_at: Timestamp,
temperature: f32,
) -> anyhow::Result<usize> {
) -> anyhow::Result<StageSummary> {
let positions = candidates
.iter()
.enumerate()
.map(|(index, candidate)| (candidate.article.id, index))
.collect::<HashMap<_, _>>();
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::<HashMap<_, _>>();
for row in rows {
let id = row.get::<i64, _>("article_id");
if !pool_ids.contains(&id) {
if !pool.contains(&id) {
continue;
}
let rejected =
row.get::<Option<String>, _>("kind").as_deref() == Some(PROVIDER_REJECTED);
if row.get::<String, _>("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::<Option<f64>, _>("score") else {
@@ -398,7 +437,7 @@ pub async fn run(
why: row
.get::<Option<String>, _>("rationale")
.unwrap_or_default(),
model: llm.model.clone(),
model: row.get::<String, _>("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::<Vec<_>>();
let prompts = pending
.chunks(batch_size.max(1))
.map(|batch| {
let allowed = batch
.iter()
.map(|candidate| candidate.article.id)
.collect::<HashSet<_>>();
(allowed, build_batch_prompt(batch))
})
.collect::<Vec<_>>();
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::<Vec<Vec<TriageItem>>>()
.await;
let positions = candidates
let pending = in_pool
.iter()
.enumerate()
.map(|(index, candidate)| (candidate.article.id, index))
.collect::<HashMap<_, _>>();
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::<Vec<_>>();
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::<Vec<_>>();
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<i64>,
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<dyn ChatBackend>) -> 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<Candidate>, HashSet<ArticleId>) {
let candidates = (1..=n)
.map(|id| Candidate::new(article(id, &format!("Piece {id}"), 800), false))
.collect::<Vec<_>>();
(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<ArticleId>,
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::<i64, _>("article_id"),
row.get::<String, _>("model"),
row.get::<Option<String>, _>("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,
+46 -13
View File
@@ -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
+45 -2
View File
@@ -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]