Skip to main content

recall_worker/
config.rs

1//! The worker's settings: the environment, and nothing else, as for the
2//! server.
3//!
4//! The server it works for is `RECALL_WORKER_SERVER`, a variable of its
5//! own. It never reads `RECALL_URL`: that is the `recall` client's setting,
6//! and on any machine with Recall set up it names the public server, so a
7//! worker that fell back to it would enrol with whatever server the shell
8//! it was started from happened to point at.
9//!
10//! The three merge variables keep the names `recall-server` reads, so a
11//! deployment moving its merge to a worker copies them across unchanged.
12
13use std::net::IpAddr;
14use std::path::PathBuf;
15use std::time::Duration;
16
17use recall_wire::jobs::{MAX_LEASE_SECONDS, MAX_WAIT_SECONDS, MIN_LEASE_SECONDS};
18
19/// How long before its lease ends a merge must have finished, so its result
20/// still reaches the server in time.
21pub const LEASE_MARGIN_SECONDS: u64 = 15;
22
23/// The longest `RECALL_MERGE_TIMEOUT_MS` may be: the longest lease the
24/// server grants, less [`LEASE_MARGIN_SECONDS`]. A merge allowed to run any
25/// longer would outlast its lease, and the job would be handed out again
26/// while it was still being merged.
27pub const MAX_MERGE_TIMEOUT: Duration =
28    Duration::from_secs(MAX_LEASE_SECONDS - LEASE_MARGIN_SECONDS);
29
30/// Why the worker cannot start.
31#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
32pub enum ConfigError {
33    /// There is no server to work for.
34    #[error(
35        "RECALL_WORKER_SERVER is not set; it names the server this worker merges for, \
36         such as http://recall-server:8787"
37    )]
38    MissingServer,
39    /// Only the client's `RECALL_URL` is set, which the worker never reads.
40    #[error(
41        "RECALL_WORKER_SERVER is not set. RECALL_URL is, but that is the recall client's \
42         setting and the worker never reads it: on a machine with Recall set up it names \
43         the public server. Set RECALL_WORKER_SERVER to the server this worker merges for, \
44         such as http://recall-server:8787"
45    )]
46    OnlyClientUrl,
47    /// `RECALL_WORKER_SERVER` is not an `http` or `https` URL.
48    #[error("RECALL_WORKER_SERVER must be an http:// or https:// URL, got {0}")]
49    BadServer(String),
50    /// There is nowhere to keep the worker's key.
51    #[error(
52        "RECALL_WORKER_DIR is not set; it names the directory that holds the worker's key, \
53         such as /data (the image sets it)"
54    )]
55    MissingDir,
56    /// `RECALL_MERGE_TIMEOUT_MS` would let a merge outlast its lease.
57    #[error(
58        "RECALL_MERGE_TIMEOUT_MS is {0}, and may be at most {max}: a merge must finish \
59         {LEASE_MARGIN_SECONDS}s before the longest lease the server grants ({MAX_LEASE_SECONDS}s) \
60         ends, or the job is handed out again while it is still being merged",
61        max = MAX_MERGE_TIMEOUT.as_millis()
62    )]
63    TimeoutTooLong(u64),
64}
65
66/// What `recall-worker` needs.
67///
68/// | Field | Variable | Default |
69/// |---|---|---|
70/// | [`server`] | `RECALL_WORKER_SERVER` | *required*; `RECALL_URL` is never read |
71/// | [`data_dir`] | `RECALL_WORKER_DIR` | *required*; the image sets `/data` |
72/// | [`name`] | `RECALL_WORKER_NAME` | `worker` |
73/// | [`claude_bin`] | `RECALL_CLAUDE_BIN` | `claude` |
74/// | [`merge_timeout`] | `RECALL_MERGE_TIMEOUT_MS` | 45s, at most [`MAX_MERGE_TIMEOUT`] |
75/// | [`claude_status_interval`] | `RECALL_CLAUDE_STATUS_INTERVAL_MS` | 30m |
76/// | [`lease_seconds`] | `RECALL_WORKER_LEASE_SECONDS` | 120 |
77/// | [`eval_stale_days`] | `RECALL_EVAL_STALE_DAYS` | 90 |
78///
79/// [`server`]: Config::server
80/// [`data_dir`]: Config::data_dir
81/// [`name`]: Config::name
82/// [`claude_bin`]: Config::claude_bin
83/// [`merge_timeout`]: Config::merge_timeout
84/// [`claude_status_interval`]: Config::claude_status_interval
85/// [`lease_seconds`]: Config::lease_seconds
86/// [`eval_stale_days`]: Config::eval_stale_days
87#[derive(Debug, Clone)]
88pub struct Config {
89    /// The server, such as `http://recall-server:8787` from inside the
90    /// compose network, or the public address from another host. Recorded
91    /// in the worker's identity on first start; the worker refuses to run
92    /// against any other.
93    pub server: String,
94    /// Where the worker keeps its device key and id. Never the server's
95    /// volume: whoever reads this file can act as the worker. No default,
96    /// so a worker started by hand never writes a key somewhere nobody
97    /// chose.
98    pub data_dir: PathBuf,
99    /// The name it enrols as, which the owner sees in the device list.
100    pub name: String,
101    /// The `claude` binary. Never the Anthropic API.
102    pub claude_bin: String,
103    /// How long one merge may take before it is abandoned and reported as
104    /// an error, which the server retries later.
105    pub merge_timeout: Duration,
106    /// How often to re-check that the CLI is present and logged in.
107    pub claude_status_interval: Duration,
108    /// How long a claimed job is the worker's before the server hands it
109    /// to someone else. Longer than the merge timeout, so a merge that
110    /// runs its full course still reports in time.
111    pub lease_seconds: u64,
112    /// How long one claim waits for a job.
113    pub wait_seconds: u64,
114    /// How many days a file must go unchanged, naming a path or a command,
115    /// before an evaluation reports it as `stale` for the owner to
116    /// confirm.
117    pub eval_stale_days: u64,
118}
119
120impl Default for Config {
121    fn default() -> Self {
122        Self {
123            server: String::new(),
124            data_dir: PathBuf::new(),
125            name: "worker".to_string(),
126            claude_bin: "claude".to_string(),
127            merge_timeout: Duration::from_secs(45),
128            claude_status_interval: Duration::from_secs(30 * 60),
129            lease_seconds: 120,
130            wait_seconds: 25,
131            eval_stale_days: 90,
132        }
133    }
134}
135
136impl Config {
137    /// Reads the real process environment.
138    pub fn from_env() -> Result<Self, ConfigError> {
139        Self::from_lookup(|key| std::env::var(key).ok())
140    }
141
142    /// Reads configuration through `lookup`, so tests need not set
143    /// process-wide variables. An unparseable number falls back to its
144    /// default, as the server's do.
145    pub fn from_lookup<F>(lookup: F) -> Result<Self, ConfigError>
146    where
147        F: Fn(&str) -> Option<String>,
148    {
149        let get = |key: &str| lookup(key).filter(|v| !v.trim().is_empty());
150        let num = |key: &str| get(key).and_then(|v| v.trim().parse::<u64>().ok());
151        let defaults = Config::default();
152
153        let server = match (get("RECALL_WORKER_SERVER"), get("RECALL_URL")) {
154            (Some(server), _) => server,
155            (None, Some(_)) => return Err(ConfigError::OnlyClientUrl),
156            (None, None) => return Err(ConfigError::MissingServer),
157        };
158        let server = server.trim().trim_end_matches('/').to_string();
159        if !(server.starts_with("http://") || server.starts_with("https://")) {
160            return Err(ConfigError::BadServer(server));
161        }
162        let data_dir = get("RECALL_WORKER_DIR")
163            .map(|d| PathBuf::from(d.trim()))
164            .ok_or(ConfigError::MissingDir)?;
165        let merge_timeout = match num("RECALL_MERGE_TIMEOUT_MS").filter(|ms| *ms > 0) {
166            Some(ms) if Duration::from_millis(ms) > MAX_MERGE_TIMEOUT => {
167                return Err(ConfigError::TimeoutTooLong(ms))
168            }
169            Some(ms) => Duration::from_millis(ms),
170            None => defaults.merge_timeout,
171        };
172        let lease_seconds = num("RECALL_WORKER_LEASE_SECONDS")
173            .unwrap_or(defaults.lease_seconds)
174            // A lease shorter than a merge can take would hand every slow
175            // merge to the next claim while it is still running. The
176            // timeout is capped above, so this never passes the clamp.
177            .max(merge_timeout.as_millis().div_ceil(1000) as u64 + LEASE_MARGIN_SECONDS)
178            .clamp(MIN_LEASE_SECONDS, MAX_LEASE_SECONDS);
179        Ok(Config {
180            server,
181            data_dir,
182            name: get("RECALL_WORKER_NAME")
183                .map(|n| n.trim().to_string())
184                .unwrap_or(defaults.name),
185            claude_bin: get("RECALL_CLAUDE_BIN").unwrap_or(defaults.claude_bin),
186            merge_timeout,
187            claude_status_interval: num("RECALL_CLAUDE_STATUS_INTERVAL_MS")
188                .filter(|ms| *ms > 0)
189                .map(Duration::from_millis)
190                .unwrap_or(defaults.claude_status_interval),
191            lease_seconds,
192            wait_seconds: defaults.wait_seconds.min(MAX_WAIT_SECONDS),
193            eval_stale_days: num("RECALL_EVAL_STALE_DAYS")
194                .filter(|d| *d > 0)
195                .unwrap_or(defaults.eval_stale_days),
196        })
197    }
198
199    /// A warning when [`server`](Config::server) is reached over plain
200    /// `http` across a network, and [`None`] when that is fine.
201    ///
202    /// Requests are signed either way, so nobody on the path can act as the
203    /// worker. But a claim's answer carries both versions of a conflicting
204    /// file, and over `http` anyone on the path reads them. That is fine on
205    /// loopback, and on the compose file's own `backend` network, where the
206    /// server is a single-label service name such as `recall-server`;
207    /// anywhere else it wants `https`.
208    pub fn plaintext_warning(&self) -> Option<String> {
209        let rest = self.server.strip_prefix("http://")?;
210        let authority = rest.split(['/', '?', '#']).next().unwrap_or_default();
211        let authority = authority.rsplit('@').next().unwrap_or_default();
212        let host = match authority.strip_prefix('[') {
213            // [::1]:8787
214            Some(v6) => v6.split(']').next().unwrap_or_default(),
215            None => authority.split(':').next().unwrap_or_default(),
216        };
217        let loopback = host.eq_ignore_ascii_case("localhost")
218            || host.parse::<IpAddr>().is_ok_and(|ip| ip.is_loopback());
219        let service = !host.is_empty() && !host.contains('.') && host.parse::<IpAddr>().is_err();
220        if loopback || service {
221            return None;
222        }
223        Some(format!(
224            "RECALL_WORKER_SERVER is plain http:// to {host}. Requests are signed, but each job \
225             carries both versions of a file, readable by anyone on the path; use https:// \
226             unless {host} is on a network you trust"
227        ))
228    }
229}
230
231#[cfg(test)]
232mod tests {
233    use super::*;
234
235    fn env<'a>(pairs: &'a [(&'a str, &'a str)]) -> impl Fn(&str) -> Option<String> + 'a {
236        move |key| {
237            pairs
238                .iter()
239                .find(|(k, _)| *k == key)
240                .map(|(_, v)| (*v).to_string())
241        }
242    }
243
244    const DIR: (&str, &str) = ("RECALL_WORKER_DIR", "/data");
245
246    #[test]
247    fn needs_a_server_and_takes_the_servers_merge_variables() {
248        assert_eq!(
249            Config::from_lookup(env(&[DIR])).unwrap_err(),
250            ConfigError::MissingServer
251        );
252        assert!(matches!(
253            Config::from_lookup(env(&[("RECALL_WORKER_SERVER", "recall-server:8787"), DIR])),
254            Err(ConfigError::BadServer(_))
255        ));
256        let cfg = Config::from_lookup(env(&[
257            ("RECALL_WORKER_SERVER", "http://recall-server:8787/"),
258            DIR,
259            ("RECALL_CLAUDE_BIN", "/opt/claude"),
260            ("RECALL_MERGE_TIMEOUT_MS", "60000"),
261        ]))
262        .unwrap();
263        assert_eq!(cfg.server, "http://recall-server:8787");
264        assert_eq!(cfg.claude_bin, "/opt/claude");
265        assert_eq!(cfg.merge_timeout, Duration::from_secs(60));
266        assert_eq!(cfg.data_dir, PathBuf::from("/data"));
267        assert_eq!(cfg.name, "worker");
268    }
269
270    /// The client's variable is never the worker's server: a shell with
271    /// Recall set up has it naming the public server.
272    #[test]
273    fn the_clients_recall_url_is_never_read() {
274        assert_eq!(
275            Config::from_lookup(env(&[("RECALL_URL", "https://recall.example.com"), DIR]))
276                .unwrap_err(),
277            ConfigError::OnlyClientUrl
278        );
279        let cfg = Config::from_lookup(env(&[
280            ("RECALL_URL", "https://recall.example.com"),
281            ("RECALL_WORKER_SERVER", "http://127.0.0.1:8787"),
282            DIR,
283        ]))
284        .unwrap();
285        assert_eq!(cfg.server, "http://127.0.0.1:8787");
286    }
287
288    /// No directory is assumed: a worker run by hand must be told where its
289    /// key goes.
290    #[test]
291    fn the_data_directory_has_no_default() {
292        assert_eq!(
293            Config::from_lookup(env(&[("RECALL_WORKER_SERVER", "http://s")])).unwrap_err(),
294            ConfigError::MissingDir
295        );
296        assert_eq!(Config::default().data_dir, PathBuf::new());
297    }
298
299    /// The lease always outlasts a merge, and stays inside what the server
300    /// accepts.
301    #[test]
302    fn the_lease_outlasts_the_merge_timeout() {
303        let cfg = Config::from_lookup(env(&[
304            ("RECALL_WORKER_SERVER", "http://s"),
305            DIR,
306            ("RECALL_MERGE_TIMEOUT_MS", "300000"),
307            ("RECALL_WORKER_LEASE_SECONDS", "60"),
308        ]))
309        .unwrap();
310        assert_eq!(cfg.lease_seconds, 315);
311        let cfg = Config::from_lookup(env(&[
312            ("RECALL_WORKER_SERVER", "http://s"),
313            DIR,
314            ("RECALL_WORKER_LEASE_SECONDS", "100000"),
315        ]))
316        .unwrap();
317        assert_eq!(cfg.lease_seconds, MAX_LEASE_SECONDS);
318    }
319
320    /// A merge timeout the longest lease cannot cover is refused, rather
321    /// than left to run past its lease.
322    #[test]
323    fn a_merge_timeout_longer_than_any_lease_is_refused() {
324        let with = |ms: &'static str| {
325            Config::from_lookup(env(&[
326                ("RECALL_WORKER_SERVER", "http://s"),
327                DIR,
328                ("RECALL_MERGE_TIMEOUT_MS", ms),
329            ]))
330        };
331        let longest = with("585000").unwrap();
332        assert_eq!(longest.merge_timeout, MAX_MERGE_TIMEOUT);
333        assert_eq!(longest.lease_seconds, MAX_LEASE_SECONDS);
334        assert!(
335            longest.merge_timeout + Duration::from_secs(LEASE_MARGIN_SECONDS)
336                <= Duration::from_secs(longest.lease_seconds)
337        );
338        assert_eq!(
339            with("585001").unwrap_err(),
340            ConfigError::TimeoutTooLong(585_001)
341        );
342    }
343
344    #[test]
345    fn plain_http_is_quiet_only_on_loopback_and_the_compose_network() {
346        let warns = |server: &str| {
347            Config {
348                server: server.to_string(),
349                ..Config::default()
350            }
351            .plaintext_warning()
352            .is_some()
353        };
354        for quiet in [
355            "http://recall-server:8787",
356            "http://localhost:8787",
357            "http://127.0.0.1:8787",
358            "http://[::1]:8787",
359            "https://recall.example.com",
360        ] {
361            assert!(!warns(quiet), "{quiet}");
362        }
363        for loud in [
364            "http://recall.example.com",
365            "http://10.0.0.5:8787",
366            "http://[2001:db8::1]:8787",
367            "http://user@recall.example.com/prefix",
368        ] {
369            assert!(warns(loud), "{loud}");
370        }
371    }
372}