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