recall_wire/jobs.rs
1//! `/v1/jobs`: the merge queue, and the worker that drains it.
2//!
3//! A push whose base is stale, on a server with a worker enrolled, is stored
4//! at once (last-write-wins) and a `merge` job is queued holding both
5//! versions. The worker, `recall-worker`, is an enrolled device with the
6//! [`SCOPE_WORKER`](crate::devices::SCOPE_WORKER) scope and no inbound port:
7//! it long-polls [`CLAIM_PATH`], merges with the local `claude` CLI, and
8//! posts the result to [`result_path`]. Nothing a hook does ever waits on
9//! it; the merged file arrives with a later pull.
10//!
11//! A claim takes a **lease**: the job is the worker's until
12//! `lease_expires_at`, and a result counts only under the `lease_id` it was
13//! handed. An expired lease puts the job back in the queue, and the old
14//! holder's late result is refused with `409`, so a worker that stalled
15//! cannot overwrite the one that took over (a fencing token).
16//!
17//! Applying a merge is a compare-and-swap on the file's hash: if another
18//! push landed while the job ran, the merged content is not written over it
19//! but queued again against the new version, as a follow-up job.
20
21use serde::{Deserialize, Serialize};
22
23/// `POST`: wait for a job, and lease it. A worker device only.
24pub const CLAIM_PATH: &str = "/v1/jobs/claim";
25
26/// `GET`: jobs, newest first, without file content. Admin only; takes an
27/// optional `?state=`.
28pub const JOBS_PATH: &str = "/v1/jobs";
29
30/// `POST`: hand back the result of the job `id`, or the error it met.
31pub fn result_path(id: &str) -> String {
32 format!("{JOBS_PATH}/{id}/result")
33}
34
35/// `POST`: queue the failed job `id` again. Admin only.
36pub fn retry_path(id: &str) -> String {
37 format!("{JOBS_PATH}/{id}/retry")
38}
39
40/// Reconcile two versions of one file.
41pub const KIND_MERGE: &str = "merge";
42
43/// Look over memory and report what needs the owner's attention: see
44/// [`crate::evaluations`].
45pub const KIND_EVALUATE: &str = "evaluate";
46
47/// Waiting to be claimed, or to reach its `not_before` after a failed
48/// attempt.
49pub const STATE_QUEUED: &str = "queued";
50/// Claimed, under a lease that has not ended.
51pub const STATE_LEASED: &str = "leased";
52/// Finished: its result was applied, or superseded by a follow-up.
53pub const STATE_DONE: &str = "done";
54/// Out of attempts, or its result could not be applied. Kept, with any
55/// unapplied result, until someone retries it.
56pub const STATE_FAILED: &str = "failed";
57
58/// Every state a job can be in, in the order a job moves through them.
59pub const STATES: [&str; 4] = [STATE_QUEUED, STATE_LEASED, STATE_DONE, STATE_FAILED];
60
61/// The longest a claim may wait for a job, in seconds. Long enough that an
62/// idle worker costs about two requests a minute; short enough to stay well
63/// inside every proxy's idle timeout.
64pub const MAX_WAIT_SECONDS: u64 = 30;
65
66/// The shortest lease a claim may ask for, in seconds.
67pub const MIN_LEASE_SECONDS: u64 = 30;
68
69/// The longest lease a claim may ask for, in seconds.
70pub const MAX_LEASE_SECONDS: u64 = 600;
71
72fn default_lease_seconds() -> u64 {
73 120
74}
75
76/// Body of `POST /v1/jobs/claim`.
77#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
78pub struct ClaimRequest {
79 /// The kinds this worker can do, such as [`KIND_MERGE`]. Empty is
80 /// allowed: the claim then waits and answers `{"job": null}`, which is
81 /// how a worker whose `claude` CLI is not logged in still reports that
82 /// it is alive, and why, without taking jobs it would only fail.
83 pub kinds: Vec<String>,
84 /// How long to wait for a job before answering `{"job": null}`: 0 to
85 /// [`MAX_WAIT_SECONDS`].
86 #[serde(default)]
87 pub wait_seconds: u64,
88 /// How long the job is this worker's once claimed:
89 /// [`MIN_LEASE_SECONDS`] to [`MAX_LEASE_SECONDS`].
90 #[serde(default = "default_lease_seconds")]
91 pub lease_seconds: u64,
92 /// The worker's own check of its `claude` CLI. This is how `/health`
93 /// keeps reporting the CLI once it runs somewhere other than the API
94 /// process.
95 #[serde(default, skip_serializing_if = "Option::is_none")]
96 pub claude_cli: Option<ClaudeCliReport>,
97}
98
99/// A worker's check of its `claude` CLI, as it reports it with each claim.
100#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
101pub struct ClaudeCliReport {
102 /// When the check ran.
103 pub checked_at: String,
104 /// Whether the binary was found and runnable.
105 pub available: bool,
106 /// Whether it has a usable login.
107 pub logged_in: bool,
108 /// Why the check failed, when it did; empty otherwise.
109 #[serde(default)]
110 pub error: String,
111}
112
113/// `POST /v1/jobs/claim`'s answer: a leased job, or `null` when none
114/// arrived within `wait_seconds`.
115#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
116pub struct ClaimResponse {
117 /// The job, now leased to the caller.
118 pub job: Option<Job>,
119}
120
121/// One leased job.
122#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
123pub struct Job {
124 /// `job_…`.
125 pub id: String,
126 /// What to do, such as [`KIND_MERGE`]; the member of the same name
127 /// holds the input.
128 pub kind: String,
129 /// The fencing token: a result counts only under this.
130 pub lease_id: String,
131 /// When the lease ends and the job goes back to the queue.
132 pub lease_expires_at: String,
133 /// Which attempt this is, from 1.
134 pub attempt: u32,
135 /// The input of a [`KIND_MERGE`] job.
136 #[serde(default, skip_serializing_if = "Option::is_none")]
137 pub merge: Option<MergeInput>,
138 /// The input of a [`KIND_EVALUATE`] job.
139 #[serde(default, skip_serializing_if = "Option::is_none")]
140 pub evaluate: Option<EvaluateInput>,
141}
142
143/// What a merge job reconciles.
144#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
145pub struct MergeInput {
146 /// The project the file belongs to.
147 pub project_key: String,
148 /// The file.
149 pub file_path: String,
150 /// The version the push displaced.
151 pub stored: MergeSide,
152 /// The version the push stored. Its hash is the one the file must still
153 /// have for the merged result to be applied.
154 pub incoming: MergeSide,
155}
156
157/// One version of a file in a merge job.
158#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
159pub struct MergeSide {
160 /// [`content_sha256`](crate::content_sha256) of `content`.
161 pub sha256: String,
162 /// The file's exact bytes.
163 pub content: String,
164 /// The machine that wrote this version.
165 pub source_env: String,
166 /// When it was written.
167 pub updated_at: String,
168}
169
170/// What an evaluate job looks at, as the claim that leases it carries it.
171///
172/// The files come with the claim, read when it is made, rather than by the
173/// worker pulling each project: a worker then reads memory only while it
174/// holds an evaluation the owner asked for, and only the projects asked
175/// for, and the worker scope stays what it was, the job routes and nothing
176/// else.
177#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
178pub struct EvaluateInput {
179 /// The run this job makes the report for: `eval_…`.
180 pub evaluation_id: String,
181 /// The projects asked for; empty for every project.
182 pub projects: Vec<String>,
183 /// Whether to run the contradiction check, which asks `claude`.
184 pub contradictions: bool,
185 /// Every live file of the projects asked for (every project, when none
186 /// was named) and of every global scope, ordered by project and path.
187 /// Tombstones are left out.
188 #[serde(default)]
189 pub files: Vec<EvaluateFile>,
190}
191
192/// One file, as an evaluate job carries it.
193#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
194pub struct EvaluateFile {
195 /// The project, or scope, it belongs to.
196 pub project_key: String,
197 /// The file.
198 pub file_path: String,
199 /// Its exact bytes.
200 pub content: String,
201 /// When it was last written.
202 pub updated_at: String,
203}
204
205/// An evaluate job's result: see [`crate::evaluations`].
206#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
207pub struct EvaluateResult {
208 /// What was found, each holding nothing but enums, a file, lines and
209 /// related files; the server refuses a finding with any other key.
210 pub findings: Vec<crate::evaluations::Finding>,
211 /// Everything that quotes a note, as a
212 /// [`Details`](crate::evaluations::Details) object.
213 pub details: serde_json::Value,
214}
215
216/// Body of `POST /v1/jobs/{id}/result`: exactly one of `merge`, `evaluate`
217/// and `error`.
218#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
219pub struct ResultRequest {
220 /// The lease the job was claimed under.
221 pub lease_id: String,
222 /// A merge job's result.
223 #[serde(default, skip_serializing_if = "Option::is_none")]
224 pub merge: Option<MergeResult>,
225 /// Why the job could not be done. It is retried after 1, 5 and 30
226 /// minutes, then marked failed.
227 #[serde(default, skip_serializing_if = "Option::is_none")]
228 pub error: Option<String>,
229 /// An evaluate job's result.
230 #[serde(default, skip_serializing_if = "Option::is_none")]
231 pub evaluate: Option<EvaluateResult>,
232}
233
234impl ResultRequest {
235 /// How many of `merge`, `evaluate` and `error` it carries: exactly one
236 /// is a result the server records.
237 pub fn members(&self) -> usize {
238 [
239 self.merge.is_some(),
240 self.evaluate.is_some(),
241 self.error.is_some(),
242 ]
243 .into_iter()
244 .filter(|m| *m)
245 .count()
246 }
247}
248
249/// A merge job's result.
250#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
251pub struct MergeResult {
252 /// The merged file.
253 pub content: String,
254}
255
256/// `POST /v1/jobs/{id}/result`'s answer: what became of the job.
257#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
258pub struct ResultResponse {
259 /// The job.
260 pub id: String,
261 /// Its state now.
262 pub state: String,
263 /// Whether a merge result was written to the file. `false` for an
264 /// error, for a result the file had moved on from, and for an
265 /// evaluation, which writes no file.
266 pub applied: bool,
267 /// When the file changed while the job ran: the job that merges this
268 /// result with the newer version.
269 pub follow_up: Option<String>,
270}
271
272/// One job, as `GET /v1/jobs` lists it and as a retry answers. No file
273/// content: the listing is for seeing what the queue is doing.
274#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
275pub struct JobSummary {
276 /// `job_…`.
277 pub id: String,
278 /// Such as [`KIND_MERGE`].
279 pub kind: String,
280 /// One of [`STATES`].
281 pub state: String,
282 /// The project of the file it concerns.
283 pub project_key: String,
284 /// The file it concerns.
285 pub file_path: String,
286 /// How many times it has been handed out.
287 pub attempt: u32,
288 /// When it was queued.
289 pub created_at: String,
290 /// When its state last changed.
291 pub updated_at: String,
292 /// The last error it met, or why its result was not applied.
293 pub error: Option<String>,
294 /// The follow-up job its result was queued into, if any.
295 pub follow_up: Option<String>,
296}
297
298/// `GET /v1/jobs`.
299#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
300pub struct JobList {
301 /// At most 200, newest first.
302 pub jobs: Vec<JobSummary>,
303}
304
305#[cfg(test)]
306mod tests {
307 use super::*;
308
309 #[test]
310 fn a_claim_defaults_to_no_wait_and_a_two_minute_lease() {
311 let claim: ClaimRequest = serde_json::from_str(r#"{"kinds":["merge"]}"#).unwrap();
312 assert_eq!((claim.wait_seconds, claim.lease_seconds), (0, 120));
313 assert_eq!(claim.claude_cli, None);
314 }
315
316 /// An empty answer says `null`, not nothing, so a client reads the
317 /// absence of a job rather than guessing at it.
318 #[test]
319 fn no_job_is_null() {
320 assert_eq!(
321 serde_json::to_string(&ClaimResponse::default()).unwrap(),
322 r#"{"job":null}"#
323 );
324 }
325
326 #[test]
327 fn a_result_sends_only_the_member_it_has() {
328 let ok = ResultRequest {
329 lease_id: "lse_a".into(),
330 merge: Some(MergeResult {
331 content: "merged".into(),
332 }),
333 error: None,
334 evaluate: None,
335 };
336 assert_eq!(
337 serde_json::to_string(&ok).unwrap(),
338 r#"{"lease_id":"lse_a","merge":{"content":"merged"}}"#
339 );
340 let failed = ResultRequest {
341 lease_id: "lse_a".into(),
342 merge: None,
343 error: Some("claude merge timed out after 45s".into()),
344 evaluate: None,
345 };
346 assert_eq!(
347 serde_json::to_string(&failed).unwrap(),
348 r#"{"lease_id":"lse_a","error":"claude merge timed out after 45s"}"#
349 );
350 }
351
352 #[test]
353 fn paths() {
354 assert_eq!(result_path("job_a"), "/v1/jobs/job_a/result");
355 assert_eq!(retry_path("job_a"), "/v1/jobs/job_a/retry");
356 }
357}