1use 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
19pub const LEASE_MARGIN_SECONDS: u64 = 15;
22
23pub const MAX_MERGE_TIMEOUT: Duration =
28 Duration::from_secs(MAX_LEASE_SECONDS - LEASE_MARGIN_SECONDS);
29
30#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
32pub enum ConfigError {
33 #[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 #[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 #[error("RECALL_WORKER_SERVER must be an http:// or https:// URL, got {0}")]
49 BadServer(String),
50 #[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 #[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#[derive(Debug, Clone)]
88pub struct Config {
89 pub server: String,
94 pub data_dir: PathBuf,
99 pub name: String,
101 pub claude_bin: String,
103 pub merge_timeout: Duration,
106 pub claude_status_interval: Duration,
108 pub lease_seconds: u64,
112 pub wait_seconds: u64,
114 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 pub fn from_env() -> Result<Self, ConfigError> {
139 Self::from_lookup(|key| std::env::var(key).ok())
140 }
141
142 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 .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 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 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 #[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 #[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 #[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 #[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}