Skip to main content

recall_server/store/
jobs.rs

1//! The merge queue: jobs, their leases, and applying their results.
2//!
3//! Every function that decides something takes `now`, so tests can move
4//! the clock rather than wait for it, and every one that changes more than
5//! one row does it in one transaction under the store's lock: a job is
6//! never claimed twice, and a merged file is never written without its job
7//! recording it.
8
9use std::time::Duration;
10
11use anyhow::{Context, Result};
12use recall_wire::jobs::{
13    KIND_EVALUATE, KIND_MERGE, MAX_LEASE_SECONDS, STATE_DONE, STATE_FAILED, STATE_LEASED,
14    STATE_QUEUED,
15};
16use recall_wire::{
17    content_sha256, Job, JobSummary, MergeInput, MergeSide, QueueStatus, ResultRequest,
18    ResultResponse,
19};
20use rusqlite::{Connection, OptionalExtension, Row};
21use time::OffsetDateTime;
22
23use super::evaluations;
24use super::{read_file, write_file, Outcome, Store};
25use crate::audit::leaf;
26use crate::format_timestamp;
27
28/// Created with the other tables, every time the store opens.
29///
30/// `kind` has no `CHECK`: later kinds (sealing, evaluating) arrive without
31/// rebuilding the table, which is what `devices.scope` needed.
32pub(super) const SCHEMA: &str = "
33    CREATE TABLE IF NOT EXISTS jobs (
34        id               TEXT PRIMARY KEY,
35        kind             TEXT NOT NULL,
36        state            TEXT NOT NULL
37                         CHECK (state IN ('queued', 'leased', 'done', 'failed')),
38        project_key      TEXT NOT NULL,
39        file_path        TEXT NOT NULL,
40        -- The job's input as JSON: for a merge, both versions. Emptied to
41        -- '{}' once the job is done, since the file itself then holds what
42        -- mattered.
43        payload          TEXT NOT NULL,
44        -- The current lease while leased; afterwards the lease its result
45        -- came under, which is what makes posting that result again a
46        -- no-op rather than a 409. NULL once a lease ran out unanswered.
47        lease_id         TEXT,
48        lease_expires_at TEXT,
49        -- How many times it has been handed out.
50        attempt          INTEGER NOT NULL DEFAULT 0,
51        -- Not handed out before this: the retry delay after a failure.
52        not_before       TEXT NOT NULL,
53        -- The job whose result this one merges again, and how far down
54        -- that chain it is, from 0.
55        parent_id        TEXT,
56        link             INTEGER NOT NULL DEFAULT 0,
57        error            TEXT,
58        -- A result that was not applied, kept so it can be recovered: the
59        -- merged version, as JSON.
60        result           TEXT,
61        applied          INTEGER NOT NULL DEFAULT 0,
62        follow_up        TEXT,
63        created_at       TEXT NOT NULL,
64        updated_at       TEXT NOT NULL
65    );
66    CREATE INDEX IF NOT EXISTS jobs_by_state ON jobs (state, not_before, created_at);
67";
68
69/// How many jobs may wait or be held at once. A queue a worker is not
70/// draining must not grow without end; past this, a stale push is stored
71/// as last-write-wins, as a failed merge always was, and says so in
72/// `/health`.
73pub const MAX_OPEN_JOBS: usize = 1000;
74
75/// How many times a job is handed out before it is marked failed: once,
76/// then again 1, 5 and 30 minutes after each failed attempt.
77pub const MAX_ATTEMPTS: u32 = 4;
78
79/// How long a job waits after its first, second and third failed attempt.
80const RETRY_AFTER: [Duration; 3] = [
81    Duration::from_secs(60),
82    Duration::from_secs(5 * 60),
83    Duration::from_secs(30 * 60),
84];
85
86/// How many follow-ups one conflict may lead to. A file that changes
87/// faster than it can be merged stops being chased here: the newest push
88/// stands, and the last result waits in a failed job.
89pub const MAX_LINKS: u32 = 3;
90
91/// The most of an error a job keeps, in bytes. The error is the worker's to
92/// write and a result may be megabytes; what went wrong fits in far less,
93/// and the listing and the server's log need no more.
94pub const MAX_ERROR_BYTES: usize = 500;
95
96/// Why a job whose merge came back empty is retried rather than applied.
97const EMPTY_RESULT: &str = "the merge came back empty from two versions that were not";
98
99/// `s`, cut to at most `max` bytes, on a character boundary.
100pub fn clip(s: &str, max: usize) -> &str {
101    if s.len() <= max {
102        return s;
103    }
104    let mut end = max;
105    while !s.is_char_boundary(end) {
106        end -= 1;
107    }
108    &s[..end]
109}
110
111const JOB_COLUMNS: &str = "id, kind, state, project_key, file_path, payload, lease_id, \
112     lease_expires_at, attempt, parent_id, link, error, result, applied, follow_up, \
113     created_at, updated_at";
114
115/// A job row.
116struct JobRow {
117    id: String,
118    kind: String,
119    state: String,
120    project_key: String,
121    file_path: String,
122    payload: String,
123    lease_id: Option<String>,
124    lease_expires_at: Option<String>,
125    attempt: u32,
126    link: u32,
127    error: Option<String>,
128    result: Option<String>,
129    applied: bool,
130    follow_up: Option<String>,
131    created_at: String,
132    updated_at: String,
133}
134
135fn row_from(r: &Row<'_>) -> rusqlite::Result<JobRow> {
136    Ok(JobRow {
137        id: r.get(0)?,
138        kind: r.get(1)?,
139        state: r.get(2)?,
140        project_key: r.get(3)?,
141        file_path: r.get(4)?,
142        payload: r.get(5)?,
143        lease_id: r.get(6)?,
144        lease_expires_at: r.get(7)?,
145        attempt: r.get(8)?,
146        link: r.get(10)?,
147        error: r.get(11)?,
148        result: r.get(12)?,
149        applied: r.get::<_, i64>(13)? != 0,
150        follow_up: r.get(14)?,
151        created_at: r.get(15)?,
152        updated_at: r.get(16)?,
153    })
154}
155
156impl JobRow {
157    fn summary(&self) -> JobSummary {
158        JobSummary {
159            id: self.id.clone(),
160            kind: self.kind.clone(),
161            state: self.state.clone(),
162            project_key: self.project_key.clone(),
163            file_path: self.file_path.clone(),
164            attempt: self.attempt,
165            created_at: self.created_at.clone(),
166            updated_at: self.updated_at.clone(),
167            error: self.error.clone(),
168            follow_up: self.follow_up.clone(),
169        }
170    }
171
172    fn outcome(&self) -> ResultResponse {
173        ResultResponse {
174            id: self.id.clone(),
175            state: self.state.clone(),
176            applied: self.applied,
177            follow_up: self.follow_up.clone(),
178        }
179    }
180
181    fn merge_input(&self) -> Result<MergeInput> {
182        serde_json::from_str(&self.payload)
183            .with_context(|| format!("job {} has a payload that does not read", self.id))
184    }
185}
186
187fn get_job(conn: &Connection, id: &str) -> Result<Option<JobRow>> {
188    Ok(conn
189        .query_row(
190            &format!("SELECT {JOB_COLUMNS} FROM jobs WHERE id = ?1"),
191            (id,),
192            row_from,
193        )
194        .optional()?)
195}
196
197fn ts(at: OffsetDateTime) -> String {
198    format_timestamp(at)
199}
200
201/// The version of a file a row holds, as a job carries it.
202fn side_of(content: &str, source_env: &str, updated_at: &str) -> MergeSide {
203    MergeSide {
204        sha256: content_sha256(content),
205        content: content.to_string(),
206        source_env: source_env.to_string(),
207        updated_at: updated_at.to_string(),
208    }
209}
210
211/// A job that ran out of attempts, or whose result could not be applied.
212///
213/// It is said two ways, because `/health` answers anyone:
214/// [`Failure::public`] names the job and nothing else, never a project or a
215/// path, and [`Failure::logged`], for the server's own log, names the file
216/// and the error as well.
217#[derive(Debug, Clone, PartialEq, Eq)]
218pub struct Failure {
219    /// The job.
220    pub id: String,
221    /// What kind of job: only a merge's failure is `/health`'s business.
222    pub kind: String,
223    /// The project its file belongs to.
224    pub project_key: String,
225    /// The file.
226    pub file_path: String,
227    /// What became of it, naming no project and no file, such as `failed
228    /// after 4 attempts`.
229    pub what: String,
230    /// The last error it met, at most [`MAX_ERROR_BYTES`]; empty for none.
231    pub error: String,
232}
233
234impl Failure {
235    /// What `/health` shows as the last merge error: the job, what became
236    /// of it, and where to look.
237    pub fn public(&self) -> String {
238        format!(
239            "{} job {} {}; see GET /v1/jobs?state=failed",
240            self.kind, self.id, self.what
241        )
242    }
243
244    /// What the server's log says: the file and the error too, for a
245    /// merge. For any other job, only what became of it: an evaluation's
246    /// error is the worker's own text, which the log is not where to keep.
247    pub fn logged(&self) -> String {
248        if self.kind != KIND_MERGE {
249            return format!("{} job {} {}", self.kind, self.id, self.what);
250        }
251        let mut line = format!(
252            "merge job {} for {}/{} {}",
253            self.id, self.project_key, self.file_path, self.what
254        );
255        if !self.error.is_empty() {
256            line.push_str(": ");
257            line.push_str(&self.error);
258        }
259        line
260    }
261}
262
263/// What queueing a merge came to.
264#[derive(Debug, Clone, PartialEq, Eq)]
265pub enum Queued {
266    /// The incoming version is stored, and this job will merge it with the
267    /// one it displaced.
268    Queued(String),
269    /// The incoming version is stored, and there was nothing to merge it
270    /// with: the file was new, a tombstone, or already this content.
271    Nothing,
272    /// The incoming version is stored, and the queue was full, so no job
273    /// was made: last-write-wins, as a failed merge.
274    Full,
275}
276
277/// What a result came to, when it was recorded.
278#[derive(Debug, Clone, PartialEq, Eq)]
279pub struct Settled {
280    /// The answer for the worker.
281    pub response: ResultResponse,
282    /// Whether this call wrote the merged file.
283    pub applied: bool,
284    /// Whether this call queued something a claim may now take: a
285    /// follow-up.
286    pub queued: bool,
287    /// When this call marked the job failed, the failure: what `/health`
288    /// shows as the last merge error.
289    pub failed: Option<Box<Failure>>,
290}
291
292/// What posting a result came to.
293#[derive(Debug, Clone, PartialEq, Eq)]
294pub enum Settlement {
295    /// Recorded, or recorded already under this lease and left alone.
296    Recorded(Settled),
297    /// No job has that id.
298    NotFound,
299    /// The lease is not the job's current one, or has run out: someone
300    /// else has, or will have, the job.
301    LeaseEnded,
302    /// A result of one kind for a job of another: the kind the result
303    /// was.
304    WrongKind(&'static str),
305    /// An evaluation's report that cannot be recorded, and why: a finding
306    /// names a file the store does not hold, two share an id, or there
307    /// are too many. Nothing was changed.
308    Invalid(String),
309}
310
311/// What retrying a job came to.
312#[derive(Debug, Clone, PartialEq, Eq)]
313pub enum Retried {
314    /// Queued again.
315    Queued(JobSummary),
316    /// No job has that id.
317    NotFound,
318    /// It has not failed, so there is nothing to retry.
319    NotFailed(String),
320}
321
322/// After a failed attempt: back in the queue after the delay for that
323/// attempt, or failed for good. Answers the failure when it is final.
324/// `error` is kept to [`MAX_ERROR_BYTES`].
325fn after_failure(
326    conn: &Connection,
327    job: &JobRow,
328    error: &str,
329    now: OffsetDateTime,
330    keep_lease: bool,
331) -> Result<Option<Failure>> {
332    let error = clip(error, MAX_ERROR_BYTES);
333    let lease = if keep_lease {
334        job.lease_id.clone()
335    } else {
336        None
337    };
338    let retry = (job.attempt as usize)
339        .checked_sub(1)
340        .and_then(|i| RETRY_AFTER.get(i))
341        .filter(|_| job.attempt < MAX_ATTEMPTS);
342    match retry {
343        Some(delay) => {
344            conn.execute(
345                "UPDATE jobs SET state = 'queued', lease_id = ?2, lease_expires_at = NULL,
346                     not_before = ?3, error = ?4, applied = 0, updated_at = ?5
347                 WHERE id = ?1",
348                (&job.id, &lease, ts(now + *delay), error, ts(now)),
349            )?;
350            Ok(None)
351        }
352        None => {
353            let kept = format!("{error} (gave up after {} attempts)", job.attempt);
354            conn.execute(
355                "UPDATE jobs SET state = 'failed', lease_id = ?2, lease_expires_at = NULL,
356                     error = ?3, applied = 0, updated_at = ?4
357                 WHERE id = ?1",
358                (&job.id, &lease, &kept, ts(now)),
359            )?;
360            Ok(Some(Failure {
361                id: job.id.clone(),
362                kind: job.kind.clone(),
363                project_key: job.project_key.clone(),
364                file_path: job.file_path.clone(),
365                what: format!("failed after {} attempts", job.attempt),
366                error: error.to_string(),
367            }))
368        }
369    }
370}
371
372fn insert_merge_job(
373    conn: &Connection,
374    id: &str,
375    input: &MergeInput,
376    parent: Option<(&str, u32)>,
377    now: OffsetDateTime,
378) -> Result<()> {
379    let now = ts(now);
380    conn.execute(
381        "INSERT INTO jobs (id, kind, state, project_key, file_path, payload, not_before,
382                           parent_id, link, created_at, updated_at)
383         VALUES (?1, ?2, 'queued', ?3, ?4, ?5, ?6, ?7, ?8, ?6, ?6)",
384        (
385            id,
386            KIND_MERGE,
387            &input.project_key,
388            &input.file_path,
389            serde_json::to_string(input)?,
390            &now,
391            parent.map(|(p, _)| p),
392            parent.map_or(0, |(_, link)| link),
393        ),
394    )?;
395    Ok(())
396}
397
398impl Store {
399    /// Stores `incoming` as the file's content, last-write-wins, and queues
400    /// a merge of it with whatever it displaced, in one transaction with
401    /// the push's leaf, which `build_leaf` makes knowing what was queued:
402    /// the job's `stored` side is read under the same lock as the write, so
403    /// it is exactly the version this push replaced. Answers what was
404    /// queued and the write's `updated_at`, which is its leaf's `at`, as
405    /// every audited write's is; `incoming.updated_at` is not used.
406    pub fn write_and_queue_merge_audited(
407        &self,
408        project_key: &str,
409        file_path: &str,
410        incoming: &MergeSide,
411        job_id: &str,
412        now: OffsetDateTime,
413        build_leaf: impl FnOnce(u64, &str, &Queued) -> Vec<u8>,
414    ) -> Result<(Queued, String)> {
415        let mut stamped = None;
416        let queued = self.audited(
417            |tx, at| {
418                stamped = Some(at.to_string());
419                let incoming = MergeSide {
420                    updated_at: at.to_string(),
421                    ..incoming.clone()
422                };
423                let displaced = read_file(tx, project_key, file_path)?
424                    .filter(|e| !e.deleted && e.content != incoming.content);
425                write_file(
426                    tx,
427                    project_key,
428                    file_path,
429                    &incoming.content,
430                    &incoming.source_env,
431                    at,
432                )?;
433                let Some(displaced) = displaced else {
434                    return Ok(Outcome::Commit(Queued::Nothing));
435                };
436                let open: i64 = tx.query_row(
437                    "SELECT COUNT(*) FROM jobs WHERE state IN ('queued', 'leased')",
438                    [],
439                    |r| r.get(0),
440                )?;
441                if open as usize >= MAX_OPEN_JOBS {
442                    return Ok(Outcome::Commit(Queued::Full));
443                }
444                let input = MergeInput {
445                    project_key: project_key.to_string(),
446                    file_path: file_path.to_string(),
447                    stored: side_of(
448                        &displaced.content,
449                        &displaced.source_env,
450                        &displaced.updated_at,
451                    ),
452                    incoming,
453                };
454                insert_merge_job(tx, job_id, &input, None, now)?;
455                Ok(Outcome::Commit(Queued::Queued(job_id.to_string())))
456            },
457            build_leaf,
458        )?;
459        Ok((queued, stamped.context("the write ran")?))
460    }
461
462    /// Puts every job whose lease ran out back in the queue, or fails it
463    /// when that was its last attempt. Answers each job it failed.
464    ///
465    /// Each job whose lease ran out gets a `job_result` leaf, which
466    /// `build_leaf` makes, in the same transaction: the attempt ended with
467    /// no result, and the job waits for another or has failed. A call that
468    /// finds no lease run out changes nothing and appends nothing.
469    pub fn expire_leases_audited(
470        &self,
471        now: OffsetDateTime,
472        build_leaf: impl Fn(u64, &str, &leaf::JobChange<'_>) -> Vec<u8>,
473    ) -> Result<Vec<Failure>> {
474        let (failed, _) = self.audited_each(
475            |tx, _| {
476                let expired: Vec<JobRow> = {
477                    let mut stmt = tx.prepare(&format!(
478                        "SELECT {JOB_COLUMNS} FROM jobs WHERE state = 'leased' AND lease_expires_at <= ?1"
479                    ))?;
480                    let rows = stmt.query_map((ts(now),), row_from)?;
481                    rows.collect::<rusqlite::Result<_>>()?
482                };
483                if expired.is_empty() {
484                    return Ok(Outcome::Refuse((Vec::new(), Vec::new())));
485                }
486                let mut failed = Vec::new();
487                let mut ended = Vec::new();
488                for job in expired {
489                    let failure = after_failure(
490                        tx,
491                        &job,
492                        "the worker did not report back before its lease ended",
493                        now,
494                        false,
495                    )?;
496                    let state = if failure.is_some() {
497                        STATE_FAILED
498                    } else {
499                        STATE_QUEUED
500                    };
501                    failed.extend(failure);
502                    ended.push((job.id, job.project_key, job.file_path, state));
503                }
504                Ok(Outcome::Commit((failed, ended)))
505            },
506            |seq, at, (_, ended)| {
507                ended
508                    .iter()
509                    .zip(seq..)
510                    .map(|((job_id, project_key, file_path, state), seq)| {
511                        build_leaf(
512                            seq,
513                            at,
514                            &leaf::JobChange {
515                                job_id,
516                                project_key,
517                                file_path,
518                                state,
519                                stored_sha256: None,
520                                follow_up: None,
521                            },
522                        )
523                    })
524                    .collect()
525            },
526        )?;
527        Ok(failed)
528    }
529
530    /// Leases the oldest queued job of one of `kinds` that may run now,
531    /// for `lease` from `now`, under `lease_id`, and appends the
532    /// `job_claim` leaf `build_leaf` makes in the same transaction. Finding
533    /// nothing to lease changes nothing and appends nothing.
534    pub fn claim_job_audited(
535        &self,
536        kinds: &[String],
537        lease_id: &str,
538        lease: Duration,
539        now: OffsetDateTime,
540        build_leaf: impl FnOnce(u64, &str, &Job) -> Vec<u8>,
541    ) -> Result<Option<Job>> {
542        if kinds.is_empty() {
543            return Ok(None);
544        }
545        self.audited(
546            |tx, _| {
547                // Kinds come from the request, so they are bound, never
548                // spliced.
549                let marks = (0..kinds.len())
550                    .map(|i| format!("?{}", i + 2))
551                    .collect::<Vec<_>>()
552                    .join(", ");
553                let now_text = ts(now);
554                let mut params: Vec<&dyn rusqlite::ToSql> = vec![&now_text];
555                for kind in kinds {
556                    params.push(kind);
557                }
558                let job = tx
559                    .query_row(
560                        &format!(
561                            "SELECT {JOB_COLUMNS} FROM jobs
562                             WHERE state = 'queued' AND not_before <= ?1 AND kind IN ({marks})
563                             ORDER BY created_at, id LIMIT 1"
564                        ),
565                        params.as_slice(),
566                        row_from,
567                    )
568                    .optional()?;
569                let Some(job) = job else {
570                    return Ok(Outcome::Refuse(None));
571                };
572                // An evaluation reads every project and may ask claude once
573                // for each, so it gets the longest lease there is, whatever
574                // the claim asked for.
575                let lease = if job.kind == KIND_EVALUATE {
576                    lease.max(Duration::from_secs(MAX_LEASE_SECONDS))
577                } else {
578                    lease
579                };
580                let expires = ts(now + lease);
581                tx.execute(
582                    "UPDATE jobs SET state = 'leased', lease_id = ?2, lease_expires_at = ?3,
583                         attempt = attempt + 1, updated_at = ?4
584                     WHERE id = ?1",
585                    (&job.id, lease_id, &expires, &now_text),
586                )?;
587                let merge = match job.kind.as_str() {
588                    KIND_MERGE => Some(job.merge_input()?),
589                    _ => None,
590                };
591                let evaluate = match job.kind.as_str() {
592                    KIND_EVALUATE => Some(evaluations::evaluate_input(tx, &job.payload)?),
593                    _ => None,
594                };
595                Ok(Outcome::Commit(Some(Job {
596                    id: job.id,
597                    kind: job.kind,
598                    lease_id: lease_id.to_string(),
599                    lease_expires_at: expires,
600                    attempt: job.attempt + 1,
601                    merge,
602                    evaluate,
603                })))
604            },
605            |seq, at, job| build_leaf(seq, at, job.as_ref().expect("a leaf only for a lease")),
606        )
607    }
608
609    /// Records a worker's result for job `id`, with the `job_result` leaf
610    /// `build_leaf` makes from what it changed, in one transaction.
611    ///
612    /// A result counts only under the job's current, unexpired lease. The
613    /// same result posted again under the lease that settled it changes
614    /// nothing, answers as the first did, and appends nothing.
615    ///
616    /// A merge is applied as a compare-and-swap: only if the file still
617    /// has the hash of the version the job stored, and then attributed to
618    /// `worker`. If another push landed meanwhile, the merged content
619    /// becomes the `stored` side of a follow-up job, `follow_up_id`,
620    /// against the newer version, at most [`MAX_LINKS`] deep. An empty
621    /// merge of two versions that were not empty is never applied: it is
622    /// taken as an error, and retried.
623    pub fn settle_job_audited(
624        &self,
625        id: &str,
626        result: &ResultRequest,
627        worker: &str,
628        follow_up_id: &str,
629        now: OffsetDateTime,
630        build_leaf: impl FnOnce(u64, &str, &leaf::JobChange<'_>) -> Vec<u8>,
631    ) -> Result<Settlement> {
632        let (settlement, _) = self.audited(
633            |tx, _| {
634                let refuse = |s: Settlement| Ok(Outcome::Refuse((s, None)));
635                let Some(job) = get_job(tx, id)? else {
636                    return refuse(Settlement::NotFound);
637                };
638                if job.lease_id.as_deref() != Some(result.lease_id.as_str()) {
639                    return refuse(Settlement::LeaseEnded);
640                }
641                if job.state != STATE_LEASED {
642                    // Settled already, under this very lease: the repeat
643                    // of a result that was recorded.
644                    return refuse(Settlement::Recorded(Settled {
645                        response: job.outcome(),
646                        applied: false,
647                        queued: false,
648                        failed: None,
649                    }));
650                }
651                if job
652                    .lease_expires_at
653                    .as_deref()
654                    .is_none_or(|at| at <= ts(now).as_str())
655                {
656                    return refuse(Settlement::LeaseEnded);
657                }
658
659                let mut stored = None;
660                if result.members() != 1 {
661                    anyhow::bail!("a result carries exactly one of merge, evaluate and error");
662                }
663                let settled = match (&result.merge, &result.error) {
664                    (None, None) => {
665                        let report = result.evaluate.as_ref().context("checked above")?;
666                        if job.kind != KIND_EVALUATE {
667                            return refuse(Settlement::WrongKind(KIND_EVALUATE));
668                        }
669                        if let Some(why) = evaluations::check_files(tx, &report.findings)? {
670                            return refuse(Settlement::Invalid(why));
671                        }
672                        if !report.details.is_object() {
673                            return refuse(Settlement::Invalid(
674                                "details is not an object".to_string(),
675                            ));
676                        }
677                        evaluations::record_report(
678                            tx,
679                            &job.id,
680                            &job.payload,
681                            &report.findings,
682                            &report.details,
683                            now,
684                        )?;
685                        finish(tx, &job, STATE_DONE, false, None, None, None, now)?;
686                        Settled {
687                            response: ResultResponse::default(),
688                            applied: false,
689                            queued: false,
690                            failed: None,
691                        }
692                    }
693                    (_, Some(error)) => {
694                        let failed = after_failure(tx, &job, error, now, true)?.map(Box::new);
695                        Settled {
696                            response: ResultResponse::default(),
697                            applied: false,
698                            queued: false,
699                            failed,
700                        }
701                    }
702                    (Some(merged), None) => {
703                        if job.kind != KIND_MERGE {
704                            return refuse(Settlement::WrongKind(KIND_MERGE));
705                        }
706                        let input = job.merge_input()?;
707                        // Nothing from two versions that had something is
708                        // not a merge but a malfunction, as the inline merge
709                        // treats it: retried like an error, and never
710                        // written, where it would replace both machines'
711                        // notes with an empty file.
712                        let emptied = merged.content.trim().is_empty()
713                            && !(input.stored.content.trim().is_empty()
714                                && input.incoming.content.trim().is_empty());
715                        if emptied {
716                            Settled {
717                                response: ResultResponse::default(),
718                                applied: false,
719                                queued: false,
720                                failed: after_failure(tx, &job, EMPTY_RESULT, now, true)?
721                                    .map(Box::new),
722                            }
723                        } else {
724                            let settled = apply_merge(
725                                tx,
726                                &job,
727                                &input,
728                                &merged.content,
729                                worker,
730                                follow_up_id,
731                                now,
732                            )?;
733                            if settled.applied {
734                                stored = Some(content_sha256(&merged.content));
735                            }
736                            settled
737                        }
738                    }
739                };
740                let response = get_job(tx, id)?
741                    .context("the job was read above")?
742                    .outcome();
743                let change = (job.project_key, job.file_path, stored);
744                Ok(Outcome::Commit((
745                    Settlement::Recorded(Settled {
746                        response,
747                        ..settled
748                    }),
749                    Some(change),
750                )))
751            },
752            |seq, at, (settlement, change)| {
753                let (Settlement::Recorded(s), Some((project_key, file_path, stored))) =
754                    (settlement, change)
755                else {
756                    unreachable!("a leaf only for a recorded result");
757                };
758                build_leaf(
759                    seq,
760                    at,
761                    &leaf::JobChange {
762                        job_id: &s.response.id,
763                        project_key,
764                        file_path,
765                        state: &s.response.state,
766                        stored_sha256: stored.as_deref(),
767                        follow_up: s.response.follow_up.as_deref(),
768                    },
769                )
770            },
771        )?;
772        Ok(settlement)
773    }
774
775    /// Jobs, newest first, at most `limit`, in `state` when given.
776    pub fn jobs(&self, state: Option<&str>, limit: usize) -> Result<Vec<JobSummary>> {
777        let conn = self.lock();
778        let mut stmt = conn.prepare(&format!(
779            "SELECT {JOB_COLUMNS} FROM jobs WHERE ?1 IS NULL OR state = ?1
780             ORDER BY created_at DESC, id DESC LIMIT ?2"
781        ))?;
782        let rows = stmt.query_map((state, limit as i64), row_from)?;
783        Ok(rows
784            .map(|r| r.map(|job| job.summary()))
785            .collect::<rusqlite::Result<_>>()?)
786    }
787
788    /// Queues a failed job again, with its attempts counted afresh, and
789    /// appends the `job_retry` leaf `build_leaf` makes in the same
790    /// transaction.
791    ///
792    /// A job that failed because the file kept changing kept its result:
793    /// it is queued as a merge of that result with the file as it is now,
794    /// which is the merge that was never finished.
795    pub fn retry_job_audited(
796        &self,
797        id: &str,
798        now: OffsetDateTime,
799        build_leaf: impl FnOnce(u64, &str, &JobSummary) -> Vec<u8>,
800    ) -> Result<Retried> {
801        self.audited(
802            |tx, _| {
803                let Some(job) = get_job(tx, id)? else {
804                    return Ok(Outcome::Refuse(Retried::NotFound));
805                };
806                if job.state != STATE_FAILED {
807                    return Ok(Outcome::Refuse(Retried::NotFailed(job.state)));
808                }
809                let mut payload = job.payload.clone();
810                if let (Some(kept), KIND_MERGE) = (&job.result, job.kind.as_str()) {
811                    let kept: MergeSide = serde_json::from_str(kept)?;
812                    if let Some(file) =
813                        read_file(tx, &job.project_key, &job.file_path)?.filter(|e| !e.deleted)
814                    {
815                        payload = serde_json::to_string(&MergeInput {
816                            project_key: job.project_key.clone(),
817                            file_path: job.file_path.clone(),
818                            stored: kept,
819                            incoming: side_of(&file.content, &file.source_env, &file.updated_at),
820                        })?;
821                    }
822                }
823                let now = ts(now);
824                tx.execute(
825                    "UPDATE jobs SET state = 'queued', payload = ?2, lease_id = NULL,
826                         lease_expires_at = NULL, attempt = 0, not_before = ?3, link = 0,
827                         error = NULL, result = NULL, applied = 0, follow_up = NULL, updated_at = ?3
828                     WHERE id = ?1",
829                    (id, &payload, &now),
830                )?;
831                let job = get_job(tx, id)?.context("the job was read above")?;
832                Ok(Outcome::Commit(Retried::Queued(job.summary())))
833            },
834            |seq, at, retried| match retried {
835                Retried::Queued(job) => build_leaf(seq, at, job),
836                _ => unreachable!("a leaf only for a retry"),
837            },
838        )
839    }
840
841    /// Makes every open job claimable at once: a leased one is released,
842    /// its holder being gone, and a queued one's retry delay is dropped.
843    /// Answers how many there are.
844    ///
845    /// For the server draining the queue itself once no worker is left to:
846    /// a revoked worker's lease would otherwise hold its job until it ran
847    /// out, and a job waiting out a delay the worker's failure set has no
848    /// worker left to wait for. Attempts still count, so a job that keeps
849    /// failing here is failed all the same.
850    ///
851    /// Not in the audit log: it changes no file, and what it frees is
852    /// recorded around it — the `revoke` of the worker that held a lease
853    /// before, and the server's own `job_claim` of each job after.
854    ///
855    /// Merge jobs only: an evaluation waits for a worker, since the server
856    /// does not evaluate, and a lease a revoked worker held on one runs out
857    /// by itself.
858    pub fn release_open_jobs(&self, now: OffsetDateTime) -> Result<usize> {
859        Ok(self.lock().execute(
860            "UPDATE jobs SET state = 'queued', lease_id = NULL, lease_expires_at = NULL,
861                 not_before = ?1
862             WHERE kind = 'merge' AND state IN ('queued', 'leased')",
863            (ts(now),),
864        )?)
865    }
866
867    /// Marks every open job failed, with `why` as its error, for a queue
868    /// nothing is left to drain. Each keeps its input, so a retry merges it
869    /// once something can. Answers their ids.
870    ///
871    /// Each job gets a `job_result` leaf, which `build_leaf` makes, in the
872    /// same transaction, as a lease that ran out on a last attempt does.
873    pub fn fail_open_jobs_audited(
874        &self,
875        why: &str,
876        now: OffsetDateTime,
877        build_leaf: impl Fn(u64, &str, &leaf::JobChange<'_>) -> Vec<u8>,
878    ) -> Result<Vec<String>> {
879        let jobs = self.audited_each(
880            |tx, _| {
881                let jobs = {
882                    let mut stmt = tx.prepare(
883                        "SELECT id, project_key, file_path FROM jobs
884                         WHERE kind = 'merge' AND state IN ('queued', 'leased')
885                         ORDER BY created_at, id",
886                    )?;
887                    let rows = stmt.query_map([], |r| {
888                        Ok((
889                            r.get::<_, String>(0)?,
890                            r.get::<_, String>(1)?,
891                            r.get::<_, String>(2)?,
892                        ))
893                    })?;
894                    rows.collect::<rusqlite::Result<Vec<_>>>()?
895                };
896                if jobs.is_empty() {
897                    return Ok(Outcome::Refuse(jobs));
898                }
899                tx.execute(
900                    "UPDATE jobs SET state = 'failed', lease_id = NULL, lease_expires_at = NULL,
901                         error = ?1, applied = 0, updated_at = ?2
902                     WHERE kind = 'merge' AND state IN ('queued', 'leased')",
903                    (clip(why, MAX_ERROR_BYTES), ts(now)),
904                )?;
905                Ok(Outcome::Commit(jobs))
906            },
907            |seq, at, jobs| {
908                jobs.iter()
909                    .zip(seq..)
910                    .map(|((job_id, project_key, file_path), seq)| {
911                        build_leaf(
912                            seq,
913                            at,
914                            &leaf::JobChange {
915                                job_id,
916                                project_key,
917                                file_path,
918                                state: STATE_FAILED,
919                                stored_sha256: None,
920                                follow_up: None,
921                            },
922                        )
923                    })
924                    .collect()
925            },
926        )?;
927        Ok(jobs.into_iter().map(|(id, _, _)| id).collect())
928    }
929
930    /// What `/health` says about the merge queue. Merge jobs only: an
931    /// evaluation waiting is not a merge held up, and `GET
932    /// /v1/evaluations` shows it.
933    pub fn queue_status(&self) -> Result<QueueStatus> {
934        let conn = self.lock();
935        let count = |state: &str| -> Result<u64> {
936            let n: i64 = conn.query_row(
937                "SELECT COUNT(*) FROM jobs WHERE kind = 'merge' AND state = ?1",
938                (state,),
939                |r| r.get(0),
940            )?;
941            Ok(n as u64)
942        };
943        Ok(QueueStatus {
944            queued: count(STATE_QUEUED)?,
945            leased: count(STATE_LEASED)?,
946            failed: count(STATE_FAILED)?,
947            oldest_queued_at: conn.query_row(
948                "SELECT MIN(created_at) FROM jobs WHERE kind = 'merge' AND state = 'queued'",
949                [],
950                |r| r.get(0),
951            )?,
952        })
953    }
954
955    /// Removes jobs that finished before `before`. Failed jobs stay until
956    /// someone retries them: each may hold a result nobody has seen.
957    ///
958    /// Not in the audit log: the leaves of the jobs removed stay in it.
959    pub fn prune_jobs(&self, before: &str) -> Result<usize> {
960        Ok(self.lock().execute(
961            "DELETE FROM jobs WHERE state = 'done' AND updated_at < ?1",
962            (before,),
963        )?)
964    }
965}
966
967/// Applies a merged result to its file, as a compare-and-swap: written
968/// only if the file still has the hash of the version the job stored, and
969/// then attributed to `worker`. If another push landed meanwhile, the
970/// merged content becomes the `stored` side of a follow-up job,
971/// `follow_up_id`, against the newer version, at most [`MAX_LINKS`] deep.
972#[allow(clippy::too_many_arguments)]
973fn apply_merge(
974    tx: &Connection,
975    job: &JobRow,
976    input: &MergeInput,
977    merged: &str,
978    worker: &str,
979    follow_up_id: &str,
980    now: OffsetDateTime,
981) -> Result<Settled> {
982    let now_text = ts(now);
983    let merged_side = side_of(merged, worker, &now_text);
984    let current = read_file(tx, &input.project_key, &input.file_path)?.filter(|e| !e.deleted);
985    let settled = |applied, queued, failed| Settled {
986        response: ResultResponse::default(),
987        applied,
988        queued,
989        failed,
990    };
991    Ok(match current {
992        // The compare-and-swap: nothing landed since the push that queued
993        // this, so the merge replaces it.
994        Some(file) if content_sha256(&file.content) == input.incoming.sha256 => {
995            write_file(
996                tx,
997                &input.project_key,
998                &input.file_path,
999                merged,
1000                worker,
1001                &now_text,
1002            )?;
1003            finish(tx, job, STATE_DONE, true, None, None, None, now)?;
1004            settled(true, false, None)
1005        }
1006        // Something else landed meanwhile, and says exactly what the
1007        // merge came to.
1008        Some(file) if file.content == merged => {
1009            finish(tx, job, STATE_DONE, true, None, None, None, now)?;
1010            settled(false, false, None)
1011        }
1012        // Something else landed: merge the result with it, unless this
1013        // conflict has been chased far enough.
1014        Some(file) if job.link < MAX_LINKS => {
1015            let next = MergeInput {
1016                project_key: input.project_key.clone(),
1017                file_path: input.file_path.clone(),
1018                stored: merged_side,
1019                incoming: side_of(&file.content, &file.source_env, &file.updated_at),
1020            };
1021            insert_merge_job(tx, follow_up_id, &next, Some((&job.id, job.link + 1)), now)?;
1022            finish(
1023                tx,
1024                job,
1025                STATE_DONE,
1026                false,
1027                None,
1028                None,
1029                Some(follow_up_id),
1030                now,
1031            )?;
1032            settled(false, true, None)
1033        }
1034        Some(_) => {
1035            let error = "the file kept changing while it was merged; the newest push stands, \
1036                         and this result is kept in the job";
1037            finish(
1038                tx,
1039                job,
1040                STATE_FAILED,
1041                false,
1042                Some(error),
1043                Some(&merged_side),
1044                None,
1045                now,
1046            )?;
1047            settled(
1048                false,
1049                false,
1050                Some(Box::new(Failure {
1051                    id: job.id.clone(),
1052                    kind: job.kind.clone(),
1053                    project_key: job.project_key.clone(),
1054                    file_path: job.file_path.clone(),
1055                    what: "was not applied: the file kept changing while it was merged".to_string(),
1056                    error: String::new(),
1057                })),
1058            )
1059        }
1060        // Deleted meanwhile: the delete said to discard it, as a push after
1061        // a delete is never merged either. The result is kept in the job
1062        // all the same. A delete closes its file's open jobs as it lands
1063        // (see [`close_for_delete`]), so this is a job retried after one.
1064        None => {
1065            finish(
1066                tx,
1067                job,
1068                STATE_DONE,
1069                false,
1070                Some("the file was deleted while it was merged; the delete stands"),
1071                Some(&merged_side),
1072                None,
1073                now,
1074            )?;
1075            settled(false, false, None)
1076        }
1077    })
1078}
1079
1080/// Closes every merge job still waiting on, or held for, one file, as the
1081/// delete of that file lands, in the delete's transaction.
1082///
1083/// A delete says to discard the file, and a push after a delete is never
1084/// merged with what was there. A job left open would say otherwise: its
1085/// result, landing on the file a later push created, would not match the
1086/// hash it was queued against, and would be chased onto the new file as a
1087/// follow-up, bringing the deleted notes back into it. Closed now, it is
1088/// `done` and unapplied, and a result its holder still posts is answered as
1089/// one already recorded, changing nothing.
1090pub(super) fn close_for_delete(
1091    conn: &Connection,
1092    project_key: &str,
1093    file_path: &str,
1094    now: &str,
1095) -> Result<usize> {
1096    Ok(conn.execute(
1097        "UPDATE jobs SET state = 'done', payload = '{}', lease_expires_at = NULL, applied = 0,
1098             error = ?3, updated_at = ?4
1099         WHERE kind = 'merge' AND project_key = ?1 AND file_path = ?2
1100           AND state IN ('queued', 'leased')",
1101        (
1102            project_key,
1103            file_path,
1104            "the file was deleted before it was merged; the delete stands",
1105            now,
1106        ),
1107    )?)
1108}
1109
1110/// Why `recall-server admin` closes a job, which decides what it says and
1111/// what it keeps. See [`close_for_admin`].
1112pub(super) enum AdminClose<'a> {
1113    /// `remove`: the file's key is gone.
1114    Removed,
1115    /// `rename`: the file is under `to` now.
1116    Renamed {
1117        /// The key the rows moved to.
1118        to: &'a str,
1119    },
1120    /// `restore`: the file holds the backup's version now.
1121    Restored,
1122}
1123
1124/// Closes every job still waiting on, or held for, the rows an admin change
1125/// touches, in the change's own transaction: every file under
1126/// `project_key`, or only `file_path` when one is given. Answers how many
1127/// it closed.
1128///
1129/// Closed as a delete closes them ([`close_for_delete`]): `done`,
1130/// unapplied, the lease kept, so a result its holder still posts is
1131/// answered as one already recorded and changes nothing. Left open, each
1132/// would be settled against a row the change moved, removed or replaced:
1133///
1134/// - after a remove, its result would be chased onto a file a machine still
1135///   syncing under the key pushes later, as a follow-up, bringing the
1136///   removed notes back into it;
1137/// - after a restore, the file no longer has the hash the job was queued
1138///   against, so its result, a merge of the versions the restore replaced,
1139///   would be chased onto the restored file as a follow-up, and undo the
1140///   restore once that is merged;
1141/// - after a rename, the job still names the old key, and its versions are
1142///   the old key's rows.
1143///
1144/// A rename could instead move a job to the new key. It does not, because
1145/// what makes a result safe to apply is the compare-and-swap against the
1146/// version the job was queued with, under the key it was queued with, and a
1147/// rename is the owner saying the old key's history ends here: a job
1148/// re-keyed would be applied to a key no push to it ever queued, on the
1149/// strength of a hash taken under another, and one whose file has moved on
1150/// since (a push after it, a result applied) would be chased across keys.
1151/// Closing costs one merge, and nothing is lost: the file keeps the newer
1152/// version the push stored, the version it displaced stays in the job
1153/// (a rename and a restore keep the job's input; only a remove drops it,
1154/// as a delete does), and the backup every change takes holds both.
1155pub(super) fn close_for_admin(
1156    conn: &Connection,
1157    project_key: &str,
1158    file_path: Option<&str>,
1159    why: &AdminClose<'_>,
1160    now: &str,
1161) -> Result<usize> {
1162    let (error, keep_input) = match why {
1163        AdminClose::Removed => (
1164            "removed by admin: `recall-server admin remove` deleted this file's project key \
1165             before it was merged; the remove stands"
1166                .to_string(),
1167            false,
1168        ),
1169        AdminClose::Renamed { to } => (
1170            format!(
1171                "renamed by admin: `recall-server admin rename` moved this file to {to:?} before \
1172                 it was merged; its versions are kept in this job"
1173            ),
1174            true,
1175        ),
1176        AdminClose::Restored => (
1177            "restored by admin: `recall-server admin restore` put back a backup's version of this \
1178             file before it was merged; the restore stands, and the versions this job held are \
1179             kept in it"
1180                .to_string(),
1181            true,
1182        ),
1183    };
1184    Ok(conn.execute(
1185        "UPDATE jobs SET state = 'done',
1186             payload = CASE WHEN ?4 THEN payload ELSE '{}' END,
1187             lease_expires_at = NULL, applied = 0, error = ?3, updated_at = ?5
1188         WHERE project_key = ?1 AND (?2 IS NULL OR file_path = ?2)
1189           AND state IN ('queued', 'leased')",
1190        (
1191            project_key,
1192            file_path,
1193            clip(&error, MAX_ERROR_BYTES),
1194            keep_input,
1195            now,
1196        ),
1197    )?)
1198}
1199
1200/// Marks a job settled, keeping the lease it was settled under so the same
1201/// result posted again is recognised. A done job's input is dropped: the
1202/// file holds what mattered now, and the job need not hold a second copy.
1203#[allow(clippy::too_many_arguments)]
1204fn finish(
1205    conn: &Connection,
1206    job: &JobRow,
1207    state: &str,
1208    applied: bool,
1209    error: Option<&str>,
1210    kept: Option<&MergeSide>,
1211    follow_up: Option<&str>,
1212    now: OffsetDateTime,
1213) -> Result<()> {
1214    let kept = kept.map(serde_json::to_string).transpose()?;
1215    let payload = if state == STATE_DONE {
1216        "{}".to_string()
1217    } else {
1218        job.payload.clone()
1219    };
1220    conn.execute(
1221        "UPDATE jobs SET state = ?2, applied = ?3, error = ?4, result = ?5, follow_up = ?6,
1222             payload = ?7, lease_expires_at = NULL, updated_at = ?8
1223         WHERE id = ?1",
1224        (
1225            &job.id,
1226            state,
1227            applied as i64,
1228            error,
1229            kept,
1230            follow_up,
1231            payload,
1232            ts(now),
1233        ),
1234    )?;
1235    Ok(())
1236}
1237
1238#[cfg(test)]
1239mod tests {
1240    use super::*;
1241    use crate::store::test_leaf;
1242    use recall_wire::MergeResult;
1243
1244    /// The audited writes, with a leaf these tests do not look at, under
1245    /// the names the tests read best with. A file written here keeps the
1246    /// `updated_at` a test gives it, since these tests move the clock.
1247    impl Store {
1248        fn upsert(&self, pk: &str, fp: &str, content: &str, env: &str, at: &str) -> Result<()> {
1249            self.audited(
1250                |tx, _| {
1251                    write_file(tx, pk, fp, content, env, at)?;
1252                    Ok(Outcome::Commit(()))
1253                },
1254                |seq, at, ()| test_leaf(seq, at),
1255            )
1256        }
1257
1258        fn tombstone(&self, pk: &str, fp: &str, env: &str, at: &str) -> Result<()> {
1259            self.audited(
1260                |tx, _| {
1261                    tx.execute(
1262                        "INSERT INTO memory_files
1263                             (project_key, file_path, content, source_env, updated_at, deleted)
1264                         VALUES (?1, ?2, '', ?3, ?4, 1)
1265                         ON CONFLICT(project_key, file_path) DO UPDATE SET
1266                             source_env = excluded.source_env,
1267                             updated_at = excluded.updated_at,
1268                             deleted = 1",
1269                        (pk, fp, env, at),
1270                    )?;
1271                    close_for_delete(tx, pk, fp, at)?;
1272                    Ok(Outcome::Commit(()))
1273                },
1274                |seq, at, ()| test_leaf(seq, at),
1275            )
1276        }
1277
1278        fn write_and_queue_merge(
1279            &self,
1280            pk: &str,
1281            fp: &str,
1282            incoming: &MergeSide,
1283            job_id: &str,
1284            now: OffsetDateTime,
1285        ) -> Result<Queued> {
1286            self.write_and_queue_merge_audited(pk, fp, incoming, job_id, now, |seq, at, _| {
1287                test_leaf(seq, at)
1288            })
1289            .map(|(queued, _)| queued)
1290        }
1291
1292        fn claim_job(
1293            &self,
1294            kinds: &[String],
1295            lease_id: &str,
1296            lease: Duration,
1297            now: OffsetDateTime,
1298        ) -> Result<Option<Job>> {
1299            self.claim_job_audited(kinds, lease_id, lease, now, |seq, at, _| test_leaf(seq, at))
1300        }
1301
1302        fn settle_job(
1303            &self,
1304            id: &str,
1305            result: &ResultRequest,
1306            worker: &str,
1307            follow_up_id: &str,
1308            now: OffsetDateTime,
1309        ) -> Result<Settlement> {
1310            self.settle_job_audited(id, result, worker, follow_up_id, now, |seq, at, _| {
1311                test_leaf(seq, at)
1312            })
1313        }
1314
1315        fn retry_job(&self, id: &str, now: OffsetDateTime) -> Result<Retried> {
1316            self.retry_job_audited(id, now, |seq, at, _| test_leaf(seq, at))
1317        }
1318
1319        fn expire_leases(&self, now: OffsetDateTime) -> Result<Vec<Failure>> {
1320            self.expire_leases_audited(now, |seq, at, _| test_leaf(seq, at))
1321        }
1322
1323        fn fail_open_jobs(&self, why: &str, now: OffsetDateTime) -> Result<Vec<String>> {
1324            self.fail_open_jobs_audited(why, now, |seq, at, _| test_leaf(seq, at))
1325        }
1326    }
1327
1328    const P: &str = "acme/app";
1329    const F: &str = "topics/auth.md";
1330
1331    fn at(secs: i64) -> OffsetDateTime {
1332        OffsetDateTime::from_unix_timestamp(1_790_000_000 + secs).unwrap()
1333    }
1334
1335    fn side(content: &str, by: &str, secs: i64) -> MergeSide {
1336        side_of(content, by, &ts(at(secs)))
1337    }
1338
1339    fn merge_kinds() -> Vec<String> {
1340        vec![KIND_MERGE.to_string()]
1341    }
1342
1343    fn merged(lease: &str, content: &str) -> ResultRequest {
1344        ResultRequest {
1345            lease_id: lease.into(),
1346            merge: Some(MergeResult {
1347                content: content.into(),
1348            }),
1349            error: None,
1350            evaluate: None,
1351        }
1352    }
1353
1354    fn failed(lease: &str, why: &str) -> ResultRequest {
1355        ResultRequest {
1356            lease_id: lease.into(),
1357            merge: None,
1358            error: Some(why.into()),
1359            evaluate: None,
1360        }
1361    }
1362
1363    fn recorded(s: Settlement) -> Settled {
1364        match s {
1365            Settlement::Recorded(s) => s,
1366            other => panic!("not recorded: {other:?}"),
1367        }
1368    }
1369
1370    /// A store holding `A` for the file, then a push of `B` queued as job
1371    /// `job_1`.
1372    fn conflicted() -> Store {
1373        let st = Store::open_in_memory().unwrap();
1374        st.upsert(P, F, "A", "laptop", &ts(at(0))).unwrap();
1375        assert_eq!(
1376            st.write_and_queue_merge(P, F, &side("B", "cloud", 1), "job_1", at(1))
1377                .unwrap(),
1378            Queued::Queued("job_1".into())
1379        );
1380        st
1381    }
1382
1383    fn content(st: &Store) -> (String, String) {
1384        let e = st.get(P, F).unwrap().unwrap();
1385        (e.content, e.source_env)
1386    }
1387
1388    #[test]
1389    fn a_stale_push_is_stored_at_once_and_queued_with_both_versions() {
1390        let st = conflicted();
1391        assert_eq!(content(&st), ("B".into(), "cloud".into()));
1392        let job = st
1393            .claim_job(&merge_kinds(), "lse_1", Duration::from_secs(120), at(2))
1394            .unwrap()
1395            .unwrap();
1396        assert_eq!((job.id.as_str(), job.attempt), ("job_1", 1));
1397        assert_eq!(job.lease_expires_at, ts(at(122)));
1398        let m = job.merge.unwrap();
1399        assert_eq!(
1400            (m.stored.content.as_str(), m.stored.source_env.as_str()),
1401            ("A", "laptop")
1402        );
1403        assert_eq!(m.stored.sha256, content_sha256("A"));
1404        // Stamped as the file it stored was: with its leaf's `at`.
1405        let stored = st.get(P, F).unwrap().unwrap();
1406        assert_eq!(
1407            m.incoming,
1408            MergeSide {
1409                updated_at: stored.updated_at,
1410                ..side("B", "cloud", 1)
1411            }
1412        );
1413        // Leased: nobody else gets it.
1414        assert!(st
1415            .claim_job(&merge_kinds(), "lse_2", Duration::from_secs(120), at(3))
1416            .unwrap()
1417            .is_none());
1418    }
1419
1420    #[test]
1421    fn nothing_is_queued_without_something_to_merge() {
1422        let st = Store::open_in_memory().unwrap();
1423        assert_eq!(
1424            st.write_and_queue_merge(P, F, &side("B", "cloud", 1), "job_1", at(1))
1425                .unwrap(),
1426            Queued::Nothing
1427        );
1428        assert_eq!(content(&st).0, "B");
1429        assert!(st.jobs(None, 10).unwrap().is_empty());
1430    }
1431
1432    #[test]
1433    fn a_claim_asks_for_kinds_and_takes_the_oldest() {
1434        let st = conflicted();
1435        st.upsert(P, "other.md", "X", "laptop", &ts(at(0))).unwrap();
1436        st.write_and_queue_merge(P, "other.md", &side("Y", "cloud", 5), "job_0", at(5))
1437            .unwrap();
1438        assert!(st
1439            .claim_job(&[], "lse", Duration::from_secs(60), at(6))
1440            .unwrap()
1441            .is_none());
1442        assert!(st
1443            .claim_job(&["seal".into()], "lse", Duration::from_secs(60), at(6))
1444            .unwrap()
1445            .is_none());
1446        let first = st
1447            .claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(6))
1448            .unwrap()
1449            .unwrap();
1450        assert_eq!(
1451            first.id, "job_1",
1452            "created first, though its id sorts later"
1453        );
1454    }
1455
1456    /// The property the lease exists for: a job a worker took and never
1457    /// finished comes back, and the old holder's result no longer counts.
1458    #[test]
1459    fn an_expired_lease_is_released_and_its_holder_fenced_off() {
1460        let st = conflicted();
1461        st.claim_job(&merge_kinds(), "lse_old", Duration::from_secs(60), at(2))
1462            .unwrap()
1463            .unwrap();
1464        // Still inside the lease: nothing to release.
1465        assert!(st.expire_leases(at(61)).unwrap().is_empty());
1466        assert_eq!(st.queue_status().unwrap().leased, 1);
1467        // Past it: queued again, for a minute from now.
1468        assert!(st.expire_leases(at(62)).unwrap().is_empty());
1469        let q = st.queue_status().unwrap();
1470        assert_eq!((q.queued, q.leased), (1, 0));
1471        assert!(st
1472            .claim_job(&merge_kinds(), "lse_new", Duration::from_secs(60), at(100))
1473            .unwrap()
1474            .is_none());
1475        let again = st
1476            .claim_job(&merge_kinds(), "lse_new", Duration::from_secs(60), at(122))
1477            .unwrap()
1478            .unwrap();
1479        assert_eq!(again.attempt, 2);
1480        // The first holder reports late: refused, and nothing changes.
1481        assert_eq!(
1482            st.settle_job(
1483                "job_1",
1484                &merged("lse_old", "late"),
1485                "worker",
1486                "job_f",
1487                at(123)
1488            )
1489            .unwrap(),
1490            Settlement::LeaseEnded
1491        );
1492        assert_eq!(content(&st).0, "B");
1493        // The current holder's result counts.
1494        let s = recorded(
1495            st.settle_job(
1496                "job_1",
1497                &merged("lse_new", "AB"),
1498                "worker",
1499                "job_f",
1500                at(124),
1501            )
1502            .unwrap(),
1503        );
1504        assert!(s.applied && s.response.applied);
1505        assert_eq!(content(&st), ("AB".into(), "worker".into()));
1506    }
1507
1508    /// A lease that ran out is over even before a sweep has noticed.
1509    #[test]
1510    fn a_result_after_the_lease_ends_is_refused_even_unswept() {
1511        let st = conflicted();
1512        st.claim_job(&merge_kinds(), "lse_1", Duration::from_secs(60), at(2))
1513            .unwrap();
1514        assert_eq!(
1515            st.settle_job("job_1", &merged("lse_1", "AB"), "worker", "job_f", at(62))
1516                .unwrap(),
1517            Settlement::LeaseEnded
1518        );
1519        assert_eq!(content(&st).0, "B");
1520    }
1521
1522    #[test]
1523    fn a_result_is_applied_once_and_a_repeat_changes_nothing() {
1524        let st = conflicted();
1525        st.claim_job(&merge_kinds(), "lse_1", Duration::from_secs(60), at(2))
1526            .unwrap();
1527        let first = recorded(
1528            st.settle_job("job_1", &merged("lse_1", "AB"), "worker", "job_f", at(3))
1529                .unwrap(),
1530        );
1531        assert_eq!(
1532            first.response,
1533            ResultResponse {
1534                id: "job_1".into(),
1535                state: "done".into(),
1536                applied: true,
1537                follow_up: None
1538            }
1539        );
1540        // Someone pushes after the merge; the repeat must not undo it.
1541        st.upsert(P, F, "C", "laptop", &ts(at(4))).unwrap();
1542        let again = recorded(
1543            st.settle_job("job_1", &merged("lse_1", "AB"), "worker", "job_g", at(5))
1544                .unwrap(),
1545        );
1546        assert_eq!(again.response, first.response);
1547        assert!(!again.applied);
1548        assert_eq!(content(&st).0, "C");
1549        assert_eq!(st.jobs(None, 10).unwrap().len(), 1, "no follow-up either");
1550    }
1551
1552    /// The compare-and-swap: a push that landed while the job ran is not
1553    /// overwritten, and the merge is chased onto it instead.
1554    #[test]
1555    fn a_result_for_a_file_that_moved_on_becomes_a_follow_up() {
1556        let st = conflicted();
1557        st.claim_job(&merge_kinds(), "lse_1", Duration::from_secs(60), at(2))
1558            .unwrap();
1559        st.upsert(P, F, "C", "laptop", &ts(at(3))).unwrap();
1560        let s = recorded(
1561            st.settle_job("job_1", &merged("lse_1", "AB"), "worker", "job_2", at(4))
1562                .unwrap(),
1563        );
1564        assert_eq!(
1565            (s.applied, s.queued, s.response.follow_up.as_deref()),
1566            (false, true, Some("job_2"))
1567        );
1568        assert_eq!(s.response.state, "done");
1569        assert_eq!(content(&st).0, "C", "the push in between stands");
1570        let next = st
1571            .claim_job(&merge_kinds(), "lse_2", Duration::from_secs(60), at(5))
1572            .unwrap()
1573            .unwrap();
1574        let m = next.merge.unwrap();
1575        assert_eq!(
1576            (m.stored.content.as_str(), m.stored.source_env.as_str()),
1577            ("AB", "worker")
1578        );
1579        assert_eq!(m.incoming.content, "C");
1580        recorded(
1581            st.settle_job("job_2", &merged("lse_2", "ABC"), "worker", "job_3", at(6))
1582                .unwrap(),
1583        );
1584        assert_eq!(content(&st).0, "ABC");
1585    }
1586
1587    #[test]
1588    fn a_chase_stops_after_three_links_and_keeps_the_result() {
1589        let st = conflicted();
1590        let mut id = "job_1".to_string();
1591        for link in 0..=MAX_LINKS {
1592            let t = 10 * (i64::from(link) + 1);
1593            let lease = format!("lse_{link}");
1594            let job = st
1595                .claim_job(&merge_kinds(), &lease, Duration::from_secs(60), at(t))
1596                .unwrap()
1597                .unwrap();
1598            assert_eq!(job.id, id);
1599            st.upsert(P, F, &format!("push {link}"), "laptop", &ts(at(t + 1)))
1600                .unwrap();
1601            let next = format!("job_next_{link}");
1602            let s = recorded(
1603                st.settle_job(&id, &merged(&lease, "merged"), "worker", &next, at(t + 2))
1604                    .unwrap(),
1605            );
1606            if link < MAX_LINKS {
1607                assert_eq!(s.response.follow_up.as_deref(), Some(next.as_str()));
1608                id = next;
1609            } else {
1610                assert_eq!(s.response.state, "failed");
1611                assert!(s.failed.is_some());
1612            }
1613        }
1614        assert_eq!(content(&st).0, format!("push {MAX_LINKS}"));
1615        // Retrying merges the kept result with the file as it is now.
1616        let Retried::Queued(job) = st.retry_job(&id, at(50)).unwrap() else {
1617            panic!("not retried")
1618        };
1619        assert_eq!((job.state.as_str(), job.attempt), ("queued", 0));
1620        let m = st
1621            .claim_job(&merge_kinds(), "lse_r", Duration::from_secs(60), at(51))
1622            .unwrap()
1623            .unwrap()
1624            .merge
1625            .unwrap();
1626        assert_eq!(
1627            (m.stored.content.as_str(), m.incoming.content.as_str()),
1628            ("merged", "push 3")
1629        );
1630    }
1631
1632    #[test]
1633    fn errors_retry_after_one_five_and_thirty_minutes_then_fail() {
1634        let st = conflicted();
1635        let mut t = 2;
1636        for (attempt, wait) in [(1, 60), (2, 300), (3, 1800)] {
1637            let lease = format!("lse_{attempt}");
1638            let job = st
1639                .claim_job(&merge_kinds(), &lease, Duration::from_secs(60), at(t))
1640                .unwrap()
1641                .unwrap();
1642            assert_eq!(job.attempt, attempt);
1643            let s = recorded(
1644                st.settle_job(
1645                    "job_1",
1646                    &failed(&lease, "claude timed out"),
1647                    "w",
1648                    "j",
1649                    at(t),
1650                )
1651                .unwrap(),
1652            );
1653            assert_eq!(s.response.state, "queued");
1654            assert!(s.failed.is_none());
1655            // A repeat of the same error is the same answer, and no change.
1656            assert_eq!(
1657                recorded(
1658                    st.settle_job(
1659                        "job_1",
1660                        &failed(&lease, "claude timed out"),
1661                        "w",
1662                        "j",
1663                        at(t)
1664                    )
1665                    .unwrap()
1666                )
1667                .response
1668                .state,
1669                "queued"
1670            );
1671            assert!(st
1672                .claim_job(
1673                    &merge_kinds(),
1674                    "x",
1675                    Duration::from_secs(60),
1676                    at(t + wait - 1)
1677                )
1678                .unwrap()
1679                .is_none());
1680            t += wait;
1681        }
1682        st.claim_job(&merge_kinds(), "lse_4", Duration::from_secs(60), at(t))
1683            .unwrap()
1684            .unwrap();
1685        let s = recorded(
1686            st.settle_job(
1687                "job_1",
1688                &failed("lse_4", "claude timed out"),
1689                "w",
1690                "j",
1691                at(t),
1692            )
1693            .unwrap(),
1694        );
1695        assert_eq!(s.response.state, "failed");
1696        let failure = s.failed.unwrap();
1697        assert!(failure.error.contains("claude timed out"));
1698        assert_eq!(failure.what, "failed after 4 attempts");
1699        assert_eq!(content(&st).0, "B", "last-write-wins stands");
1700        let listed = st.jobs(Some("failed"), 10).unwrap();
1701        assert_eq!(listed[0].attempt, 4);
1702        assert!(listed[0]
1703            .error
1704            .as_deref()
1705            .unwrap()
1706            .contains("gave up after 4"));
1707    }
1708
1709    #[test]
1710    fn a_lease_that_runs_out_on_the_last_attempt_fails_the_job() {
1711        let st = conflicted();
1712        st.lock()
1713            .execute("UPDATE jobs SET attempt = 3", [])
1714            .unwrap();
1715        st.claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(2))
1716            .unwrap()
1717            .unwrap();
1718        let failed = st.expire_leases(at(100)).unwrap();
1719        assert_eq!(failed.len(), 1);
1720        assert_eq!(st.queue_status().unwrap().failed, 1);
1721    }
1722
1723    #[test]
1724    fn a_delete_meanwhile_stands() {
1725        let st = conflicted();
1726        st.claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(2))
1727            .unwrap();
1728        st.tombstone(P, F, "laptop", &ts(at(3))).unwrap();
1729        let s = recorded(
1730            st.settle_job("job_1", &merged("lse", "AB"), "worker", "job_2", at(4))
1731                .unwrap(),
1732        );
1733        assert_eq!((s.response.state.as_str(), s.applied), ("done", false));
1734        assert!(st.get(P, F).unwrap().unwrap().deleted);
1735    }
1736
1737    #[test]
1738    fn the_queue_is_bounded_and_a_full_one_still_stores_the_push() {
1739        let st = Store::open_in_memory().unwrap();
1740        {
1741            let conn = st.lock();
1742            for i in 0..MAX_OPEN_JOBS {
1743                conn.execute(
1744                    "INSERT INTO jobs (id, kind, state, project_key, file_path, payload,
1745                                       not_before, created_at, updated_at)
1746                     VALUES (?1, 'merge', 'queued', 'p', 'f', '{}', 'x', 'x', 'x')",
1747                    (format!("job_{i}"),),
1748                )
1749                .unwrap();
1750            }
1751        }
1752        st.upsert(P, F, "A", "laptop", &ts(at(0))).unwrap();
1753        assert_eq!(
1754            st.write_and_queue_merge(P, F, &side("B", "cloud", 1), "job_x", at(1))
1755                .unwrap(),
1756            Queued::Full
1757        );
1758        assert_eq!(content(&st).0, "B");
1759    }
1760
1761    #[test]
1762    fn only_a_failed_job_is_retried() {
1763        let st = conflicted();
1764        assert!(matches!(
1765            st.retry_job("job_1", at(2)).unwrap(),
1766            Retried::NotFailed(_)
1767        ));
1768        assert_eq!(st.retry_job("job_none", at(2)).unwrap(), Retried::NotFound);
1769    }
1770
1771    /// Done jobs go once old enough; failed ones stay however old, since
1772    /// each may hold a result nobody has seen, and open ones stay too.
1773    #[test]
1774    fn finished_jobs_are_pruned_and_failed_ones_kept() {
1775        let st = conflicted();
1776        st.upsert(P, "other.md", "X", "laptop", &ts(at(0))).unwrap();
1777        st.write_and_queue_merge(P, "other.md", &side("Y", "cloud", 1), "job_f", at(1))
1778            .unwrap();
1779        st.upsert(P, "third.md", "X", "laptop", &ts(at(0))).unwrap();
1780        st.write_and_queue_merge(P, "third.md", &side("Y", "cloud", 1), "job_q", at(1))
1781            .unwrap();
1782        st.claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(2))
1783            .unwrap();
1784        st.settle_job("job_1", &merged("lse", "AB"), "worker", "j", at(3))
1785            .unwrap();
1786        st.fail_open_jobs("nothing can merge this", at(3)).unwrap();
1787        st.upsert(P, "fourth.md", "X", "laptop", &ts(at(0)))
1788            .unwrap();
1789        st.write_and_queue_merge(P, "fourth.md", &side("Y", "cloud", 4), "job_o", at(4))
1790            .unwrap();
1791        assert_eq!(st.prune_jobs(&ts(at(3))).unwrap(), 0);
1792        assert_eq!(st.prune_jobs(&ts(at(1_000_000))).unwrap(), 1);
1793        let mut left: Vec<(String, String)> = st
1794            .jobs(None, 10)
1795            .unwrap()
1796            .into_iter()
1797            .map(|j| (j.id, j.state))
1798            .collect();
1799        left.sort();
1800        assert_eq!(
1801            left,
1802            vec![
1803                ("job_f".into(), "failed".into()),
1804                ("job_o".into(), "queued".into()),
1805                ("job_q".into(), "failed".into()),
1806            ]
1807        );
1808    }
1809
1810    /// The fence holds as soon as the lease runs out, before anyone has
1811    /// claimed the job again: the old holder's late result is refused, not
1812    /// taken as the repeat of one already recorded.
1813    #[test]
1814    fn an_expired_lease_fences_its_holder_before_anyone_claims_again() {
1815        let st = conflicted();
1816        st.claim_job(&merge_kinds(), "lse_old", Duration::from_secs(60), at(2))
1817            .unwrap()
1818            .unwrap();
1819        assert!(st.expire_leases(at(62)).unwrap().is_empty());
1820        assert_eq!(
1821            st.settle_job(
1822                "job_1",
1823                &merged("lse_old", "late"),
1824                "worker",
1825                "job_f",
1826                at(63)
1827            )
1828            .unwrap(),
1829            Settlement::LeaseEnded
1830        );
1831        assert_eq!(content(&st).0, "B");
1832        assert_eq!(st.jobs(Some("queued"), 10).unwrap().len(), 1);
1833    }
1834
1835    /// A push onto a deleted file is not a conflict: the delete said to
1836    /// discard what was there.
1837    #[test]
1838    fn a_push_onto_a_deleted_file_queues_nothing() {
1839        let st = Store::open_in_memory().unwrap();
1840        st.upsert(P, F, "A", "laptop", &ts(at(0))).unwrap();
1841        st.tombstone(P, F, "laptop", &ts(at(1))).unwrap();
1842        assert_eq!(
1843            st.write_and_queue_merge(P, F, &side("B", "cloud", 2), "job_1", at(2))
1844                .unwrap(),
1845            Queued::Nothing
1846        );
1847        assert_eq!(content(&st).0, "B");
1848        assert!(st.jobs(None, 10).unwrap().is_empty());
1849    }
1850
1851    /// Deleted, then made again, while its merge waited: the old notes must
1852    /// not come back into the new file, whether the job was still queued or
1853    /// already held by a worker.
1854    #[test]
1855    fn a_delete_closes_the_files_open_jobs() {
1856        // Queued when the delete landed.
1857        let st = conflicted();
1858        st.tombstone(P, F, "laptop", &ts(at(2))).unwrap();
1859        st.upsert(P, F, "fresh", "laptop", &ts(at(3))).unwrap();
1860        assert!(st
1861            .claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(4))
1862            .unwrap()
1863            .is_none());
1864        let job = &st.jobs(None, 10).unwrap()[0];
1865        assert_eq!(job.state, "done");
1866        assert!(job.error.as_deref().unwrap().contains("deleted"));
1867
1868        // Held by a worker when the delete landed.
1869        let st = conflicted();
1870        st.claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(2))
1871            .unwrap()
1872            .unwrap();
1873        st.tombstone(P, F, "laptop", &ts(at(3))).unwrap();
1874        st.upsert(P, F, "fresh", "laptop", &ts(at(4))).unwrap();
1875        let s = recorded(
1876            st.settle_job("job_1", &merged("lse", "A and B"), "worker", "job_2", at(5))
1877                .unwrap(),
1878        );
1879        assert_eq!(
1880            (s.response.state.as_str(), s.applied, s.queued),
1881            ("done", false, false)
1882        );
1883        assert_eq!(s.response.follow_up, None);
1884        assert_eq!(content(&st).0, "fresh");
1885        assert_eq!(st.jobs(None, 10).unwrap().len(), 1, "no follow-up");
1886        // Another file's job is left alone.
1887        let st = conflicted();
1888        st.tombstone(P, "other.md", "laptop", &ts(at(2))).unwrap();
1889        assert_eq!(st.queue_status().unwrap().queued, 1);
1890    }
1891
1892    /// An empty merge of two versions that had content is a malfunction:
1893    /// retried like an error, and never written over the file.
1894    #[test]
1895    fn an_empty_merge_of_non_empty_versions_is_retried_not_written() {
1896        let st = conflicted();
1897        st.claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(2))
1898            .unwrap()
1899            .unwrap();
1900        let s = recorded(
1901            st.settle_job("job_1", &merged("lse", " \n"), "worker", "job_2", at(3))
1902                .unwrap(),
1903        );
1904        assert_eq!((s.response.state.as_str(), s.applied), ("queued", false));
1905        assert_eq!(content(&st), ("B".into(), "cloud".into()));
1906        assert!(st.jobs(None, 10).unwrap()[0]
1907            .error
1908            .as_deref()
1909            .unwrap()
1910            .contains("empty"));
1911
1912        // Two empty versions may merge to nothing.
1913        let st = Store::open_in_memory().unwrap();
1914        st.upsert(P, F, "", "laptop", &ts(at(0))).unwrap();
1915        st.write_and_queue_merge(P, F, &side("\n", "cloud", 1), "job_1", at(1))
1916            .unwrap();
1917        st.claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(2))
1918            .unwrap()
1919            .unwrap();
1920        let s = recorded(
1921            st.settle_job("job_1", &merged("lse", ""), "worker", "job_2", at(3))
1922                .unwrap(),
1923        );
1924        assert!(s.applied);
1925    }
1926
1927    /// What `/health` shows of a failure names the job, never the project
1928    /// or the file; the worker's error is kept short.
1929    #[test]
1930    fn a_failure_is_public_without_its_file() {
1931        let st = conflicted();
1932        st.lock()
1933            .execute("UPDATE jobs SET attempt = 3", [])
1934            .unwrap();
1935        st.claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(2))
1936            .unwrap()
1937            .unwrap();
1938        let long = format!("{P}/{F}: {}", "é".repeat(2000));
1939        let s = recorded(
1940            st.settle_job("job_1", &failed("lse", &long), "w", "j", at(3))
1941                .unwrap(),
1942        );
1943        let failure = s.failed.unwrap();
1944        let public = failure.public();
1945        assert!(public.starts_with("merge job job_1 failed after 4 attempts"));
1946        assert!(!public.contains(P) && !public.contains(F), "{public}");
1947        assert!(failure.logged().contains(F));
1948        assert!(failure.error.len() <= MAX_ERROR_BYTES);
1949        let kept = st.jobs(None, 1).unwrap()[0].error.clone().unwrap();
1950        assert!(kept.len() <= MAX_ERROR_BYTES + 40, "{}", kept.len());
1951
1952        // Not applied because the file kept changing: the same.
1953        let st = conflicted();
1954        st.lock()
1955            .execute("UPDATE jobs SET link = ?1", (MAX_LINKS,))
1956            .unwrap();
1957        st.claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(2))
1958            .unwrap();
1959        st.upsert(P, F, "C", "laptop", &ts(at(3))).unwrap();
1960        let public = recorded(
1961            st.settle_job("job_1", &merged("lse", "AB"), "w", "j", at(4))
1962                .unwrap(),
1963        )
1964        .failed
1965        .unwrap()
1966        .public();
1967        assert!(!public.contains(P) && !public.contains(F), "{public}");
1968    }
1969
1970    /// An evaluation's error is the worker's own text: the server's log
1971    /// says what became of the job, and not that.
1972    #[test]
1973    fn an_evaluations_failure_is_logged_without_the_workers_text() {
1974        let failure = Failure {
1975            id: "job_e".into(),
1976            kind: KIND_EVALUATE.into(),
1977            project_key: String::new(),
1978            file_path: String::new(),
1979            what: "failed after 4 attempts".into(),
1980            error: "the worker said SENTINEL".into(),
1981        };
1982        assert_eq!(
1983            failure.logged(),
1984            "evaluate job job_e failed after 4 attempts"
1985        );
1986        assert!(!failure.public().contains("SENTINEL"));
1987    }
1988
1989    #[test]
1990    fn clip_keeps_whole_characters() {
1991        assert_eq!(clip("abc", 5), "abc");
1992        assert_eq!(clip("abcdef", 3), "abc");
1993        assert_eq!(clip("aé", 2), "a");
1994    }
1995
1996    /// With no worker left, every open job can be taken at once, a revoked
1997    /// holder's lease and a retry delay notwithstanding; or failed, keeping
1998    /// its input for a retry.
1999    #[test]
2000    fn open_jobs_are_released_or_failed_for_a_drain() {
2001        let st = conflicted();
2002        st.upsert(P, "other.md", "X", "laptop", &ts(at(0))).unwrap();
2003        st.write_and_queue_merge(P, "other.md", &side("Y", "cloud", 1), "job_2", at(1))
2004            .unwrap();
2005        st.claim_job(&merge_kinds(), "lse_gone", Duration::from_secs(600), at(2))
2006            .unwrap()
2007            .unwrap();
2008        st.claim_job(&merge_kinds(), "lse_x", Duration::from_secs(600), at(2))
2009            .unwrap()
2010            .unwrap();
2011        st.settle_job("job_2", &failed("lse_x", "boom"), "w", "j", at(3))
2012            .unwrap();
2013        assert_eq!(st.release_open_jobs(at(4)).unwrap(), 2);
2014        // The revoked holder's lease is over.
2015        assert_eq!(
2016            st.settle_job("job_1", &merged("lse_gone", "AB"), "w", "j", at(5))
2017                .unwrap(),
2018            Settlement::LeaseEnded
2019        );
2020        // Both claimable now, the retry delay dropped.
2021        for lease in ["lse_a", "lse_b"] {
2022            assert!(st
2023                .claim_job(&merge_kinds(), lease, Duration::from_secs(60), at(5))
2024                .unwrap()
2025                .is_some());
2026        }
2027
2028        let st = conflicted();
2029        let ids = st.fail_open_jobs("no worker, and no CLI", at(2)).unwrap();
2030        assert_eq!(ids, vec!["job_1".to_string()]);
2031        assert_eq!(st.queue_status().unwrap().failed, 1);
2032        let Retried::Queued(job) = st.retry_job("job_1", at(3)).unwrap() else {
2033            panic!("not retried")
2034        };
2035        assert_eq!(job.state, "queued");
2036        let m = st
2037            .claim_job(&merge_kinds(), "lse", Duration::from_secs(60), at(4))
2038            .unwrap()
2039            .unwrap()
2040            .merge
2041            .unwrap();
2042        assert_eq!(
2043            (m.stored.content.as_str(), m.incoming.content.as_str()),
2044            ("A", "B")
2045        );
2046    }
2047}