use std::net::IpAddr;
use std::path::PathBuf;
use std::time::Duration;
use recall_wire::jobs::{MAX_LEASE_SECONDS, MAX_WAIT_SECONDS, MIN_LEASE_SECONDS};
pub const LEASE_MARGIN_SECONDS: u64 = 15;
pub const MAX_MERGE_TIMEOUT: Duration =
Duration::from_secs(MAX_LEASE_SECONDS - LEASE_MARGIN_SECONDS);
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum ConfigError {
#[error(
"RECALL_WORKER_SERVER is not set; it names the server this worker merges for, \
such as http://recall-server:8787"
)]
MissingServer,
#[error(
"RECALL_WORKER_SERVER is not set. RECALL_URL is, but that is the recall client's \
setting and the worker never reads it: on a machine with Recall set up it names \
the public server. Set RECALL_WORKER_SERVER to the server this worker merges for, \
such as http://recall-server:8787"
)]
OnlyClientUrl,
#[error("RECALL_WORKER_SERVER must be an http:// or https:// URL, got {0}")]
BadServer(String),
#[error(
"RECALL_WORKER_DIR is not set; it names the directory that holds the worker's key, \
such as /data (the image sets it)"
)]
MissingDir,
#[error(
"RECALL_MERGE_TIMEOUT_MS is {0}, and may be at most {max}: a merge must finish \
{LEASE_MARGIN_SECONDS}s before the longest lease the server grants ({MAX_LEASE_SECONDS}s) \
ends, or the job is handed out again while it is still being merged",
max = MAX_MERGE_TIMEOUT.as_millis()
)]
TimeoutTooLong(u64),
}
#[derive(Debug, Clone)]
pub struct Config {
pub server: String,
pub data_dir: PathBuf,
pub name: String,
pub claude_bin: String,
pub merge_timeout: Duration,
pub claude_status_interval: Duration,
pub lease_seconds: u64,
pub wait_seconds: u64,
pub eval_stale_days: u64,
}
impl Default for Config {
fn default() -> Self {
Self {
server: String::new(),
data_dir: PathBuf::new(),
name: "worker".to_string(),
claude_bin: "claude".to_string(),
merge_timeout: Duration::from_secs(45),
claude_status_interval: Duration::from_secs(30 * 60),
lease_seconds: 120,
wait_seconds: 25,
eval_stale_days: 90,
}
}
}
impl Config {
pub fn from_env() -> Result<Self, ConfigError> {
Self::from_lookup(|key| std::env::var(key).ok())
}
pub fn from_lookup<F>(lookup: F) -> Result<Self, ConfigError>
where
F: Fn(&str) -> Option<String>,
{
let get = |key: &str| lookup(key).filter(|v| !v.trim().is_empty());
let num = |key: &str| get(key).and_then(|v| v.trim().parse::<u64>().ok());
let defaults = Config::default();
let server = match (get("RECALL_WORKER_SERVER"), get("RECALL_URL")) {
(Some(server), _) => server,
(None, Some(_)) => return Err(ConfigError::OnlyClientUrl),
(None, None) => return Err(ConfigError::MissingServer),
};
let server = server.trim().trim_end_matches('/').to_string();
if !(server.starts_with("http://") || server.starts_with("https://")) {
return Err(ConfigError::BadServer(server));
}
let data_dir = get("RECALL_WORKER_DIR")
.map(|d| PathBuf::from(d.trim()))
.ok_or(ConfigError::MissingDir)?;
let merge_timeout = match num("RECALL_MERGE_TIMEOUT_MS").filter(|ms| *ms > 0) {
Some(ms) if Duration::from_millis(ms) > MAX_MERGE_TIMEOUT => {
return Err(ConfigError::TimeoutTooLong(ms))
}
Some(ms) => Duration::from_millis(ms),
None => defaults.merge_timeout,
};
let lease_seconds = num("RECALL_WORKER_LEASE_SECONDS")
.unwrap_or(defaults.lease_seconds)
.max(merge_timeout.as_millis().div_ceil(1000) as u64 + LEASE_MARGIN_SECONDS)
.clamp(MIN_LEASE_SECONDS, MAX_LEASE_SECONDS);
Ok(Config {
server,
data_dir,
name: get("RECALL_WORKER_NAME")
.map(|n| n.trim().to_string())
.unwrap_or(defaults.name),
claude_bin: get("RECALL_CLAUDE_BIN").unwrap_or(defaults.claude_bin),
merge_timeout,
claude_status_interval: num("RECALL_CLAUDE_STATUS_INTERVAL_MS")
.filter(|ms| *ms > 0)
.map(Duration::from_millis)
.unwrap_or(defaults.claude_status_interval),
lease_seconds,
wait_seconds: defaults.wait_seconds.min(MAX_WAIT_SECONDS),
eval_stale_days: num("RECALL_EVAL_STALE_DAYS")
.filter(|d| *d > 0)
.unwrap_or(defaults.eval_stale_days),
})
}
pub fn plaintext_warning(&self) -> Option<String> {
let rest = self.server.strip_prefix("http://")?;
let authority = rest.split(['/', '?', '#']).next().unwrap_or_default();
let authority = authority.rsplit('@').next().unwrap_or_default();
let host = match authority.strip_prefix('[') {
Some(v6) => v6.split(']').next().unwrap_or_default(),
None => authority.split(':').next().unwrap_or_default(),
};
let loopback = host.eq_ignore_ascii_case("localhost")
|| host.parse::<IpAddr>().is_ok_and(|ip| ip.is_loopback());
let service = !host.is_empty() && !host.contains('.') && host.parse::<IpAddr>().is_err();
if loopback || service {
return None;
}
Some(format!(
"RECALL_WORKER_SERVER is plain http:// to {host}. Requests are signed, but each job \
carries both versions of a file, readable by anyone on the path; use https:// \
unless {host} is on a network you trust"
))
}
}
#[cfg(test)]
mod tests {
use super::*;
fn env<'a>(pairs: &'a [(&'a str, &'a str)]) -> impl Fn(&str) -> Option<String> + 'a {
move |key| {
pairs
.iter()
.find(|(k, _)| *k == key)
.map(|(_, v)| (*v).to_string())
}
}
const DIR: (&str, &str) = ("RECALL_WORKER_DIR", "/data");
#[test]
fn needs_a_server_and_takes_the_servers_merge_variables() {
assert_eq!(
Config::from_lookup(env(&[DIR])).unwrap_err(),
ConfigError::MissingServer
);
assert!(matches!(
Config::from_lookup(env(&[("RECALL_WORKER_SERVER", "recall-server:8787"), DIR])),
Err(ConfigError::BadServer(_))
));
let cfg = Config::from_lookup(env(&[
("RECALL_WORKER_SERVER", "http://recall-server:8787/"),
DIR,
("RECALL_CLAUDE_BIN", "/opt/claude"),
("RECALL_MERGE_TIMEOUT_MS", "60000"),
]))
.unwrap();
assert_eq!(cfg.server, "http://recall-server:8787");
assert_eq!(cfg.claude_bin, "/opt/claude");
assert_eq!(cfg.merge_timeout, Duration::from_secs(60));
assert_eq!(cfg.data_dir, PathBuf::from("/data"));
assert_eq!(cfg.name, "worker");
}
#[test]
fn the_clients_recall_url_is_never_read() {
assert_eq!(
Config::from_lookup(env(&[("RECALL_URL", "https://recall.example.com"), DIR]))
.unwrap_err(),
ConfigError::OnlyClientUrl
);
let cfg = Config::from_lookup(env(&[
("RECALL_URL", "https://recall.example.com"),
("RECALL_WORKER_SERVER", "http://127.0.0.1:8787"),
DIR,
]))
.unwrap();
assert_eq!(cfg.server, "http://127.0.0.1:8787");
}
#[test]
fn the_data_directory_has_no_default() {
assert_eq!(
Config::from_lookup(env(&[("RECALL_WORKER_SERVER", "http://s")])).unwrap_err(),
ConfigError::MissingDir
);
assert_eq!(Config::default().data_dir, PathBuf::new());
}
#[test]
fn the_lease_outlasts_the_merge_timeout() {
let cfg = Config::from_lookup(env(&[
("RECALL_WORKER_SERVER", "http://s"),
DIR,
("RECALL_MERGE_TIMEOUT_MS", "300000"),
("RECALL_WORKER_LEASE_SECONDS", "60"),
]))
.unwrap();
assert_eq!(cfg.lease_seconds, 315);
let cfg = Config::from_lookup(env(&[
("RECALL_WORKER_SERVER", "http://s"),
DIR,
("RECALL_WORKER_LEASE_SECONDS", "100000"),
]))
.unwrap();
assert_eq!(cfg.lease_seconds, MAX_LEASE_SECONDS);
}
#[test]
fn a_merge_timeout_longer_than_any_lease_is_refused() {
let with = |ms: &'static str| {
Config::from_lookup(env(&[
("RECALL_WORKER_SERVER", "http://s"),
DIR,
("RECALL_MERGE_TIMEOUT_MS", ms),
]))
};
let longest = with("585000").unwrap();
assert_eq!(longest.merge_timeout, MAX_MERGE_TIMEOUT);
assert_eq!(longest.lease_seconds, MAX_LEASE_SECONDS);
assert!(
longest.merge_timeout + Duration::from_secs(LEASE_MARGIN_SECONDS)
<= Duration::from_secs(longest.lease_seconds)
);
assert_eq!(
with("585001").unwrap_err(),
ConfigError::TimeoutTooLong(585_001)
);
}
#[test]
fn plain_http_is_quiet_only_on_loopback_and_the_compose_network() {
let warns = |server: &str| {
Config {
server: server.to_string(),
..Config::default()
}
.plaintext_warning()
.is_some()
};
for quiet in [
"http://recall-server:8787",
"http://localhost:8787",
"http://127.0.0.1:8787",
"http://[::1]:8787",
"https://recall.example.com",
] {
assert!(!warns(quiet), "{quiet}");
}
for loud in [
"http://recall.example.com",
"http://10.0.0.5:8787",
"http://[2001:db8::1]:8787",
"http://user@recall.example.com/prefix",
] {
assert!(warns(loud), "{loud}");
}
}
}