1use 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
28pub(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
69pub const MAX_OPEN_JOBS: usize = 1000;
74
75pub const MAX_ATTEMPTS: u32 = 4;
78
79const RETRY_AFTER: [Duration; 3] = [
81 Duration::from_secs(60),
82 Duration::from_secs(5 * 60),
83 Duration::from_secs(30 * 60),
84];
85
86pub const MAX_LINKS: u32 = 3;
90
91pub const MAX_ERROR_BYTES: usize = 500;
95
96const EMPTY_RESULT: &str = "the merge came back empty from two versions that were not";
98
99pub 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
115struct 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
201fn 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#[derive(Debug, Clone, PartialEq, Eq)]
218pub struct Failure {
219 pub id: String,
221 pub kind: String,
223 pub project_key: String,
225 pub file_path: String,
227 pub what: String,
230 pub error: String,
232}
233
234impl Failure {
235 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 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#[derive(Debug, Clone, PartialEq, Eq)]
265pub enum Queued {
266 Queued(String),
269 Nothing,
272 Full,
275}
276
277#[derive(Debug, Clone, PartialEq, Eq)]
279pub struct Settled {
280 pub response: ResultResponse,
282 pub applied: bool,
284 pub queued: bool,
287 pub failed: Option<Box<Failure>>,
290}
291
292#[derive(Debug, Clone, PartialEq, Eq)]
294pub enum Settlement {
295 Recorded(Settled),
297 NotFound,
299 LeaseEnded,
302 WrongKind(&'static str),
305 Invalid(String),
309}
310
311#[derive(Debug, Clone, PartialEq, Eq)]
313pub enum Retried {
314 Queued(JobSummary),
316 NotFound,
318 NotFailed(String),
320}
321
322fn 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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#[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 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 Some(file) if file.content == merged => {
1009 finish(tx, job, STATE_DONE, true, None, None, None, now)?;
1010 settled(false, false, None)
1011 }
1012 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 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
1080pub(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
1110pub(super) enum AdminClose<'a> {
1113 Removed,
1115 Renamed {
1117 to: &'a str,
1119 },
1120 Restored,
1122}
1123
1124pub(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#[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 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 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 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 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 #[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 assert!(st.expire_leases(at(61)).unwrap().is_empty());
1466 assert_eq!(st.queue_status().unwrap().leased, 1);
1467 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 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 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 #[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 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 #[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 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 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 #[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 #[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 #[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 #[test]
1855 fn a_delete_closes_the_files_open_jobs() {
1856 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 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 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 #[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 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 #[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 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 #[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 #[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 assert_eq!(
2016 st.settle_job("job_1", &merged("lse_gone", "AB"), "w", "j", at(5))
2017 .unwrap(),
2018 Settlement::LeaseEnded
2019 );
2020 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}