1use 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
24pub(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
65pub const MAX_OPEN_JOBS: usize = 1000;
70
71pub const MAX_ATTEMPTS: u32 = 4;
74
75const RETRY_AFTER: [Duration; 3] = [
77 Duration::from_secs(60),
78 Duration::from_secs(5 * 60),
79 Duration::from_secs(30 * 60),
80];
81
82pub const MAX_LINKS: u32 = 3;
86
87pub const MAX_ERROR_BYTES: usize = 500;
91
92const EMPTY_RESULT: &str = "the merge came back empty from two versions that were not";
94
95pub 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
111struct 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
197fn 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#[derive(Debug, Clone, PartialEq, Eq)]
214pub struct Failure {
215 pub id: String,
217 pub project_key: String,
219 pub file_path: String,
221 pub what: String,
224 pub error: String,
226}
227
228impl Failure {
229 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 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#[derive(Debug, Clone, PartialEq, Eq)]
254pub enum Queued {
255 Queued(String),
258 Nothing,
261 Full,
264}
265
266#[derive(Debug, Clone, PartialEq, Eq)]
268pub struct Settled {
269 pub response: ResultResponse,
271 pub applied: bool,
273 pub queued: bool,
276 pub failed: Option<Box<Failure>>,
279}
280
281#[derive(Debug, Clone, PartialEq, Eq)]
283pub enum Settlement {
284 Recorded(Settled),
286 NotFound,
288 LeaseEnded,
291 WrongKind,
293}
294
295#[derive(Debug, Clone, PartialEq, Eq)]
297pub enum Retried {
298 Queued(JobSummary),
300 NotFound,
302 NotFailed(String),
304}
305
306fn 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 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 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 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 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 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 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 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 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 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 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 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 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 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#[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 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 Some(file) if file.content == merged => {
941 finish(tx, job, STATE_DONE, true, None, None, None, now)?;
942 settled(false, false, None)
943 }
944 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 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
1011pub(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
1041pub(super) enum AdminClose<'a> {
1044 Removed,
1046 Renamed {
1048 to: &'a str,
1050 },
1051 Restored,
1053}
1054
1055pub(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#[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 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 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 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 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 #[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 assert!(st.expire_leases(at(61)).unwrap().is_empty());
1395 assert_eq!(st.queue_status().unwrap().leased, 1);
1396 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 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 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 #[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 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 #[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 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 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 #[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 #[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 #[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 #[test]
1784 fn a_delete_closes_the_files_open_jobs() {
1785 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 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 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 #[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 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 #[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 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 #[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 assert_eq!(
1926 st.settle_job("job_1", &merged("lse_gone", "AB"), "w", "j", at(5))
1927 .unwrap(),
1928 Settlement::LeaseEnded
1929 );
1930 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}