1use 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
52const UNSIGNED_TIMEOUT: Duration = Duration::from_secs(30);
54
55const RESULT_TIMEOUT: Duration = Duration::from_secs(60);
57
58const RESULT_TRIES: u32 = 5;
61
62const MAX_BACKOFF: Duration = Duration::from_secs(60);
65
66const NOT_LOGGED_IN_RECHECK: Duration = Duration::from_secs(60);
70
71const FAST_POLL_FOR: Duration = Duration::from_secs(60);
75
76const SLOW_POLL: Duration = Duration::from_secs(30);
78
79#[derive(Debug, thiserror::Error)]
82pub enum Fatal {
83 #[error("{0}")]
85 Identity(String),
86 #[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 made_for: String,
95 configured: String,
97 path: String,
99 },
100 #[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 #[error("{0}")]
110 NotRecall(String),
111 #[error("the owner denied this worker's enrolment. Delete {0} and restart to ask again")]
113 Denied(String),
114 #[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 #[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 #[error("{0}")]
130 Other(String),
131}
132
133#[derive(Debug, Clone, PartialEq, Eq)]
135pub enum Step {
136 Idle,
138 Done(ResultResponse),
140 Rejected(String),
143}
144
145pub struct Worker {
147 cfg: Config,
148 api: Api,
149 id: Identity,
150 merger: Merger,
151 status: Status,
152 status_at: Option<Instant>,
153 announced: Option<String>,
155 evaluations: bool,
158}
159
160fn log(message: &str) {
161 eprintln!("recall-worker: {message}");
162}
163
164fn 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
176fn 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 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 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 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 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 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 fn fatal(&self, e: &ApiError) -> Option<Fatal> {
335 match e.status()? {
336 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 #[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 (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 let result = w.work(&job(side("A"), side("A"))).await;
918 assert_eq!(result.merge.unwrap().content, "A");
919 }
920
921 #[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 #[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 #[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 #[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 #[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 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 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 #[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 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 #[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 #[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}