use std::future::Future;
use std::time::{Duration, Instant};
use recall_wire::devices::{
ACCESS_DENIED, AUTHORIZATION_PENDING, ENROLL_PATH, ENROLL_POLL_PATH, EXPIRED_TOKEN,
INVALID_GRANT, POLL_INTERVAL_SECONDS, SCOPE_WORKER, SLOW_DOWN,
};
use recall_wire::discovery::CAPABILITY_EVALUATION;
use recall_wire::discovery::CAPABILITY_MERGE_QUEUE;
use recall_wire::jobs::{self, KIND_EVALUATE, KIND_MERGE};
use recall_wire::{
ClaimRequest, ClaimResponse, ClaudeCliReport, Discovery, EnrollPending, EnrollPollRequest,
EnrollPollResponse, EnrollRequest, EvaluateResult, Job, MergeResult, ResultRequest,
ResultResponse, DISCOVERY_PATH,
};
use reqwest::StatusCode;
use crate::api::{Api, ApiError};
use crate::config::Config;
use crate::evaluate::{self, Settings, CONTRADICTION_TIMEOUT};
use crate::identity::Identity;
use crate::merge::{Merger, Status};
const UNSIGNED_TIMEOUT: Duration = Duration::from_secs(30);
const RESULT_TIMEOUT: Duration = Duration::from_secs(60);
const RESULT_TRIES: u32 = 5;
const MAX_BACKOFF: Duration = Duration::from_secs(60);
const NOT_LOGGED_IN_RECHECK: Duration = Duration::from_secs(60);
const FAST_POLL_FOR: Duration = Duration::from_secs(60);
const SLOW_POLL: Duration = Duration::from_secs(30);
#[derive(Debug, thiserror::Error)]
pub enum Fatal {
#[error("{0}")]
Identity(String),
#[error(
"{path} was made for {made_for}, not {configured}, and a worker's key is never \
offered to a second server. If RECALL_WORKER_SERVER is wrong, correct it; to move \
this worker to {configured}, revoke it on {made_for}, delete {path} and restart"
)]
OtherServer {
made_for: String,
configured: String,
path: String,
},
#[error(
"{0} was written by an earlier recall-worker that did not record which server it \
enrolled with, so it cannot be checked against RECALL_WORKER_SERVER. Revoke the \
device it enrolled as (or deny its code), delete {0} and restart"
)]
Unbound(String),
#[error("{0}")]
NotRecall(String),
#[error("the owner denied this worker's enrolment. Delete {0} and restart to ask again")]
Denied(String),
#[error(
"this device was approved with the {0} scope, not worker, so it cannot claim jobs. \
Revoke it, delete {1} and restart, then approve the new code as a worker \
(recall devices approve <code> --worker)"
)]
WrongScope(String, String),
#[error(
"the server refused this worker ({0}). If it was revoked on purpose, nothing is wrong; \
to enrol again, delete {1} and restart"
)]
Refused(String, String),
#[error("{0}")]
Other(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Step {
Idle,
Done(ResultResponse),
Rejected(String),
}
pub struct Worker {
cfg: Config,
api: Api,
id: Identity,
merger: Merger,
status: Status,
status_at: Option<Instant>,
announced: Option<String>,
evaluations: bool,
}
fn log(message: &str) {
eprintln!("recall-worker: {message}");
}
fn is_final(e: &ApiError) -> bool {
e.status().is_some_and(|s| {
s.is_redirection()
|| (s.is_client_error()
&& s != StatusCode::REQUEST_TIMEOUT
&& s != StatusCode::TOO_MANY_REQUESTS)
})
}
fn next_wait(e: &ApiError, backoff: Duration) -> Duration {
match e {
ApiError::Status {
retry_after: Some(s),
..
} => Duration::from_secs(*s),
_ => (backoff * 2).clamp(Duration::from_secs(1), MAX_BACKOFF),
}
}
impl Worker {
pub fn new(cfg: Config) -> Result<Self, Fatal> {
if cfg.data_dir.as_os_str().is_empty() {
return Err(Fatal::Other(
"no data directory: set RECALL_WORKER_DIR to where the worker's key is kept".into(),
));
}
let path = Identity::path(&cfg.data_dir).display().to_string();
let mut id = Identity::load_or_create(&cfg.data_dir, &cfg.server)
.map_err(|e| Fatal::Identity(format!("cannot use {path}: {e}")))?;
match id.server.as_deref() {
Some(server) if server == cfg.server => {}
Some(server) => {
return Err(Fatal::OtherServer {
made_for: server.to_string(),
configured: cfg.server.clone(),
path,
})
}
None if id.device_id.is_none() && id.enrollment_id.is_none() => {
id.server = Some(cfg.server.clone());
id.save(&cfg.data_dir)
.map_err(|e| Fatal::Identity(format!("cannot write {path}: {e}")))?;
}
None => return Err(Fatal::Unbound(path)),
}
let api = Api::new(&cfg.server).map_err(|e| Fatal::Other(e.to_string()))?;
let merger = Merger::new(cfg.claude_bin.clone(), cfg.merge_timeout);
Ok(Self {
cfg,
api,
id,
merger,
status: Status::default(),
status_at: None,
announced: None,
evaluations: false,
})
}
pub fn identity(&self) -> &Identity {
&self.id
}
fn identity_path(&self) -> String {
Identity::path(&self.cfg.data_dir).display().to_string()
}
fn save(&self) -> Result<(), Fatal> {
self.id
.save(&self.cfg.data_dir)
.map_err(|e| Fatal::Identity(format!("cannot write {}: {e}", self.identity_path())))
}
pub async fn run(mut self, shutdown: impl Future<Output = ()>) -> Result<(), Fatal> {
tokio::pin!(shutdown);
log(&format!(
"{} for {}, key fingerprint {}",
crate::user_agent(),
self.cfg.server,
self.id.fingerprint()
));
if let Some(warning) = self.cfg.plaintext_warning() {
log(&warning);
}
tokio::select! {
_ = &mut shutdown => return Ok(()),
checked = self.check_server() => checked?,
}
if self.id.device_id.is_none() {
tokio::select! {
_ = &mut shutdown => return Ok(()),
enrolled = self.enrol() => enrolled?,
}
}
log(&format!(
"enrolled as {}; waiting for jobs",
self.id.device_id.as_deref().unwrap_or("")
));
let mut backoff = Duration::ZERO;
loop {
if !backoff.is_zero() {
tokio::select! {
_ = &mut shutdown => return Ok(()),
_ = tokio::time::sleep(backoff) => {}
}
}
let step = tokio::select! {
_ = &mut shutdown => return Ok(()),
step = self.step() => step,
};
backoff = match step {
Ok(Step::Idle) => Duration::ZERO,
Ok(Step::Done(outcome)) => {
log(&format!(
"job {}: {}{}",
outcome.id,
outcome.state,
match (outcome.applied, &outcome.follow_up) {
(true, _) => ", merged file stored".to_string(),
(false, Some(next)) => {
format!(", the file changed meanwhile; follow-up {next}")
}
(false, None) => String::new(),
}
));
Duration::ZERO
}
Ok(Step::Rejected(why)) => {
log(&why);
Duration::ZERO
}
Err(e) => {
if let Some(fatal) = self.fatal(&e) {
return Err(fatal);
}
if is_final(&e) && e.status() != Some(StatusCode::UNAUTHORIZED) {
return Err(Fatal::NotRecall(format!(
"{} refused a claim ({e}), and asking again will not change that. \
If the server was rolled back to a release without the merge \
queue, restart the worker once it has the queue again",
self.cfg.server
)));
}
let wait = next_wait(&e, backoff);
log(&format!("claim failed ({e}); trying again in {wait:?}"));
wait
}
};
}
}
fn fatal(&self, e: &ApiError) -> Option<Fatal> {
match e.status()? {
StatusCode::UNAUTHORIZED
if e.message().contains("unknown device") || e.message().contains("revoked") =>
{
Some(Fatal::Refused(
e.message().to_string(),
self.identity_path(),
))
}
StatusCode::FORBIDDEN => Some(Fatal::Refused(
e.message().to_string(),
self.identity_path(),
)),
_ => None,
}
}
pub async fn check_server(&mut self) -> Result<(), Fatal> {
let server = &self.cfg.server;
let mut backoff = Duration::ZERO;
loop {
let answer: Result<Discovery, ApiError> =
self.api.get(DISCOVERY_PATH, UNSIGNED_TIMEOUT).await;
match answer {
Ok(doc) if doc.can(CAPABILITY_MERGE_QUEUE) => {
self.evaluations = doc.can(CAPABILITY_EVALUATION);
return Ok(());
}
Ok(doc) => {
return Err(Fatal::NotRecall(format!(
"{server} is Recall {}, which has no merge queue, so there is nothing \
for a worker to do. Upgrade it, or stop the worker",
doc.server.version
)))
}
Err(ApiError::Body(why)) => {
return Err(Fatal::NotRecall(format!(
"{server} did not answer {DISCOVERY_PATH} with Recall's discovery \
document ({why}); RECALL_WORKER_SERVER must name a Recall server"
)))
}
Err(e) if is_final(&e) => {
return Err(Fatal::NotRecall(format!(
"{server} answered {DISCOVERY_PATH} with {e}: it is not a Recall server, \
or one older than 0.4.1, which has no merge queue. Check \
RECALL_WORKER_SERVER"
)))
}
Err(e) => {
backoff = next_wait(&e, backoff);
log(&format!(
"cannot reach {server} ({e}); trying again in {backoff:?}"
));
tokio::time::sleep(backoff).await;
}
}
}
}
pub async fn enrol(&mut self) -> Result<(), Fatal> {
loop {
match self.id.user_code.clone() {
Some(code) if self.id.enrollment_id.is_some() => self.announce(&code),
_ => self.start_enrolment().await?,
}
match self.wait_for_approval().await? {
Some(device_id) => {
self.id.device_id = Some(device_id);
self.id.enrollment_id = None;
self.id.user_code = None;
self.save()?;
return Ok(());
}
None => {
self.id.enrollment_id = None;
self.id.user_code = None;
self.save()?;
}
}
}
}
async fn start_enrolment(&mut self) -> Result<(), Fatal> {
let req = EnrollRequest {
name: self.cfg.name.clone(),
public_key: self.id.public_key(),
agent: crate::user_agent(),
authkey: None,
};
let mut backoff = Duration::ZERO;
let pending: EnrollPending = loop {
match self.api.post(ENROLL_PATH, &req, UNSIGNED_TIMEOUT).await {
Ok(pending) => break pending,
Err(e) if e.status() == Some(StatusCode::CONFLICT) => {
return Err(Fatal::Other(format!(
"{}. Set RECALL_WORKER_NAME to another name, or revoke the old worker",
e.message()
)))
}
Err(e) if is_final(&e) => {
return Err(Fatal::Other(format!(
"{} refused the enrolment ({e}); asking again would be refused \
the same way",
self.cfg.server
)))
}
Err(e) => {
backoff = next_wait(&e, backoff);
log(&format!(
"cannot reach the server to enrol ({e}); trying again in {backoff:?}"
));
tokio::time::sleep(backoff).await;
}
}
};
self.id.enrollment_id = Some(pending.enrollment_id);
self.id.user_code = Some(pending.user_code.clone());
self.save()?;
self.announce(&pending.user_code);
Ok(())
}
fn announce(&mut self, code: &str) {
if self.announced.as_deref() == Some(code) {
return;
}
self.announced = Some(code.to_string());
log(&format!(
"waiting for approval of code {code} as a worker, key fingerprint {}",
self.id.fingerprint()
));
log(&format!(
"approve it from an admin device: recall devices approve {code} --worker \
--fingerprint {}",
self.id.fingerprint()
));
log(&format!(
"or with the operator token: POST /v1/devices/approve \
{{\"user_code\":\"{code}\",\"scope\":\"worker\",\"fingerprint\":\"{}\"}} \
(see deploy/README.md, \"The merge worker\")",
self.id.fingerprint()
));
}
fn poll_interval(since: Duration, slowed: Duration) -> Duration {
let base = if since < FAST_POLL_FOR {
Duration::from_secs(POLL_INTERVAL_SECONDS)
} else {
SLOW_POLL
};
base + slowed
}
async fn wait_for_approval(&mut self) -> Result<Option<String>, Fatal> {
let Some(enrollment_id) = self.id.enrollment_id.clone() else {
return Ok(None);
};
let since = Instant::now();
let mut slowed = Duration::ZERO;
let req = EnrollPollRequest { enrollment_id };
loop {
tokio::time::sleep(Self::poll_interval(since.elapsed(), slowed)).await;
let polled: Result<EnrollPollResponse, ApiError> = self
.api
.post(ENROLL_POLL_PATH, &req, UNSIGNED_TIMEOUT)
.await;
match polled {
Ok(approved) if approved.scope == SCOPE_WORKER => {
return Ok(Some(approved.device_id))
}
Ok(approved) => {
self.id.device_id = Some(approved.device_id);
self.save()?;
return Err(Fatal::WrongScope(approved.scope, self.identity_path()));
}
Err(e) if e.status() == Some(StatusCode::BAD_REQUEST) => match e.message() {
AUTHORIZATION_PENDING => {}
SLOW_DOWN => slowed += Duration::from_secs(5),
EXPIRED_TOKEN | INVALID_GRANT => {
log("the code expired before it was approved; asking for a new one");
return Ok(None);
}
ACCESS_DENIED => return Err(Fatal::Denied(self.identity_path())),
other => log(&format!("poll refused: {other}")),
},
Err(e) if is_final(&e) => {
return Err(Fatal::Other(format!(
"{} refused the poll for this worker's approval ({e})",
self.cfg.server
)))
}
Err(e) => log(&format!("poll failed ({e}); trying again")),
}
}
}
async fn refresh_status(&mut self) {
let every = if self.status.logged_in {
self.cfg.claude_status_interval
} else {
self.cfg.claude_status_interval.min(NOT_LOGGED_IN_RECHECK)
};
if self.status_at.is_some_and(|at| at.elapsed() < every) {
return;
}
let (was, first) = (self.status.logged_in, self.status.checked_at.is_empty());
self.status = self.merger.check_status().await;
self.status_at = Some(Instant::now());
if first || self.status.logged_in != was {
log(&if self.status.logged_in {
"the claude CLI is logged in; taking merge jobs".to_string()
} else {
format!(
"the claude CLI cannot merge ({}); taking no jobs until it can. \
Log it in with: docker exec -it -u node recall-worker claude setup-token",
if self.status.error.is_empty() {
"not logged in"
} else {
&self.status.error
}
)
});
}
}
pub fn claim_request(&self) -> ClaimRequest {
let mut kinds = Vec::new();
if self.status.logged_in {
kinds.push(KIND_MERGE.to_string());
}
if self.evaluations {
kinds.push(KIND_EVALUATE.to_string());
}
ClaimRequest {
kinds,
wait_seconds: self.cfg.wait_seconds,
lease_seconds: self.cfg.lease_seconds,
claude_cli: Some(ClaudeCliReport {
checked_at: self.status.checked_at.clone(),
available: self.status.available,
logged_in: self.status.logged_in,
error: self.status.error.clone(),
}),
}
}
pub async fn step(&mut self) -> Result<Step, ApiError> {
self.refresh_status().await;
let device_id = self.id.device_id.clone().unwrap_or_default();
let claim = self.claim_request();
let answer: ClaimResponse = self
.api
.post_signed(
jobs::CLAIM_PATH,
&claim,
self.id.key(),
&device_id,
Duration::from_secs(self.cfg.wait_seconds + 30),
)
.await?;
let Some(job) = answer.job else {
return Ok(Step::Idle);
};
let result = self.work(&job).await;
self.report(&job, &result).await
}
pub async fn work(&mut self, job: &Job) -> ResultRequest {
if let (KIND_EVALUATE, Some(input)) = (job.kind.as_str(), &job.evaluate) {
return self.evaluate(job, input).await;
}
let outcome = match (job.kind.as_str(), &job.merge) {
(KIND_MERGE, Some(m)) => {
if recall_wire::content_sha256(&m.stored.content) != m.stored.sha256
|| recall_wire::content_sha256(&m.incoming.content) != m.incoming.sha256
{
Err("a version's content does not match its sha256".to_string())
} else if m.stored.content == m.incoming.content {
Ok(m.incoming.content.clone())
} else {
log(&format!(
"job {}: merging {}/{} (attempt {})",
job.id, m.project_key, m.file_path, job.attempt
));
let merged = self
.merger
.merge(&m.stored.content, &m.incoming.content)
.await;
if merged.is_err() {
self.status_at = None;
}
merged.map_err(|e| e.to_string())
}
}
(kind, _) => Err(format!("this worker does not do {kind} jobs")),
};
match outcome {
Ok(content) => ResultRequest {
lease_id: job.lease_id.clone(),
merge: Some(MergeResult { content }),
error: None,
evaluate: None,
},
Err(error) => {
log(&format!("job {}: {error}", job.id));
ResultRequest {
lease_id: job.lease_id.clone(),
merge: None,
error: Some(error),
evaluate: None,
}
}
}
}
async fn evaluate(&mut self, job: &Job, input: &recall_wire::EvaluateInput) -> ResultRequest {
log(&format!(
"job {}: evaluation {} of {} files{} (attempt {})",
job.id,
input.evaluation_id,
input.files.len(),
if input.contradictions {
", with the contradiction check"
} else {
""
},
job.attempt
));
let now = time::OffsetDateTime::now_utc();
let deadline = time::OffsetDateTime::parse(
&job.lease_expires_at,
&time::format_description::well_known::Rfc3339,
)
.ok()
.map(|end| {
let left = (end - now).max(time::Duration::ZERO);
Instant::now() + Duration::try_from(left).unwrap_or_default()
});
let settings = Settings {
now,
stale_after: Duration::from_secs(self.cfg.eval_stale_days.saturating_mul(24 * 60 * 60)),
deadline,
cli_unavailable: (!self.status.logged_in).then(|| {
if self.status.error.is_empty() {
"not logged in".to_string()
} else {
self.status.error.clone()
}
}),
};
let claude = Merger::new(self.cfg.claude_bin.clone(), CONTRADICTION_TIMEOUT);
let report = evaluate::evaluate(input, &settings, &claude).await;
log(&format!(
"job {}: {} findings {:?}{}",
job.id,
report.findings.len(),
evaluate::counts(&report),
if report.details.skipped.is_empty() {
String::new()
} else {
format!(
", {} checks skipped (see the report)",
report.details.skipped.len()
)
}
));
let details = serde_json::to_value(&report.details).unwrap_or_default();
ResultRequest {
lease_id: job.lease_id.clone(),
merge: None,
error: None,
evaluate: Some(EvaluateResult {
findings: report.findings,
details,
}),
}
}
async fn report(&self, job: &Job, result: &ResultRequest) -> Result<Step, ApiError> {
let device_id = self.id.device_id.clone().unwrap_or_default();
let mut result = result.clone();
let mut wait = Duration::from_secs(1);
let mut last = None;
let mut tries = 0;
while tries < RESULT_TRIES {
tries += 1;
let posted: Result<ResultResponse, ApiError> = self
.api
.post_signed(
&jobs::result_path(&job.id),
&result,
self.id.key(),
&device_id,
RESULT_TIMEOUT,
)
.await;
match posted {
Ok(outcome) => return Ok(Step::Done(outcome)),
Err(e)
if result.error.is_none()
&& matches!(
e.status(),
Some(StatusCode::PAYLOAD_TOO_LARGE | StatusCode::BAD_REQUEST)
) =>
{
let why = if e.status() == Some(StatusCode::PAYLOAD_TOO_LARGE) {
format!(
"the result came to {} bytes, more than the server takes",
serde_json::to_vec(&result).map_or(0, |b| b.len())
)
} else {
let said: String = e.message().chars().take(200).collect();
format!("the server refused the result: {said}")
};
log(&format!("job {}: {why}; reporting that instead", job.id));
result = ResultRequest {
lease_id: result.lease_id.clone(),
merge: None,
error: Some(why),
evaluate: None,
};
tries = 0;
}
Err(e)
if matches!(
e.status(),
Some(
StatusCode::CONFLICT
| StatusCode::NOT_FOUND
| StatusCode::BAD_REQUEST
| StatusCode::PAYLOAD_TOO_LARGE
)
) =>
{
return Ok(Step::Rejected(format!(
"job {}: the server did not take the result: {}",
job.id,
e.message()
)))
}
Err(e) if self.fatal(&e).is_some() => return Err(e),
Err(e) => {
log(&format!(
"job {}: posting the result failed ({e}); trying again in {wait:?}",
job.id
));
last = Some(e);
tokio::time::sleep(wait).await;
wait *= 2;
}
}
}
Err(last.unwrap_or_else(|| ApiError::Transport("no attempt was made".into())))
}
}
#[cfg(test)]
mod tests {
use super::*;
use recall_wire::{MergeInput, MergeSide};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
fn config(dir: &std::path::Path, server: &str) -> Config {
Config {
server: server.to_string(),
data_dir: dir.to_path_buf(),
claude_bin: "definitely-not-a-real-claude".into(),
..Config::default()
}
}
fn side(content: &str) -> MergeSide {
MergeSide {
sha256: recall_wire::content_sha256(content),
content: content.to_string(),
source_env: "laptop".into(),
updated_at: "2026-10-02T09:10:11.020Z".into(),
}
}
fn job(stored: MergeSide, incoming: MergeSide) -> Job {
Job {
id: "job_1".into(),
kind: KIND_MERGE.into(),
lease_id: "lse_1".into(),
lease_expires_at: "2026-10-02T09:16:03.118Z".into(),
attempt: 1,
merge: Some(MergeInput {
project_key: "acme/app".into(),
file_path: "topics/auth.md".into(),
stored,
incoming,
}),
evaluate: None,
}
}
fn worker() -> (tempfile::TempDir, Worker) {
let dir = tempfile::tempdir().unwrap();
let w = Worker::new(config(dir.path(), "http://127.0.0.1:9")).unwrap();
(dir, w)
}
#[tokio::test]
async fn a_version_that_does_not_match_its_hash_is_not_merged() {
let (_dir, mut w) = worker();
let mut stored = side("A");
stored.sha256 = recall_wire::content_sha256("not A");
let mut incoming = side("B");
incoming.sha256 = recall_wire::content_sha256("not B");
for (s, i) in [
(stored.clone(), side("B")),
(side("A"), incoming.clone()),
(stored.clone(), {
let mut same = side("A");
same.sha256 = stored.sha256.clone();
same
}),
] {
let result = w.work(&job(s, i)).await;
assert!(result.merge.is_none(), "{result:?}");
assert!(
result.error.as_deref().unwrap().contains("sha256"),
"{result:?}"
);
assert_eq!(result.lease_id, "lse_1");
}
let result = w.work(&job(side("A"), side("A"))).await;
assert_eq!(result.merge.unwrap().content, "A");
}
#[test]
fn a_worker_that_cannot_merge_claims_nothing() {
let (_dir, mut w) = worker();
assert!(w.claim_request().kinds.is_empty(), "never checked");
w.status = Status {
checked_at: "2026-10-02T09:13:40.002Z".into(),
available: true,
logged_in: false,
error: "not logged in".into(),
};
let claim = w.claim_request();
assert!(claim.kinds.is_empty());
assert_eq!(claim.claude_cli.unwrap().error, "not logged in");
w.status.logged_in = true;
assert_eq!(w.claim_request().kinds, vec![KIND_MERGE.to_string()]);
}
#[test]
fn a_worker_takes_evaluations_whenever_the_server_makes_them() {
let (_dir, mut w) = worker();
w.evaluations = true;
assert_eq!(w.claim_request().kinds, vec![KIND_EVALUATE.to_string()]);
w.status.logged_in = true;
assert_eq!(
w.claim_request().kinds,
vec![KIND_MERGE.to_string(), KIND_EVALUATE.to_string()]
);
}
#[tokio::test]
async fn a_failed_merge_has_the_cli_checked_again() {
let (_dir, mut w) = worker();
w.status.logged_in = true;
w.status_at = Some(Instant::now());
let result = w.work(&job(side("A"), side("B"))).await;
assert!(result.error.unwrap().contains("unavailable"));
assert!(w.status_at.is_none());
}
#[test]
fn a_code_is_polled_quickly_for_a_minute_then_slowly() {
let fast = Duration::from_secs(POLL_INTERVAL_SECONDS);
assert_eq!(Worker::poll_interval(Duration::ZERO, Duration::ZERO), fast);
assert_eq!(
Worker::poll_interval(Duration::from_secs(59), Duration::from_secs(5)),
fast + Duration::from_secs(5)
);
assert_eq!(
Worker::poll_interval(Duration::from_secs(61), Duration::ZERO),
SLOW_POLL
);
assert!(SLOW_POLL >= Duration::from_secs(30));
}
#[test]
fn an_identity_is_bound_to_its_server() {
let dir = tempfile::tempdir().unwrap();
Worker::new(config(dir.path(), "http://recall-server:8787")).unwrap();
match Worker::new(config(dir.path(), "https://recall.example.com")) {
Err(Fatal::OtherServer {
made_for,
configured,
..
}) => {
assert_eq!(made_for, "http://recall-server:8787");
assert_eq!(configured, "https://recall.example.com");
}
other => panic!("not refused: {:?}", other.err()),
}
assert!(Worker::new(config(dir.path(), "http://recall-server:8787")).is_ok());
}
#[test]
fn an_identity_from_before_servers_were_recorded() {
let dir = tempfile::tempdir().unwrap();
let mut id = Identity::from_seed([3; 32]);
id.save(dir.path()).unwrap();
Worker::new(config(dir.path(), "http://recall-server:8787")).unwrap();
assert_eq!(
Identity::load(dir.path())
.unwrap()
.unwrap()
.server
.as_deref(),
Some("http://recall-server:8787")
);
let dir = tempfile::tempdir().unwrap();
id.device_id = Some("dev_somewhere".into());
id.save(dir.path()).unwrap();
assert!(matches!(
Worker::new(config(dir.path(), "http://recall-server:8787")),
Err(Fatal::Unbound(_))
));
}
#[test]
fn a_worker_needs_a_data_directory() {
assert!(matches!(
Worker::new(Config {
server: "http://127.0.0.1:9".into(),
..Config::default()
}),
Err(Fatal::Other(_))
));
}
async fn fake_server(
answer: fn(&str, &str) -> (u16, String),
) -> (String, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
use std::sync::atomic::{AtomicUsize, Ordering};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
let count = std::sync::Arc::new(AtomicUsize::new(0));
let counted = count.clone();
tokio::spawn(async move {
while let Ok((mut conn, _)) = listener.accept().await {
let counted = counted.clone();
tokio::spawn(async move {
let mut buf = Vec::new();
let mut chunk = [0u8; 4096];
let head_end = loop {
let n = conn.read(&mut chunk).await.unwrap_or(0);
if n == 0 {
return;
}
buf.extend_from_slice(&chunk[..n]);
if let Some(i) = buf.windows(4).position(|w| w == b"\r\n\r\n") {
break i + 4;
}
};
let head = String::from_utf8_lossy(&buf[..head_end]).to_string();
let length = head
.lines()
.find_map(|l| {
let (k, v) = l.split_once(':')?;
k.eq_ignore_ascii_case("content-length")
.then(|| v.trim().parse::<usize>().ok())?
})
.unwrap_or(0);
while buf.len() < head_end + length {
let n = conn.read(&mut chunk).await.unwrap_or(0);
if n == 0 {
break;
}
buf.extend_from_slice(&chunk[..n]);
}
let mut first = head.lines().next().unwrap_or_default().split(' ');
let (method, path) = (first.next().unwrap_or(""), first.next().unwrap_or(""));
counted.fetch_add(1, Ordering::SeqCst);
let (status, body) = answer(method, path);
let resp = format!(
"HTTP/1.1 {status} X\r\ncontent-type: application/json\r\n\
content-length: {}\r\nconnection: close\r\n\r\n{body}",
body.len()
);
let _ = conn.write_all(resp.as_bytes()).await;
let _ = conn.shutdown().await;
});
}
});
(url, count)
}
#[tokio::test]
async fn a_result_too_large_to_post_becomes_an_error() {
use std::sync::{Arc, Mutex};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
let bodies: Arc<Mutex<Vec<String>>> = Arc::default();
let seen = bodies.clone();
tokio::spawn(async move {
while let Ok((mut conn, _)) = listener.accept().await {
let seen = seen.clone();
tokio::spawn(async move {
let mut buf = Vec::new();
let mut chunk = [0u8; 65536];
let head_end = loop {
let n = conn.read(&mut chunk).await.unwrap_or(0);
if n == 0 {
return;
}
buf.extend_from_slice(&chunk[..n]);
if let Some(i) = buf.windows(4).position(|w| w == b"\r\n\r\n") {
break i + 4;
}
};
let head = String::from_utf8_lossy(&buf[..head_end]).to_string();
let length = head
.lines()
.find_map(|l| {
let (k, v) = l.split_once(':')?;
k.eq_ignore_ascii_case("content-length")
.then(|| v.trim().parse::<usize>().ok())?
})
.unwrap_or(0);
while buf.len() < head_end + length {
let n = conn.read(&mut chunk).await.unwrap_or(0);
if n == 0 {
break;
}
buf.extend_from_slice(&chunk[..n]);
}
let body = String::from_utf8_lossy(&buf[head_end..]).to_string();
let (status, answer) = if body.contains("\"evaluate\"") {
(413, r#"{"error":"request body too large"}"#.to_string())
} else {
(
200,
r#"{"id":"job_1","state":"queued","applied":false,"follow_up":null}"#
.to_string(),
)
};
seen.lock().unwrap().push(body);
let resp = format!(
"HTTP/1.1 {status} X\r\ncontent-type: application/json\r\n\
content-length: {}\r\nconnection: close\r\n\r\n{answer}",
answer.len()
);
let _ = conn.write_all(resp.as_bytes()).await;
let _ = conn.shutdown().await;
});
}
});
let dir = tempfile::tempdir().unwrap();
let w = Worker::new(config(dir.path(), &url)).unwrap();
let job = Job {
id: "job_1".into(),
kind: KIND_EVALUATE.into(),
lease_id: "lse_1".into(),
lease_expires_at: "2026-10-02T09:16:03.118Z".into(),
attempt: 1,
merge: None,
evaluate: None,
};
let report = ResultRequest {
lease_id: "lse_1".into(),
merge: None,
error: None,
evaluate: Some(EvaluateResult {
findings: Vec::new(),
details: serde_json::json!({"skipped": [], "findings": {}}),
}),
};
let step = tokio::time::timeout(Duration::from_secs(10), w.report(&job, &report))
.await
.expect("the worker kept retrying a result the server will never take");
assert!(matches!(step, Ok(Step::Done(_))), "{step:?}");
let bodies = bodies.lock().unwrap();
assert_eq!(bodies.len(), 2, "{bodies:?}");
let error: serde_json::Value = serde_json::from_str(&bodies[1]).unwrap();
assert!(
error["error"]
.as_str()
.unwrap()
.contains("more than the server takes"),
"{error}"
);
assert_eq!(error["lease_id"], "lse_1");
}
fn discovery(capabilities: &str) -> String {
format!(
r#"{{"protocol":{{"current":1,"supported":[1]}},"server":{{"version":"0.4.0","build":{{"channel":"release"}}}},"min_client":"0.1.0","auth":{{"methods":["bearer"]}},"capabilities":{{{capabilities}}}}}"#
)
}
async fn run_briefly(url: &str) -> Result<(), Fatal> {
let dir = tempfile::tempdir().unwrap();
let w = Worker::new(config(dir.path(), url)).unwrap();
tokio::time::timeout(Duration::from_secs(10), w.run(std::future::pending()))
.await
.expect("the worker kept retrying something that will not change")
}
#[tokio::test]
async fn a_server_without_the_merge_queue_is_never_enrolled_with() {
let (url, count) = fake_server(|_, path| match path {
DISCOVERY_PATH => (404, r#"{"error":"not found"}"#.into()),
_ => (500, "{}".into()),
})
.await;
assert!(matches!(run_briefly(&url).await, Err(Fatal::NotRecall(_))));
assert_eq!(count.load(std::sync::atomic::Ordering::SeqCst), 1);
let (url, count) = fake_server(|_, path| match path {
DISCOVERY_PATH => (200, discovery(r#""devices":{}"#)),
_ => (500, "{}".into()),
})
.await;
match run_briefly(&url).await {
Err(Fatal::NotRecall(why)) => assert!(why.contains("no merge queue"), "{why}"),
other => panic!("not refused: {other:?}"),
}
assert_eq!(count.load(std::sync::atomic::Ordering::SeqCst), 1);
let (url, _) = fake_server(|_, _| (200, "<html>hello</html>".into())).await;
assert!(matches!(run_briefly(&url).await, Err(Fatal::NotRecall(_))));
}
#[tokio::test]
async fn an_enrolment_refused_with_404_is_final() {
let (url, count) = fake_server(|_, path| match path {
DISCOVERY_PATH => (200, discovery(r#""merge_queue":{}"#)),
_ => (404, r#"{"error":"not found"}"#.into()),
})
.await;
match run_briefly(&url).await {
Err(Fatal::Other(why)) => assert!(why.contains("refused the enrolment"), "{why}"),
other => panic!("not refused: {other:?}"),
}
assert_eq!(count.load(std::sync::atomic::Ordering::SeqCst), 2);
}
}