Skip to main content

recall_worker/
worker.rs

1//! Enrolling, then the job loop: claim, merge, report, again.
2//!
3//! Before anything else the worker reads the server's discovery document,
4//! and goes no further unless the server lists the merge queue: a worker
5//! pointed at the wrong server, or one too old to have the queue, stops
6//! there, and never asks it to approve anything.
7//!
8//! The worker enrols by the device flow, like any machine, and never with
9//! an authkey: an authkey enrols `sync` devices only, so a leaked one
10//! cannot mint a worker. The owner approves its code with the `worker`
11//! scope, after comparing the fingerprint the worker prints. Its identity
12//! file records the server it was made for, and the worker refuses to run
13//! against any other.
14//!
15//! Each claim long-polls, so an idle worker costs about two requests a
16//! minute and a merge starts within a second of the push that queued it.
17//! Every claim also carries the worker's check of its `claude` CLI, which
18//! is how `/health` still reports the CLI once it has left the API process.
19//! A worker whose CLI is not logged in claims no merges: it stays visible,
20//! and takes no merge it would only fail. It still takes evaluations, whose
21//! checks but one need no CLI; the one that does, the contradiction check,
22//! is then skipped, and the report says why (see [`crate::evaluate`]).
23//!
24//! A refusal no retry will get past (a `4xx` other than `408` and `429`,
25//! or a server with no merge queue) is [`Fatal`]: the worker says what to
26//! do and stops asking. Anything else, such as a server that is restarting,
27//! is retried with a backoff.
28
29use std::future::Future;
30use std::time::{Duration, Instant};
31
32use recall_wire::devices::{
33    ACCESS_DENIED, AUTHORIZATION_PENDING, ENROLL_PATH, ENROLL_POLL_PATH, EXPIRED_TOKEN,
34    INVALID_GRANT, POLL_INTERVAL_SECONDS, SCOPE_WORKER, SLOW_DOWN,
35};
36use recall_wire::discovery::CAPABILITY_EVALUATION;
37use recall_wire::discovery::CAPABILITY_MERGE_QUEUE;
38use recall_wire::jobs::{self, KIND_EVALUATE, KIND_MERGE};
39use recall_wire::{
40    ClaimRequest, ClaimResponse, ClaudeCliReport, Discovery, EnrollPending, EnrollPollRequest,
41    EnrollPollResponse, EnrollRequest, EvaluateResult, Job, MergeResult, ResultRequest,
42    ResultResponse, DISCOVERY_PATH,
43};
44use reqwest::StatusCode;
45
46use crate::api::{Api, ApiError};
47use crate::config::Config;
48use crate::evaluate::{self, Settings, CONTRADICTION_TIMEOUT};
49use crate::identity::Identity;
50use crate::merge::{Merger, Status};
51
52/// How long an unsigned request (discovery, enrolment) may take.
53const UNSIGNED_TIMEOUT: Duration = Duration::from_secs(30);
54
55/// How long posting a result may take. A result carries one file.
56const RESULT_TIMEOUT: Duration = Duration::from_secs(60);
57
58/// How many times a result is posted before the worker gives up on it and
59/// lets the lease run out, which puts the job back in the queue.
60const RESULT_TRIES: u32 = 5;
61
62/// The longest the worker waits between failed attempts to reach the
63/// server.
64const MAX_BACKOFF: Duration = Duration::from_secs(60);
65
66/// How often the CLI is checked while it is not logged in, whatever the
67/// configured interval: logging in on the worker should be noticed within
68/// a minute, not half an hour.
69const NOT_LOGGED_IN_RECHECK: Duration = Duration::from_secs(60);
70
71/// For how long after a code is printed it is polled at the interval the
72/// server asks for. The owner approving it is most likely at a terminal in
73/// that first minute; after it, [`SLOW_POLL`] is plenty.
74const FAST_POLL_FOR: Duration = Duration::from_secs(60);
75
76/// How often a code nobody approved in its first minute is polled.
77const SLOW_POLL: Duration = Duration::from_secs(30);
78
79/// Why the worker stopped. Each is something only a person can fix, so the
80/// worker says what to do rather than retrying forever.
81#[derive(Debug, thiserror::Error)]
82pub enum Fatal {
83    /// The identity file could not be read or written.
84    #[error("{0}")]
85    Identity(String),
86    /// The identity was made for another server.
87    #[error(
88        "{path} was made for {made_for}, not {configured}, and a worker's key is never \
89         offered to a second server. If RECALL_WORKER_SERVER is wrong, correct it; to move \
90         this worker to {configured}, revoke it on {made_for}, delete {path} and restart"
91    )]
92    OtherServer {
93        /// The server the identity records.
94        made_for: String,
95        /// `RECALL_WORKER_SERVER`.
96        configured: String,
97        /// The identity file.
98        path: String,
99    },
100    /// An identity an earlier build wrote, part way through enrolling
101    /// with a server it did not record.
102    #[error(
103        "{0} was written by an earlier recall-worker that did not record which server it \
104         enrolled with, so it cannot be checked against RECALL_WORKER_SERVER. Revoke the \
105         device it enrolled as (or deny its code), delete {0} and restart"
106    )]
107    Unbound(String),
108    /// The server is not one this worker can work for.
109    #[error("{0}")]
110    NotRecall(String),
111    /// The owner denied the enrolment.
112    #[error("the owner denied this worker's enrolment. Delete {0} and restart to ask again")]
113    Denied(String),
114    /// The owner approved it with a scope other than `worker`.
115    #[error(
116        "this device was approved with the {0} scope, not worker, so it cannot claim jobs. \
117         Revoke it, delete {1} and restart, then approve the new code as a worker \
118         (recall devices approve <code> --worker)"
119    )]
120    WrongScope(String, String),
121    /// The server no longer knows this device, or revoked it.
122    #[error(
123        "the server refused this worker ({0}). If it was revoked on purpose, nothing is wrong; \
124         to enrol again, delete {1} and restart"
125    )]
126    Refused(String, String),
127    /// Something else that will not fix itself, such as a name another
128    /// device already has.
129    #[error("{0}")]
130    Other(String),
131}
132
133/// What one claim came to.
134#[derive(Debug, Clone, PartialEq, Eq)]
135pub enum Step {
136    /// No job arrived within the wait.
137    Idle,
138    /// A job was done and its result recorded.
139    Done(ResultResponse),
140    /// A job was done, and the server would not take the result: its
141    /// lease ended, or the job is gone.
142    Rejected(String),
143}
144
145/// A worker: its settings, identity and CLI.
146pub struct Worker {
147    cfg: Config,
148    api: Api,
149    id: Identity,
150    merger: Merger,
151    status: Status,
152    status_at: Option<Instant>,
153    /// The last code the approval instructions were printed for.
154    announced: Option<String>,
155    /// Whether the server makes evaluation reports, and so has `evaluate`
156    /// jobs to claim: a server older than 0.4.5 would refuse the kind.
157    evaluations: bool,
158}
159
160fn log(message: &str) {
161    eprintln!("recall-worker: {message}");
162}
163
164/// Whether an answer is one that asking again will not change: a redirect
165/// (never followed), or a client error other than a timeout or a rate
166/// limit.
167fn is_final(e: &ApiError) -> bool {
168    e.status().is_some_and(|s| {
169        s.is_redirection()
170            || (s.is_client_error()
171                && s != StatusCode::REQUEST_TIMEOUT
172                && s != StatusCode::TOO_MANY_REQUESTS)
173    })
174}
175
176/// How long to wait before asking again after `e`: what the server asked
177/// for, or the next step of a backoff that doubles up to [`MAX_BACKOFF`].
178fn next_wait(e: &ApiError, backoff: Duration) -> Duration {
179    match e {
180        ApiError::Status {
181            retry_after: Some(s),
182            ..
183        } => Duration::from_secs(*s),
184        _ => (backoff * 2).clamp(Duration::from_secs(1), MAX_BACKOFF),
185    }
186}
187
188impl Worker {
189    /// A worker for `cfg`, with the identity kept in its data directory,
190    /// made there if this is the first start. Refuses an identity made for
191    /// another server.
192    pub fn new(cfg: Config) -> Result<Self, Fatal> {
193        if cfg.data_dir.as_os_str().is_empty() {
194            return Err(Fatal::Other(
195                "no data directory: set RECALL_WORKER_DIR to where the worker's key is kept".into(),
196            ));
197        }
198        let path = Identity::path(&cfg.data_dir).display().to_string();
199        let mut id = Identity::load_or_create(&cfg.data_dir, &cfg.server)
200            .map_err(|e| Fatal::Identity(format!("cannot use {path}: {e}")))?;
201        match id.server.as_deref() {
202            Some(server) if server == cfg.server => {}
203            Some(server) => {
204                return Err(Fatal::OtherServer {
205                    made_for: server.to_string(),
206                    configured: cfg.server.clone(),
207                    path,
208                })
209            }
210            // An earlier build's file. With nothing enrolled yet, the key
211            // has been offered to no server, and is this one's from now on.
212            None if id.device_id.is_none() && id.enrollment_id.is_none() => {
213                id.server = Some(cfg.server.clone());
214                id.save(&cfg.data_dir)
215                    .map_err(|e| Fatal::Identity(format!("cannot write {path}: {e}")))?;
216            }
217            None => return Err(Fatal::Unbound(path)),
218        }
219        let api = Api::new(&cfg.server).map_err(|e| Fatal::Other(e.to_string()))?;
220        let merger = Merger::new(cfg.claude_bin.clone(), cfg.merge_timeout);
221        Ok(Self {
222            cfg,
223            api,
224            id,
225            merger,
226            status: Status::default(),
227            status_at: None,
228            announced: None,
229            evaluations: false,
230        })
231    }
232
233    /// Its identity: key fingerprint, and device id once approved.
234    pub fn identity(&self) -> &Identity {
235        &self.id
236    }
237
238    fn identity_path(&self) -> String {
239        Identity::path(&self.cfg.data_dir).display().to_string()
240    }
241
242    fn save(&self) -> Result<(), Fatal> {
243        self.id
244            .save(&self.cfg.data_dir)
245            .map_err(|e| Fatal::Identity(format!("cannot write {}: {e}", self.identity_path())))
246    }
247
248    /// Checks the server, enrols until the owner approves, then runs jobs
249    /// until `shutdown` resolves. Returns only on shutdown, or on something
250    /// a person has to fix.
251    pub async fn run(mut self, shutdown: impl Future<Output = ()>) -> Result<(), Fatal> {
252        tokio::pin!(shutdown);
253        log(&format!(
254            "{} for {}, key fingerprint {}",
255            crate::user_agent(),
256            self.cfg.server,
257            self.id.fingerprint()
258        ));
259        if let Some(warning) = self.cfg.plaintext_warning() {
260            log(&warning);
261        }
262        tokio::select! {
263            _ = &mut shutdown => return Ok(()),
264            checked = self.check_server() => checked?,
265        }
266        if self.id.device_id.is_none() {
267            tokio::select! {
268                _ = &mut shutdown => return Ok(()),
269                enrolled = self.enrol() => enrolled?,
270            }
271        }
272        log(&format!(
273            "enrolled as {}; waiting for jobs",
274            self.id.device_id.as_deref().unwrap_or("")
275        ));
276        let mut backoff = Duration::ZERO;
277        loop {
278            if !backoff.is_zero() {
279                tokio::select! {
280                    _ = &mut shutdown => return Ok(()),
281                    _ = tokio::time::sleep(backoff) => {}
282                }
283            }
284            let step = tokio::select! {
285                _ = &mut shutdown => return Ok(()),
286                step = self.step() => step,
287            };
288            backoff = match step {
289                Ok(Step::Idle) => Duration::ZERO,
290                Ok(Step::Done(outcome)) => {
291                    log(&format!(
292                        "job {}: {}{}",
293                        outcome.id,
294                        outcome.state,
295                        match (outcome.applied, &outcome.follow_up) {
296                            (true, _) => ", merged file stored".to_string(),
297                            (false, Some(next)) => {
298                                format!(", the file changed meanwhile; follow-up {next}")
299                            }
300                            (false, None) => String::new(),
301                        }
302                    ));
303                    Duration::ZERO
304                }
305                Ok(Step::Rejected(why)) => {
306                    log(&why);
307                    Duration::ZERO
308                }
309                Err(e) => {
310                    if let Some(fatal) = self.fatal(&e) {
311                        return Err(fatal);
312                    }
313                    // Not a refusal of this device but of the claim itself,
314                    // such as the 404 of a server rolled back to a release
315                    // from before the queue. A 401 is left to retry: it is
316                    // also what a signature made just before a restart gets.
317                    if is_final(&e) && e.status() != Some(StatusCode::UNAUTHORIZED) {
318                        return Err(Fatal::NotRecall(format!(
319                            "{} refused a claim ({e}), and asking again will not change that. \
320                             If the server was rolled back to a release without the merge \
321                             queue, restart the worker once it has the queue again",
322                            self.cfg.server
323                        )));
324                    }
325                    let wait = next_wait(&e, backoff);
326                    log(&format!("claim failed ({e}); trying again in {wait:?}"));
327                    wait
328                }
329            };
330        }
331    }
332
333    /// Whether a refusal is of this device, which no retry will get past.
334    fn fatal(&self, e: &ApiError) -> Option<Fatal> {
335        match e.status()? {
336            // A 401 is also what a signature made a moment before the
337            // server restarted gets, which the next request fixes; only a
338            // device the server does not know, or revoked, is final.
339            StatusCode::UNAUTHORIZED
340                if e.message().contains("unknown device") || e.message().contains("revoked") =>
341            {
342                Some(Fatal::Refused(
343                    e.message().to_string(),
344                    self.identity_path(),
345                ))
346            }
347            StatusCode::FORBIDDEN => Some(Fatal::Refused(
348                e.message().to_string(),
349                self.identity_path(),
350            )),
351            _ => None,
352        }
353    }
354
355    /// Reads the server's discovery document, and goes on only if the
356    /// server lists the merge queue. An answer that says this is not such a
357    /// server is final; one that says nothing (the server is restarting, or
358    /// not up yet) is asked again.
359    pub async fn check_server(&mut self) -> Result<(), Fatal> {
360        let server = &self.cfg.server;
361        let mut backoff = Duration::ZERO;
362        loop {
363            let answer: Result<Discovery, ApiError> =
364                self.api.get(DISCOVERY_PATH, UNSIGNED_TIMEOUT).await;
365            match answer {
366                Ok(doc) if doc.can(CAPABILITY_MERGE_QUEUE) => {
367                    self.evaluations = doc.can(CAPABILITY_EVALUATION);
368                    return Ok(());
369                }
370                Ok(doc) => {
371                    return Err(Fatal::NotRecall(format!(
372                        "{server} is Recall {}, which has no merge queue, so there is nothing \
373                         for a worker to do. Upgrade it, or stop the worker",
374                        doc.server.version
375                    )))
376                }
377                Err(ApiError::Body(why)) => {
378                    return Err(Fatal::NotRecall(format!(
379                        "{server} did not answer {DISCOVERY_PATH} with Recall's discovery \
380                         document ({why}); RECALL_WORKER_SERVER must name a Recall server"
381                    )))
382                }
383                Err(e) if is_final(&e) => {
384                    return Err(Fatal::NotRecall(format!(
385                        "{server} answered {DISCOVERY_PATH} with {e}: it is not a Recall server, \
386                         or one older than 0.4.1, which has no merge queue. Check \
387                         RECALL_WORKER_SERVER"
388                    )))
389                }
390                Err(e) => {
391                    backoff = next_wait(&e, backoff);
392                    log(&format!(
393                        "cannot reach {server} ({e}); trying again in {backoff:?}"
394                    ));
395                    tokio::time::sleep(backoff).await;
396                }
397            }
398        }
399    }
400
401    /// Enrols by the device flow: asks for a code, prints it with the key
402    /// fingerprint, and polls until the owner approves it. A restart while
403    /// waiting keeps polling the same enrolment.
404    pub async fn enrol(&mut self) -> Result<(), Fatal> {
405        loop {
406            match self.id.user_code.clone() {
407                Some(code) if self.id.enrollment_id.is_some() => self.announce(&code),
408                _ => self.start_enrolment().await?,
409            }
410            match self.wait_for_approval().await? {
411                Some(device_id) => {
412                    self.id.device_id = Some(device_id);
413                    self.id.enrollment_id = None;
414                    self.id.user_code = None;
415                    self.save()?;
416                    return Ok(());
417                }
418                // The code expired before anyone approved it: ask for a
419                // new one.
420                None => {
421                    self.id.enrollment_id = None;
422                    self.id.user_code = None;
423                    self.save()?;
424                }
425            }
426        }
427    }
428
429    async fn start_enrolment(&mut self) -> Result<(), Fatal> {
430        let req = EnrollRequest {
431            name: self.cfg.name.clone(),
432            public_key: self.id.public_key(),
433            agent: crate::user_agent(),
434            authkey: None,
435        };
436        let mut backoff = Duration::ZERO;
437        let pending: EnrollPending = loop {
438            match self.api.post(ENROLL_PATH, &req, UNSIGNED_TIMEOUT).await {
439                Ok(pending) => break pending,
440                Err(e) if e.status() == Some(StatusCode::CONFLICT) => {
441                    return Err(Fatal::Other(format!(
442                        "{}. Set RECALL_WORKER_NAME to another name, or revoke the old worker",
443                        e.message()
444                    )))
445                }
446                // A 404 among them: a server without the enrolment route
447                // answers every retry the same way.
448                Err(e) if is_final(&e) => {
449                    return Err(Fatal::Other(format!(
450                        "{} refused the enrolment ({e}); asking again would be refused \
451                         the same way",
452                        self.cfg.server
453                    )))
454                }
455                Err(e) => {
456                    backoff = next_wait(&e, backoff);
457                    log(&format!(
458                        "cannot reach the server to enrol ({e}); trying again in {backoff:?}"
459                    ));
460                    tokio::time::sleep(backoff).await;
461                }
462            }
463        };
464        self.id.enrollment_id = Some(pending.enrollment_id);
465        self.id.user_code = Some(pending.user_code.clone());
466        self.save()?;
467        self.announce(&pending.user_code);
468        Ok(())
469    }
470
471    /// Says what the owner has to do, once per code.
472    fn announce(&mut self, code: &str) {
473        if self.announced.as_deref() == Some(code) {
474            return;
475        }
476        self.announced = Some(code.to_string());
477        log(&format!(
478            "waiting for approval of code {code} as a worker, key fingerprint {}",
479            self.id.fingerprint()
480        ));
481        log(&format!(
482            "approve it from an admin device: recall devices approve {code} --worker \
483             --fingerprint {}",
484            self.id.fingerprint()
485        ));
486        log(&format!(
487            "or with the operator token: POST /v1/devices/approve \
488             {{\"user_code\":\"{code}\",\"scope\":\"worker\",\"fingerprint\":\"{}\"}} \
489             (see deploy/README.md, \"The merge worker\")",
490            self.id.fingerprint()
491        ));
492    }
493
494    /// How long to wait before the next poll of a code printed `since`
495    /// ago, `slowed` being what the server's `slow_down` answers added.
496    fn poll_interval(since: Duration, slowed: Duration) -> Duration {
497        let base = if since < FAST_POLL_FOR {
498            Duration::from_secs(POLL_INTERVAL_SECONDS)
499        } else {
500            SLOW_POLL
501        };
502        base + slowed
503    }
504
505    /// Polls until approved (`Some(device_id)`), or until the code expires
506    /// (`None`).
507    async fn wait_for_approval(&mut self) -> Result<Option<String>, Fatal> {
508        let Some(enrollment_id) = self.id.enrollment_id.clone() else {
509            return Ok(None);
510        };
511        let since = Instant::now();
512        let mut slowed = Duration::ZERO;
513        let req = EnrollPollRequest { enrollment_id };
514        loop {
515            tokio::time::sleep(Self::poll_interval(since.elapsed(), slowed)).await;
516            let polled: Result<EnrollPollResponse, ApiError> = self
517                .api
518                .post(ENROLL_POLL_PATH, &req, UNSIGNED_TIMEOUT)
519                .await;
520            match polled {
521                Ok(approved) if approved.scope == SCOPE_WORKER => {
522                    return Ok(Some(approved.device_id))
523                }
524                Ok(approved) => {
525                    // Kept, so the owner can find and revoke it; the worker
526                    // itself will not run as anything but a worker.
527                    self.id.device_id = Some(approved.device_id);
528                    self.save()?;
529                    return Err(Fatal::WrongScope(approved.scope, self.identity_path()));
530                }
531                Err(e) if e.status() == Some(StatusCode::BAD_REQUEST) => match e.message() {
532                    AUTHORIZATION_PENDING => {}
533                    SLOW_DOWN => slowed += Duration::from_secs(5),
534                    EXPIRED_TOKEN | INVALID_GRANT => {
535                        log("the code expired before it was approved; asking for a new one");
536                        return Ok(None);
537                    }
538                    ACCESS_DENIED => return Err(Fatal::Denied(self.identity_path())),
539                    other => log(&format!("poll refused: {other}")),
540                },
541                Err(e) if is_final(&e) => {
542                    return Err(Fatal::Other(format!(
543                        "{} refused the poll for this worker's approval ({e})",
544                        self.cfg.server
545                    )))
546                }
547                Err(e) => log(&format!("poll failed ({e}); trying again")),
548            }
549        }
550    }
551
552    /// Re-checks the CLI when the last check is old enough, or when a
553    /// merge has just failed.
554    async fn refresh_status(&mut self) {
555        let every = if self.status.logged_in {
556            self.cfg.claude_status_interval
557        } else {
558            self.cfg.claude_status_interval.min(NOT_LOGGED_IN_RECHECK)
559        };
560        if self.status_at.is_some_and(|at| at.elapsed() < every) {
561            return;
562        }
563        // Never checked before: whatever it finds is news.
564        let (was, first) = (self.status.logged_in, self.status.checked_at.is_empty());
565        self.status = self.merger.check_status().await;
566        self.status_at = Some(Instant::now());
567        if first || self.status.logged_in != was {
568            log(&if self.status.logged_in {
569                "the claude CLI is logged in; taking merge jobs".to_string()
570            } else {
571                format!(
572                    "the claude CLI cannot merge ({}); taking no jobs until it can. \
573                     Log it in with: docker exec -it -u node recall-worker claude setup-token",
574                    if self.status.error.is_empty() {
575                        "not logged in"
576                    } else {
577                        &self.status.error
578                    }
579                )
580            });
581        }
582    }
583
584    /// The claim this worker sends now: every kind it can do, and that
585    /// CLI's last check. No merges while its CLI cannot merge; evaluations
586    /// whenever the server makes them, since all their checks but the
587    /// contradiction check run without the CLI.
588    pub fn claim_request(&self) -> ClaimRequest {
589        let mut kinds = Vec::new();
590        if self.status.logged_in {
591            kinds.push(KIND_MERGE.to_string());
592        }
593        if self.evaluations {
594            kinds.push(KIND_EVALUATE.to_string());
595        }
596        ClaimRequest {
597            kinds,
598            wait_seconds: self.cfg.wait_seconds,
599            lease_seconds: self.cfg.lease_seconds,
600            claude_cli: Some(ClaudeCliReport {
601                checked_at: self.status.checked_at.clone(),
602                available: self.status.available,
603                logged_in: self.status.logged_in,
604                error: self.status.error.clone(),
605            }),
606        }
607    }
608
609    /// One claim, and the job it brought, if any.
610    pub async fn step(&mut self) -> Result<Step, ApiError> {
611        self.refresh_status().await;
612        let device_id = self.id.device_id.clone().unwrap_or_default();
613        let claim = self.claim_request();
614        let answer: ClaimResponse = self
615            .api
616            .post_signed(
617                jobs::CLAIM_PATH,
618                &claim,
619                self.id.key(),
620                &device_id,
621                Duration::from_secs(self.cfg.wait_seconds + 30),
622            )
623            .await?;
624        let Some(job) = answer.job else {
625            return Ok(Step::Idle);
626        };
627        let result = self.work(&job).await;
628        self.report(&job, &result).await
629    }
630
631    /// Does one job, and says what to report.
632    ///
633    /// A merge that fails has the CLI checked again before the next claim,
634    /// rather than at the next scheduled check: a CLI that was logged out
635    /// under the worker should stop it taking jobs now, not in half an hour.
636    pub async fn work(&mut self, job: &Job) -> ResultRequest {
637        if let (KIND_EVALUATE, Some(input)) = (job.kind.as_str(), &job.evaluate) {
638            return self.evaluate(job, input).await;
639        }
640        let outcome = match (job.kind.as_str(), &job.merge) {
641            (KIND_MERGE, Some(m)) => {
642                if recall_wire::content_sha256(&m.stored.content) != m.stored.sha256
643                    || recall_wire::content_sha256(&m.incoming.content) != m.incoming.sha256
644                {
645                    Err("a version's content does not match its sha256".to_string())
646                } else if m.stored.content == m.incoming.content {
647                    // Identical inputs never reach claude: there is nothing
648                    // to reconcile, and a call would cost money to say so.
649                    Ok(m.incoming.content.clone())
650                } else {
651                    log(&format!(
652                        "job {}: merging {}/{} (attempt {})",
653                        job.id, m.project_key, m.file_path, job.attempt
654                    ));
655                    let merged = self
656                        .merger
657                        .merge(&m.stored.content, &m.incoming.content)
658                        .await;
659                    if merged.is_err() {
660                        self.status_at = None;
661                    }
662                    merged.map_err(|e| e.to_string())
663                }
664            }
665            (kind, _) => Err(format!("this worker does not do {kind} jobs")),
666        };
667        match outcome {
668            Ok(content) => ResultRequest {
669                lease_id: job.lease_id.clone(),
670                merge: Some(MergeResult { content }),
671                error: None,
672                evaluate: None,
673            },
674            Err(error) => {
675                log(&format!("job {}: {error}", job.id));
676                ResultRequest {
677                    lease_id: job.lease_id.clone(),
678                    merge: None,
679                    error: Some(error),
680                    evaluate: None,
681                }
682            }
683        }
684    }
685
686    /// Makes an evaluation's report. The contradiction check runs only
687    /// when the job asks for it, and only while the CLI is logged in; each
688    /// call it makes must end well before the lease does.
689    async fn evaluate(&mut self, job: &Job, input: &recall_wire::EvaluateInput) -> ResultRequest {
690        log(&format!(
691            "job {}: evaluation {} of {} files{} (attempt {})",
692            job.id,
693            input.evaluation_id,
694            input.files.len(),
695            if input.contradictions {
696                ", with the contradiction check"
697            } else {
698                ""
699            },
700            job.attempt
701        ));
702        let now = time::OffsetDateTime::now_utc();
703        let deadline = time::OffsetDateTime::parse(
704            &job.lease_expires_at,
705            &time::format_description::well_known::Rfc3339,
706        )
707        .ok()
708        .map(|end| {
709            let left = (end - now).max(time::Duration::ZERO);
710            Instant::now() + Duration::try_from(left).unwrap_or_default()
711        });
712        let settings = Settings {
713            now,
714            stale_after: Duration::from_secs(self.cfg.eval_stale_days.saturating_mul(24 * 60 * 60)),
715            deadline,
716            cli_unavailable: (!self.status.logged_in).then(|| {
717                if self.status.error.is_empty() {
718                    "not logged in".to_string()
719                } else {
720                    self.status.error.clone()
721                }
722            }),
723        };
724        let claude = Merger::new(self.cfg.claude_bin.clone(), CONTRADICTION_TIMEOUT);
725        let report = evaluate::evaluate(input, &settings, &claude).await;
726        log(&format!(
727            "job {}: {} findings {:?}{}",
728            job.id,
729            report.findings.len(),
730            evaluate::counts(&report),
731            if report.details.skipped.is_empty() {
732                String::new()
733            } else {
734                format!(
735                    ", {} checks skipped (see the report)",
736                    report.details.skipped.len()
737                )
738            }
739        ));
740        let details = serde_json::to_value(&report.details).unwrap_or_default();
741        ResultRequest {
742            lease_id: job.lease_id.clone(),
743            merge: None,
744            error: None,
745            evaluate: Some(EvaluateResult {
746                findings: report.findings,
747                details,
748            }),
749        }
750    }
751
752    /// Posts a result, trying again when the server could not be reached.
753    /// Posting the same result twice is safe: the server records it once.
754    async fn report(&self, job: &Job, result: &ResultRequest) -> Result<Step, ApiError> {
755        let device_id = self.id.device_id.clone().unwrap_or_default();
756        let mut result = result.clone();
757        let mut wait = Duration::from_secs(1);
758        let mut last = None;
759        let mut tries = 0;
760        while tries < RESULT_TRIES {
761            tries += 1;
762            let posted: Result<ResultResponse, ApiError> = self
763                .api
764                .post_signed(
765                    &jobs::result_path(&job.id),
766                    &result,
767                    self.id.key(),
768                    &device_id,
769                    RESULT_TIMEOUT,
770                )
771                .await;
772            match posted {
773                Ok(outcome) => return Ok(Step::Done(outcome)),
774                // A result the server will not take as it is, too large or
775                // not of the shape it takes: asking again would be answered
776                // the same way, and a worker that stopped over one job
777                // would stop merging too. The job gets an error instead,
778                // and is retried later, or fails, as any job does.
779                Err(e)
780                    if result.error.is_none()
781                        && matches!(
782                            e.status(),
783                            Some(StatusCode::PAYLOAD_TOO_LARGE | StatusCode::BAD_REQUEST)
784                        ) =>
785                {
786                    let why = if e.status() == Some(StatusCode::PAYLOAD_TOO_LARGE) {
787                        format!(
788                            "the result came to {} bytes, more than the server takes",
789                            serde_json::to_vec(&result).map_or(0, |b| b.len())
790                        )
791                    } else {
792                        let said: String = e.message().chars().take(200).collect();
793                        format!("the server refused the result: {said}")
794                    };
795                    log(&format!("job {}: {why}; reporting that instead", job.id));
796                    result = ResultRequest {
797                        lease_id: result.lease_id.clone(),
798                        merge: None,
799                        error: Some(why),
800                        evaluate: None,
801                    };
802                    tries = 0;
803                }
804                Err(e)
805                    if matches!(
806                        e.status(),
807                        Some(
808                            StatusCode::CONFLICT
809                                | StatusCode::NOT_FOUND
810                                | StatusCode::BAD_REQUEST
811                                | StatusCode::PAYLOAD_TOO_LARGE
812                        )
813                    ) =>
814                {
815                    return Ok(Step::Rejected(format!(
816                        "job {}: the server did not take the result: {}",
817                        job.id,
818                        e.message()
819                    )))
820                }
821                Err(e) if self.fatal(&e).is_some() => return Err(e),
822                Err(e) => {
823                    log(&format!(
824                        "job {}: posting the result failed ({e}); trying again in {wait:?}",
825                        job.id
826                    ));
827                    last = Some(e);
828                    tokio::time::sleep(wait).await;
829                    wait *= 2;
830                }
831            }
832        }
833        Err(last.unwrap_or_else(|| ApiError::Transport("no attempt was made".into())))
834    }
835}
836
837#[cfg(test)]
838mod tests {
839    use super::*;
840    use recall_wire::{MergeInput, MergeSide};
841    use tokio::io::{AsyncReadExt, AsyncWriteExt};
842
843    fn config(dir: &std::path::Path, server: &str) -> Config {
844        Config {
845            server: server.to_string(),
846            data_dir: dir.to_path_buf(),
847            // Never run: nothing in these tests may reach a real CLI.
848            claude_bin: "definitely-not-a-real-claude".into(),
849            ..Config::default()
850        }
851    }
852
853    fn side(content: &str) -> MergeSide {
854        MergeSide {
855            sha256: recall_wire::content_sha256(content),
856            content: content.to_string(),
857            source_env: "laptop".into(),
858            updated_at: "2026-10-02T09:10:11.020Z".into(),
859        }
860    }
861
862    fn job(stored: MergeSide, incoming: MergeSide) -> Job {
863        Job {
864            id: "job_1".into(),
865            kind: KIND_MERGE.into(),
866            lease_id: "lse_1".into(),
867            lease_expires_at: "2026-10-02T09:16:03.118Z".into(),
868            attempt: 1,
869            merge: Some(MergeInput {
870                project_key: "acme/app".into(),
871                file_path: "topics/auth.md".into(),
872                stored,
873                incoming,
874            }),
875            evaluate: None,
876        }
877    }
878
879    /// A worker for a server nothing listens on: nothing here sends a
880    /// request.
881    fn worker() -> (tempfile::TempDir, Worker) {
882        let dir = tempfile::tempdir().unwrap();
883        let w = Worker::new(config(dir.path(), "http://127.0.0.1:9")).unwrap();
884        (dir, w)
885    }
886
887    /// Either side's content must match the hash the server sent with it,
888    /// the stored side as much as the incoming one, or the worker merges
889    /// nothing and says why.
890    #[tokio::test]
891    async fn a_version_that_does_not_match_its_hash_is_not_merged() {
892        let (_dir, mut w) = worker();
893        let mut stored = side("A");
894        stored.sha256 = recall_wire::content_sha256("not A");
895        let mut incoming = side("B");
896        incoming.sha256 = recall_wire::content_sha256("not B");
897        for (s, i) in [
898            (stored.clone(), side("B")),
899            (side("A"), incoming.clone()),
900            // Identical but for the hash: without the check this would be
901            // answered as a merge without asking claude at all.
902            (stored.clone(), {
903                let mut same = side("A");
904                same.sha256 = stored.sha256.clone();
905                same
906            }),
907        ] {
908            let result = w.work(&job(s, i)).await;
909            assert!(result.merge.is_none(), "{result:?}");
910            assert!(
911                result.error.as_deref().unwrap().contains("sha256"),
912                "{result:?}"
913            );
914            assert_eq!(result.lease_id, "lse_1");
915        }
916        // Matching hashes and the same content: answered without claude.
917        let result = w.work(&job(side("A"), side("A"))).await;
918        assert_eq!(result.merge.unwrap().content, "A");
919    }
920
921    /// A worker whose CLI cannot merge asks for no kind of job, so it is
922    /// never handed one it would only fail.
923    #[test]
924    fn a_worker_that_cannot_merge_claims_nothing() {
925        let (_dir, mut w) = worker();
926        assert!(w.claim_request().kinds.is_empty(), "never checked");
927        w.status = Status {
928            checked_at: "2026-10-02T09:13:40.002Z".into(),
929            available: true,
930            logged_in: false,
931            error: "not logged in".into(),
932        };
933        let claim = w.claim_request();
934        assert!(claim.kinds.is_empty());
935        assert_eq!(claim.claude_cli.unwrap().error, "not logged in");
936        w.status.logged_in = true;
937        assert_eq!(w.claim_request().kinds, vec![KIND_MERGE.to_string()]);
938    }
939
940    /// A server that makes evaluation reports gets evaluate claims, from a
941    /// worker whose CLI cannot merge as well: every check but the
942    /// contradiction check runs without it.
943    #[test]
944    fn a_worker_takes_evaluations_whenever_the_server_makes_them() {
945        let (_dir, mut w) = worker();
946        w.evaluations = true;
947        assert_eq!(w.claim_request().kinds, vec![KIND_EVALUATE.to_string()]);
948        w.status.logged_in = true;
949        assert_eq!(
950            w.claim_request().kinds,
951            vec![KIND_MERGE.to_string(), KIND_EVALUATE.to_string()]
952        );
953    }
954
955    /// A merge that fails has the CLI checked again before the next claim.
956    #[tokio::test]
957    async fn a_failed_merge_has_the_cli_checked_again() {
958        let (_dir, mut w) = worker();
959        w.status.logged_in = true;
960        w.status_at = Some(Instant::now());
961        let result = w.work(&job(side("A"), side("B"))).await;
962        assert!(result.error.unwrap().contains("unavailable"));
963        assert!(w.status_at.is_none());
964    }
965
966    #[test]
967    fn a_code_is_polled_quickly_for_a_minute_then_slowly() {
968        let fast = Duration::from_secs(POLL_INTERVAL_SECONDS);
969        assert_eq!(Worker::poll_interval(Duration::ZERO, Duration::ZERO), fast);
970        assert_eq!(
971            Worker::poll_interval(Duration::from_secs(59), Duration::from_secs(5)),
972            fast + Duration::from_secs(5)
973        );
974        assert_eq!(
975            Worker::poll_interval(Duration::from_secs(61), Duration::ZERO),
976            SLOW_POLL
977        );
978        assert!(SLOW_POLL >= Duration::from_secs(30));
979    }
980
981    /// A key enrolled with one server is never offered to another.
982    #[test]
983    fn an_identity_is_bound_to_its_server() {
984        let dir = tempfile::tempdir().unwrap();
985        Worker::new(config(dir.path(), "http://recall-server:8787")).unwrap();
986        match Worker::new(config(dir.path(), "https://recall.example.com")) {
987            Err(Fatal::OtherServer {
988                made_for,
989                configured,
990                ..
991            }) => {
992                assert_eq!(made_for, "http://recall-server:8787");
993                assert_eq!(configured, "https://recall.example.com");
994            }
995            other => panic!("not refused: {:?}", other.err()),
996        }
997        assert!(Worker::new(config(dir.path(), "http://recall-server:8787")).is_ok());
998    }
999
1000    /// A file an earlier build wrote records no server. Unused, it is bound
1001    /// to this one; part way through an enrolment, it is refused.
1002    #[test]
1003    fn an_identity_from_before_servers_were_recorded() {
1004        let dir = tempfile::tempdir().unwrap();
1005        let mut id = Identity::from_seed([3; 32]);
1006        id.save(dir.path()).unwrap();
1007        Worker::new(config(dir.path(), "http://recall-server:8787")).unwrap();
1008        assert_eq!(
1009            Identity::load(dir.path())
1010                .unwrap()
1011                .unwrap()
1012                .server
1013                .as_deref(),
1014            Some("http://recall-server:8787")
1015        );
1016
1017        let dir = tempfile::tempdir().unwrap();
1018        id.device_id = Some("dev_somewhere".into());
1019        id.save(dir.path()).unwrap();
1020        assert!(matches!(
1021            Worker::new(config(dir.path(), "http://recall-server:8787")),
1022            Err(Fatal::Unbound(_))
1023        ));
1024    }
1025
1026    #[test]
1027    fn a_worker_needs_a_data_directory() {
1028        assert!(matches!(
1029            Worker::new(Config {
1030                server: "http://127.0.0.1:9".into(),
1031                ..Config::default()
1032            }),
1033            Err(Fatal::Other(_))
1034        ));
1035    }
1036
1037    /// A stand-in server on loopback that answers each request with what
1038    /// `answer` returns for its method and path, and counts the requests.
1039    async fn fake_server(
1040        answer: fn(&str, &str) -> (u16, String),
1041    ) -> (String, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
1042        use std::sync::atomic::{AtomicUsize, Ordering};
1043        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1044        let url = format!("http://{}", listener.local_addr().unwrap());
1045        let count = std::sync::Arc::new(AtomicUsize::new(0));
1046        let counted = count.clone();
1047        tokio::spawn(async move {
1048            while let Ok((mut conn, _)) = listener.accept().await {
1049                let counted = counted.clone();
1050                tokio::spawn(async move {
1051                    let mut buf = Vec::new();
1052                    let mut chunk = [0u8; 4096];
1053                    // The head, then as much body as it says there is.
1054                    let head_end = loop {
1055                        let n = conn.read(&mut chunk).await.unwrap_or(0);
1056                        if n == 0 {
1057                            return;
1058                        }
1059                        buf.extend_from_slice(&chunk[..n]);
1060                        if let Some(i) = buf.windows(4).position(|w| w == b"\r\n\r\n") {
1061                            break i + 4;
1062                        }
1063                    };
1064                    let head = String::from_utf8_lossy(&buf[..head_end]).to_string();
1065                    let length = head
1066                        .lines()
1067                        .find_map(|l| {
1068                            let (k, v) = l.split_once(':')?;
1069                            k.eq_ignore_ascii_case("content-length")
1070                                .then(|| v.trim().parse::<usize>().ok())?
1071                        })
1072                        .unwrap_or(0);
1073                    while buf.len() < head_end + length {
1074                        let n = conn.read(&mut chunk).await.unwrap_or(0);
1075                        if n == 0 {
1076                            break;
1077                        }
1078                        buf.extend_from_slice(&chunk[..n]);
1079                    }
1080                    let mut first = head.lines().next().unwrap_or_default().split(' ');
1081                    let (method, path) = (first.next().unwrap_or(""), first.next().unwrap_or(""));
1082                    counted.fetch_add(1, Ordering::SeqCst);
1083                    let (status, body) = answer(method, path);
1084                    let resp = format!(
1085                        "HTTP/1.1 {status} X\r\ncontent-type: application/json\r\n\
1086                         content-length: {}\r\nconnection: close\r\n\r\n{body}",
1087                        body.len()
1088                    );
1089                    let _ = conn.write_all(resp.as_bytes()).await;
1090                    let _ = conn.shutdown().await;
1091                });
1092            }
1093        });
1094        (url, count)
1095    }
1096
1097    /// A result the server will not take (413, too large) does not stop
1098    /// the worker: the job gets an error instead, which it does take.
1099    #[tokio::test]
1100    async fn a_result_too_large_to_post_becomes_an_error() {
1101        use std::sync::{Arc, Mutex};
1102        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1103        let url = format!("http://{}", listener.local_addr().unwrap());
1104        let bodies: Arc<Mutex<Vec<String>>> = Arc::default();
1105        let seen = bodies.clone();
1106        tokio::spawn(async move {
1107            while let Ok((mut conn, _)) = listener.accept().await {
1108                let seen = seen.clone();
1109                tokio::spawn(async move {
1110                    let mut buf = Vec::new();
1111                    let mut chunk = [0u8; 65536];
1112                    let head_end = loop {
1113                        let n = conn.read(&mut chunk).await.unwrap_or(0);
1114                        if n == 0 {
1115                            return;
1116                        }
1117                        buf.extend_from_slice(&chunk[..n]);
1118                        if let Some(i) = buf.windows(4).position(|w| w == b"\r\n\r\n") {
1119                            break i + 4;
1120                        }
1121                    };
1122                    let head = String::from_utf8_lossy(&buf[..head_end]).to_string();
1123                    let length = head
1124                        .lines()
1125                        .find_map(|l| {
1126                            let (k, v) = l.split_once(':')?;
1127                            k.eq_ignore_ascii_case("content-length")
1128                                .then(|| v.trim().parse::<usize>().ok())?
1129                        })
1130                        .unwrap_or(0);
1131                    while buf.len() < head_end + length {
1132                        let n = conn.read(&mut chunk).await.unwrap_or(0);
1133                        if n == 0 {
1134                            break;
1135                        }
1136                        buf.extend_from_slice(&chunk[..n]);
1137                    }
1138                    let body = String::from_utf8_lossy(&buf[head_end..]).to_string();
1139                    let (status, answer) = if body.contains("\"evaluate\"") {
1140                        (413, r#"{"error":"request body too large"}"#.to_string())
1141                    } else {
1142                        (
1143                            200,
1144                            r#"{"id":"job_1","state":"queued","applied":false,"follow_up":null}"#
1145                                .to_string(),
1146                        )
1147                    };
1148                    seen.lock().unwrap().push(body);
1149                    let resp = format!(
1150                        "HTTP/1.1 {status} X\r\ncontent-type: application/json\r\n\
1151                         content-length: {}\r\nconnection: close\r\n\r\n{answer}",
1152                        answer.len()
1153                    );
1154                    let _ = conn.write_all(resp.as_bytes()).await;
1155                    let _ = conn.shutdown().await;
1156                });
1157            }
1158        });
1159        let dir = tempfile::tempdir().unwrap();
1160        let w = Worker::new(config(dir.path(), &url)).unwrap();
1161        let job = Job {
1162            id: "job_1".into(),
1163            kind: KIND_EVALUATE.into(),
1164            lease_id: "lse_1".into(),
1165            lease_expires_at: "2026-10-02T09:16:03.118Z".into(),
1166            attempt: 1,
1167            merge: None,
1168            evaluate: None,
1169        };
1170        let report = ResultRequest {
1171            lease_id: "lse_1".into(),
1172            merge: None,
1173            error: None,
1174            evaluate: Some(EvaluateResult {
1175                findings: Vec::new(),
1176                details: serde_json::json!({"skipped": [], "findings": {}}),
1177            }),
1178        };
1179        let step = tokio::time::timeout(Duration::from_secs(10), w.report(&job, &report))
1180            .await
1181            .expect("the worker kept retrying a result the server will never take");
1182        assert!(matches!(step, Ok(Step::Done(_))), "{step:?}");
1183        let bodies = bodies.lock().unwrap();
1184        assert_eq!(bodies.len(), 2, "{bodies:?}");
1185        let error: serde_json::Value = serde_json::from_str(&bodies[1]).unwrap();
1186        assert!(
1187            error["error"]
1188                .as_str()
1189                .unwrap()
1190                .contains("more than the server takes"),
1191            "{error}"
1192        );
1193        assert_eq!(error["lease_id"], "lse_1");
1194    }
1195
1196    fn discovery(capabilities: &str) -> String {
1197        format!(
1198            r#"{{"protocol":{{"current":1,"supported":[1]}},"server":{{"version":"0.4.0","build":{{"channel":"release"}}}},"min_client":"0.1.0","auth":{{"methods":["bearer"]}},"capabilities":{{{capabilities}}}}}"#
1199        )
1200    }
1201
1202    /// Runs a worker against `url` until it stops, which must be soon.
1203    async fn run_briefly(url: &str) -> Result<(), Fatal> {
1204        let dir = tempfile::tempdir().unwrap();
1205        let w = Worker::new(config(dir.path(), url)).unwrap();
1206        tokio::time::timeout(Duration::from_secs(10), w.run(std::future::pending()))
1207            .await
1208            .expect("the worker kept retrying something that will not change")
1209    }
1210
1211    /// A server that is not Recall, or has no merge queue, is never asked
1212    /// to enrol anything.
1213    #[tokio::test]
1214    async fn a_server_without_the_merge_queue_is_never_enrolled_with() {
1215        let (url, count) = fake_server(|_, path| match path {
1216            DISCOVERY_PATH => (404, r#"{"error":"not found"}"#.into()),
1217            _ => (500, "{}".into()),
1218        })
1219        .await;
1220        assert!(matches!(run_briefly(&url).await, Err(Fatal::NotRecall(_))));
1221        assert_eq!(count.load(std::sync::atomic::Ordering::SeqCst), 1);
1222
1223        let (url, count) = fake_server(|_, path| match path {
1224            DISCOVERY_PATH => (200, discovery(r#""devices":{}"#)),
1225            _ => (500, "{}".into()),
1226        })
1227        .await;
1228        match run_briefly(&url).await {
1229            Err(Fatal::NotRecall(why)) => assert!(why.contains("no merge queue"), "{why}"),
1230            other => panic!("not refused: {other:?}"),
1231        }
1232        assert_eq!(count.load(std::sync::atomic::Ordering::SeqCst), 1);
1233
1234        let (url, _) = fake_server(|_, _| (200, "<html>hello</html>".into())).await;
1235        assert!(matches!(run_briefly(&url).await, Err(Fatal::NotRecall(_))));
1236    }
1237
1238    /// An enrolment refused with a 404 is refused for good: the worker says
1239    /// so once, rather than asking again every minute.
1240    #[tokio::test]
1241    async fn an_enrolment_refused_with_404_is_final() {
1242        let (url, count) = fake_server(|_, path| match path {
1243            DISCOVERY_PATH => (200, discovery(r#""merge_queue":{}"#)),
1244            _ => (404, r#"{"error":"not found"}"#.into()),
1245        })
1246        .await;
1247        match run_briefly(&url).await {
1248            Err(Fatal::Other(why)) => assert!(why.contains("refused the enrolment"), "{why}"),
1249            other => panic!("not refused: {other:?}"),
1250        }
1251        assert_eq!(count.load(std::sync::atomic::Ordering::SeqCst), 2);
1252    }
1253}