1use anyhow::{Context, Result};
21use recall_wire::evaluations::{
22 self, EvaluationRequest, GLOBAL_PREFIX, STATE_DONE, STATE_FAILED, STATE_QUEUED, STATE_RUNNING,
23};
24use recall_wire::jobs::{KIND_EVALUATE, STATE_LEASED};
25use recall_wire::{EvaluateFile, EvaluateInput, Evaluation, EvaluationSummary, Finding};
26use rusqlite::{Connection, OptionalExtension, Row};
27use time::OffsetDateTime;
28
29use super::{Outcome, Store};
30use crate::format_timestamp;
31
32pub(super) const SCHEMA: &str = "
34 CREATE TABLE IF NOT EXISTS evaluations (
35 id TEXT PRIMARY KEY,
36 -- 'queued' until its report is in, then 'done'. While not done,
37 -- the state reported is its job's.
38 state TEXT NOT NULL CHECK (state IN ('queued', 'done')),
39 job_id TEXT NOT NULL,
40 -- What was asked: a JSON list of project keys, empty for every
41 -- project, and whether to run the contradiction check.
42 projects TEXT NOT NULL,
43 contradictions INTEGER NOT NULL DEFAULT 0,
44 -- The report: findings as JSON, never note text, and details as
45 -- the worker sent them (plain JSON; sealed once encryption lands).
46 findings TEXT,
47 details TEXT,
48 created_at TEXT NOT NULL,
49 finished_at TEXT
50 );
51 CREATE INDEX IF NOT EXISTS evaluations_by_created ON evaluations (created_at);
52";
53
54#[derive(serde::Serialize, serde::Deserialize)]
57struct Payload {
58 evaluation_id: String,
59 projects: Vec<String>,
60 contradictions: bool,
61}
62
63#[derive(Debug, Clone, PartialEq, Eq)]
65pub enum Requested {
66 Queued(String),
68 Busy(String),
70 Full,
72}
73
74struct EvalRow {
76 id: String,
77 state: String,
78 projects: String,
79 contradictions: bool,
80 findings: Option<String>,
81 details: Option<String>,
82 created_at: String,
83 finished_at: Option<String>,
84 job_state: Option<String>,
85 job_error: Option<String>,
86 job_updated_at: Option<String>,
87}
88
89const EVAL_COLUMNS: &str = "e.id, e.state, e.projects, e.contradictions, e.findings, e.details, \
90 e.created_at, e.finished_at, j.state, j.error, j.updated_at";
91
92fn eval_from(r: &Row<'_>) -> rusqlite::Result<EvalRow> {
93 Ok(EvalRow {
94 id: r.get(0)?,
95 state: r.get(1)?,
96 projects: r.get(2)?,
97 contradictions: r.get::<_, i64>(3)? != 0,
98 findings: r.get(4)?,
99 details: r.get(5)?,
100 created_at: r.get(6)?,
101 finished_at: r.get(7)?,
102 job_state: r.get(8)?,
103 job_error: r.get(9)?,
104 job_updated_at: r.get(10)?,
105 })
106}
107
108impl EvalRow {
109 fn state(&self) -> &'static str {
111 if self.state == STATE_DONE {
112 return STATE_DONE;
113 }
114 match self.job_state.as_deref() {
115 Some(STATE_LEASED) => STATE_RUNNING,
116 Some("failed") => STATE_FAILED,
117 _ => STATE_QUEUED,
118 }
119 }
120
121 fn finished_at(&self) -> Option<String> {
122 match self.state() {
123 STATE_DONE => self.finished_at.clone(),
124 STATE_FAILED => self.job_updated_at.clone(),
125 _ => None,
126 }
127 }
128
129 fn error(&self) -> Option<String> {
130 if self.state == STATE_DONE {
131 return None;
132 }
133 self.job_error.clone()
134 }
135
136 fn projects(&self) -> Vec<String> {
137 serde_json::from_str(&self.projects).unwrap_or_default()
138 }
139
140 fn findings(&self) -> Result<Vec<Finding>> {
141 match &self.findings {
142 Some(text) => serde_json::from_str(text)
143 .with_context(|| format!("evaluation {} has findings that do not read", self.id)),
144 None => Ok(Vec::new()),
145 }
146 }
147
148 fn summary(&self) -> Result<EvaluationSummary> {
149 let mut counts = std::collections::BTreeMap::new();
150 for f in self.findings()? {
151 *counts.entry(f.kind).or_insert(0) += 1;
152 }
153 Ok(EvaluationSummary {
154 id: self.id.clone(),
155 state: self.state().to_string(),
156 created_at: self.created_at.clone(),
157 finished_at: self.finished_at(),
158 counts,
159 projects: self.projects(),
160 contradictions: self.contradictions,
161 error: self.error(),
162 })
163 }
164}
165
166fn ts(at: OffsetDateTime) -> String {
167 format_timestamp(at)
168}
169
170pub(super) fn evaluate_input(conn: &Connection, payload: &str) -> Result<EvaluateInput> {
175 let asked: Payload = serde_json::from_str(payload)
176 .context("an evaluate job has a payload that does not read")?;
177 let mut stmt = conn.prepare(
178 "SELECT project_key, file_path, content, updated_at FROM memory_files
179 WHERE deleted = 0 ORDER BY project_key, file_path",
180 )?;
181 let rows = stmt.query_map([], |r| {
182 Ok(EvaluateFile {
183 project_key: r.get(0)?,
184 file_path: r.get(1)?,
185 content: r.get(2)?,
186 updated_at: r.get(3)?,
187 })
188 })?;
189 let mut files = Vec::new();
190 for row in rows {
191 let file = row?;
192 if asked.projects.is_empty()
193 || file.project_key.starts_with(GLOBAL_PREFIX)
194 || asked.projects.contains(&file.project_key)
195 {
196 files.push(file);
197 }
198 }
199 Ok(EvaluateInput {
200 evaluation_id: asked.evaluation_id,
201 projects: asked.projects,
202 contradictions: asked.contradictions,
203 files,
204 })
205}
206
207pub(super) fn holds(conn: &Connection, project_key: &str, file_path: &str) -> Result<bool> {
210 Ok(conn
211 .query_row(
212 "SELECT 1 FROM memory_files WHERE project_key = ?1 AND file_path = ?2",
213 (project_key, file_path),
214 |_| Ok(()),
215 )
216 .optional()?
217 .is_some())
218}
219
220pub(super) fn record_report(
223 conn: &Connection,
224 job_id: &str,
225 payload: &str,
226 findings: &[Finding],
227 details: &serde_json::Value,
228 now: OffsetDateTime,
229) -> Result<()> {
230 let asked: Payload = serde_json::from_str(payload)
231 .context("an evaluate job has a payload that does not read")?;
232 let changed = conn.execute(
233 "UPDATE evaluations SET state = 'done', findings = ?3, details = ?4, finished_at = ?5
234 WHERE id = ?1 AND job_id = ?2",
235 (
236 &asked.evaluation_id,
237 job_id,
238 serde_json::to_string(findings)?,
239 serde_json::to_string(details)?,
240 ts(now),
241 ),
242 )?;
243 anyhow::ensure!(
244 changed == 1,
245 "job {job_id} is for evaluation {}, which is not there",
246 asked.evaluation_id
247 );
248 Ok(())
249}
250
251impl Store {
252 pub fn has_project(&self, project_key: &str) -> Result<bool> {
255 Ok(self
256 .lock()
257 .query_row(
258 "SELECT 1 FROM memory_files WHERE project_key = ?1 LIMIT 1",
259 (project_key,),
260 |_| Ok(()),
261 )
262 .optional()?
263 .is_some())
264 }
265
266 pub fn request_evaluation_audited(
272 &self,
273 id: &str,
274 job_id: &str,
275 req: &EvaluationRequest,
276 now: OffsetDateTime,
277 build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
278 ) -> Result<Requested> {
279 self.audited(
280 |tx, _| {
281 let busy: Option<String> = tx
284 .query_row(
285 "SELECT e.id FROM evaluations e JOIN jobs j ON j.id = e.job_id
286 WHERE j.state IN ('queued', 'leased') LIMIT 1",
287 [],
288 |r| r.get(0),
289 )
290 .optional()?;
291 if let Some(open) = busy {
292 return Ok(Outcome::Refuse(Requested::Busy(open)));
293 }
294 let open: i64 = tx.query_row(
295 "SELECT COUNT(*) FROM jobs WHERE state IN ('queued', 'leased')",
296 [],
297 |r| r.get(0),
298 )?;
299 if open as usize >= super::MAX_OPEN_JOBS {
300 return Ok(Outcome::Refuse(Requested::Full));
301 }
302 let now = ts(now);
303 let projects = serde_json::to_string(&req.projects)?;
304 let payload = serde_json::to_string(&Payload {
305 evaluation_id: id.to_string(),
306 projects: req.projects.clone(),
307 contradictions: req.contradictions,
308 })?;
309 tx.execute(
312 "INSERT INTO jobs (id, kind, state, project_key, file_path, payload, not_before,
313 created_at, updated_at)
314 VALUES (?1, ?2, 'queued', '', '', ?3, ?4, ?4, ?4)",
315 (job_id, KIND_EVALUATE, &payload, &now),
316 )?;
317 tx.execute(
318 "INSERT INTO evaluations (id, state, job_id, projects, contradictions, created_at)
319 VALUES (?1, 'queued', ?2, ?3, ?4, ?5)",
320 (id, job_id, &projects, req.contradictions as i64, &now),
321 )?;
322 Ok(Outcome::Commit(Requested::Queued(job_id.to_string())))
323 },
324 |seq, at, _| build_leaf(seq, at),
325 )
326 }
327
328 pub fn evaluations(&self, limit: usize) -> Result<Vec<EvaluationSummary>> {
330 let conn = self.lock();
331 let mut stmt = conn.prepare(&format!(
332 "SELECT {EVAL_COLUMNS} FROM evaluations e LEFT JOIN jobs j ON j.id = e.job_id
333 ORDER BY e.created_at DESC, e.id DESC LIMIT ?1"
334 ))?;
335 let rows = stmt.query_map((limit as i64,), eval_from)?;
336 let mut out = Vec::new();
337 for row in rows {
338 out.push(row?.summary()?);
339 }
340 Ok(out)
341 }
342
343 pub fn evaluation(&self, id: &str, with_details: bool) -> Result<Option<Evaluation>> {
345 let conn = self.lock();
346 let row = conn
347 .query_row(
348 &format!(
349 "SELECT {EVAL_COLUMNS} FROM evaluations e LEFT JOIN jobs j ON j.id = e.job_id
350 WHERE e.id = ?1"
351 ),
352 (id,),
353 eval_from,
354 )
355 .optional()?;
356 let Some(row) = row else {
357 return Ok(None);
358 };
359 let details = match (&row.details, with_details) {
360 (Some(text), true) => Some(
361 serde_json::from_str(text)
362 .with_context(|| format!("evaluation {id} has details that do not read"))?,
363 ),
364 _ => None,
365 };
366 Ok(Some(Evaluation {
367 id: row.id.clone(),
368 state: row.state().to_string(),
369 created_at: row.created_at.clone(),
370 finished_at: row.finished_at(),
371 findings: row.findings()?,
372 details,
373 projects: row.projects(),
374 contradictions: row.contradictions,
375 error: row.error(),
376 }))
377 }
378
379 pub fn evaluation_open(&self) -> Result<bool> {
382 Ok(self
383 .lock()
384 .query_row(
385 "SELECT 1 FROM jobs WHERE kind = ?1 AND state IN ('queued', 'leased') LIMIT 1",
386 (KIND_EVALUATE,),
387 |_| Ok(()),
388 )
389 .optional()?
390 .is_some())
391 }
392
393 pub fn prune_evaluations(&self, before: &str) -> Result<usize> {
398 Ok(self.lock().execute(
399 "DELETE FROM evaluations WHERE state = 'done' AND finished_at < ?1",
400 (before,),
401 )?)
402 }
403}
404
405pub(super) fn check_files(conn: &Connection, findings: &[Finding]) -> Result<Option<String>> {
410 let mut ids = std::collections::HashSet::new();
411 for f in findings {
412 if !ids.insert(f.id.as_str()) {
413 return Ok(Some(format!("two findings have the id {}", f.id)));
414 }
415 let named = std::iter::once((&f.project_key, &f.file_path))
416 .chain(f.related.iter().map(|r| (&r.project_key, &r.file_path)));
417 for (project_key, file_path) in named {
418 if !holds(conn, project_key, file_path)? {
419 return Ok(Some(format!(
420 "finding {} names a file this server does not hold",
421 f.id
422 )));
423 }
424 }
425 }
426 if findings.len() > evaluations::MAX_FINDINGS {
427 return Ok(Some(format!(
428 "a report may carry at most {} findings",
429 evaluations::MAX_FINDINGS
430 )));
431 }
432 Ok(None)
433}