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)]
86pub struct Config {
87 pub server: String,
92 pub data_dir: PathBuf,
97 pub name: String,
99 pub claude_bin: String,
101 pub merge_timeout: Duration,
104 pub claude_status_interval: Duration,
106 pub lease_seconds: u64,
110 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 pub fn from_env() -> Result<Self, ConfigError> {
132 Self::from_lookup(|key| std::env::var(key).ok())
133 }
134
135 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 .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 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 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 #[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 #[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 #[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 #[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}