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