1use super::*;
2use mj_core::hex::lower_hex;
3
4#[derive(Debug, Clone, PartialEq, Eq)]
5pub enum WorkerBinaryAvailability {
6 Local {
7 path: PathBuf,
8 source: String,
9 },
10 Remote {
11 url: String,
12 sha256: String,
13 triple: String,
14 },
15}
16
17#[derive(Debug)]
23pub(super) struct WorkerBinarySourceSnapshot {
24 pub(super) entries: HashMap<
25 (String, WorkerBinaryRequirement),
26 std::result::Result<WorkerBinaryAvailability, String>,
27 >,
28}
29
30pub(super) static PINNED_WORKER_BINARY_SOURCES: OnceLock<WorkerBinarySourceSnapshot> =
31 OnceLock::new();
32
33pub(super) fn packaged_worker_binary_path(directory: &Path, triple: &str) -> PathBuf {
34 directory.join(format!("mj-worker-{triple}"))
35}
36
37pub(super) fn running_executable_file_name(controller: &Path) -> Option<std::ffi::OsString> {
42 let name = controller.file_name()?;
43 #[cfg(target_os = "linux")]
44 {
45 use std::os::unix::ffi::{OsStrExt, OsStringExt};
46
47 if let Some(name) = name.as_bytes().strip_suffix(b" (deleted)") {
48 return Some(std::ffi::OsString::from_vec(name.to_vec()));
49 }
50 }
51 Some(name.to_os_string())
52}
53
54pub(super) fn worker_sibling_names(controller: &Path) -> Vec<std::ffi::OsString> {
59 use std::ffi::OsString;
60 let mut names = Vec::new();
61 if let Some(own) = running_executable_file_name(controller) {
62 names.push(own);
63 }
64 let legacy = OsString::from("hel");
65 if !names.contains(&legacy) {
66 names.push(legacy);
67 }
68 names
69}
70
71pub(super) fn select_native_worker(
75 controller: &Path,
76 is_file: impl Fn(&Path) -> bool,
77) -> Option<(PathBuf, &'static str)> {
78 let directory = controller.parent()?;
79 if let (Some(profile), Some(target_dir)) = (directory.file_name(), directory.parent()) {
80 let development_worker = target_dir.join("worker").join(profile).join("mj-worker");
81 if is_file(&development_worker) {
82 return Some((development_worker, "isolated native development worker"));
83 }
84 }
85 let packaged_worker = directory.join("mj-worker");
86 is_file(&packaged_worker).then_some((packaged_worker, "native worker beside mj"))
87}
88
89pub(super) fn select_sibling_worker(
96 controller: &Path,
97 triple: &str,
98 is_file: impl Fn(&Path) -> bool,
99) -> Option<(PathBuf, &'static str)> {
100 let directory = controller.parent()?;
101 let names = worker_sibling_names(controller);
102 let mut candidates: Vec<(PathBuf, &'static str)> = Vec::new();
103 candidates.push((
105 packaged_worker_binary_path(directory, triple),
106 "beside the mj binary",
107 ));
108 if let (Some(profile), Some(target_dir)) = (directory.file_name(), directory.parent()) {
114 candidates.push((
115 target_dir
116 .join("worker")
117 .join(triple)
118 .join(profile)
119 .join("mj-worker"),
120 "isolated development musl worker",
121 ));
122 candidates.push((
123 target_dir.join(triple).join(profile).join("mj-worker"),
124 "development musl worker",
125 ));
126 for name in &names {
127 candidates.push((
128 target_dir.join(triple).join(profile).join(name),
129 "development musl sibling",
130 ));
131 }
132 }
133 let controller_name = running_executable_file_name(controller);
138 for name in names
139 .iter()
140 .filter(|name| Some(name.as_os_str()) != controller_name.as_deref())
141 {
142 candidates.push((directory.join(name), "beside the running executable"));
143 }
144 candidates.into_iter().find(|(path, _)| is_file(path))
145}
146
147#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
148pub(super) enum WorkerBinaryRequirement {
149 PortableLinux,
150 LocalHost,
151}
152
153impl WorkerBinarySourceSnapshot {
154 pub(super) fn capture<F>(cache_root: &Path, resolve: F) -> Self
155 where
156 F: Fn(&str, WorkerBinaryRequirement) -> Result<WorkerBinaryAvailability>,
157 {
158 let mut entries = HashMap::new();
159 let mut local_cache = HashMap::<PathBuf, PathBuf>::new();
160 let architectures = [
161 (std::env::consts::ARCH, WorkerBinaryRequirement::LocalHost),
162 ("x86_64", WorkerBinaryRequirement::PortableLinux),
163 ("aarch64", WorkerBinaryRequirement::PortableLinux),
164 ];
165
166 for (arch, requirement) in architectures {
167 let pinned = match resolve(arch, requirement) {
168 Ok(WorkerBinaryAvailability::Local { path, source }) => {
169 match local_cache.get(&path).cloned().map(Ok).unwrap_or_else(|| {
170 copy_worker_source_to_cache(&path, cache_root).inspect(|cached| {
171 local_cache.insert(path.clone(), cached.clone());
172 })
173 }) {
174 Ok(cached) => Ok(WorkerBinaryAvailability::Local {
175 path: cached,
176 source,
177 }),
178 Err(error) => {
179 let error = format!(
180 "pin worker source {} for {arch} ({requirement:?}): {error:#}",
181 path.display()
182 );
183 tracing::warn!(arch, requirement = ?requirement, error = %error);
184 Err(error)
185 }
186 }
187 }
188 Ok(WorkerBinaryAvailability::Remote {
189 url,
190 sha256,
191 triple,
192 }) => Ok(WorkerBinaryAvailability::Remote {
193 url,
194 sha256,
195 triple,
196 }),
197 Err(error) => {
198 let error = format!("{error:#}");
199 tracing::debug!(
200 arch,
201 requirement = ?requirement,
202 error = %error,
203 "worker source was unavailable when the daemon started"
204 );
205 Err(error)
206 }
207 };
208 entries.insert((arch.to_owned(), requirement), pinned);
209 }
210
211 Self { entries }
212 }
213
214 pub(super) fn resolve(
215 &self,
216 arch: &str,
217 requirement: WorkerBinaryRequirement,
218 ) -> Result<WorkerBinaryAvailability> {
219 let Some(source) = self.entries.get(&(arch.to_owned(), requirement)) else {
220 bail!(
221 "worker source for {arch} ({requirement:?}) was not captured when the daemon started"
222 );
223 };
224 match source {
225 Ok(availability) => {
226 if let WorkerBinaryAvailability::Local { path, .. } = availability {
227 verify_worker_build(path).inspect_err(|error| {
228 tracing::warn!(path = %path.display(), error = %error, "rejecting pinned worker source");
229 })?;
230 }
231 Ok(availability.clone())
232 }
233 Err(error) => bail!(
234 "worker source for {arch} ({requirement:?}) was unavailable when the daemon started: {error}"
235 ),
236 }
237 }
238}
239
240pub fn pin_worker_binary_sources() -> Result<()> {
244 if PINNED_WORKER_BINARY_SOURCES.get().is_some() {
245 return Ok(());
246 }
247 let current = std::env::current_exe().context("resolve Mjolnir controller binary")?;
248 let cache_root = data_dir().join("workers").join("pinned");
249 let started = std::time::Instant::now();
250 let snapshot = WorkerBinarySourceSnapshot::capture(&cache_root, |arch, requirement| {
251 worker_binary_prerequisite_for_current(arch, requirement, ¤t, &|path| path.is_file())
252 });
253 tracing::info!(
254 elapsed_ms = started.elapsed().as_millis(),
255 "worker sources pinned"
256 );
257 let _ = PINNED_WORKER_BINARY_SOURCES.set(snapshot);
260 Ok(())
261}
262
263pub(super) fn copy_worker_source_to_cache(source: &Path, cache_root: &Path) -> Result<PathBuf> {
264 std::fs::create_dir_all(cache_root)
265 .with_context(|| format!("create pinned worker cache {}", cache_root.display()))?;
266 let mut input =
267 File::open(source).with_context(|| format!("open worker source {}", source.display()))?;
268 let metadata = input
269 .metadata()
270 .with_context(|| format!("stat worker source {}", source.display()))?;
271 let mut temporary = tempfile::NamedTempFile::new_in(cache_root)
272 .with_context(|| format!("create pinned worker staging file {}", cache_root.display()))?;
273 let mut digest = Sha256::new();
274 let mut buffer = [0_u8; 128 * 1024];
275 loop {
276 let count = input
277 .read(&mut buffer)
278 .with_context(|| format!("read worker source {}", source.display()))?;
279 if count == 0 {
280 break;
281 }
282 temporary
283 .write_all(&buffer[..count])
284 .with_context(|| format!("copy worker source {}", source.display()))?;
285 digest.update(&buffer[..count]);
286 }
287 temporary
288 .as_file_mut()
289 .sync_all()
290 .with_context(|| format!("flush pinned worker source {}", source.display()))?;
291 std::fs::set_permissions(temporary.path(), metadata.permissions())
292 .with_context(|| format!("preserve permissions for {}", source.display()))?;
293 let digest = lower_hex(digest.finalize());
294 publish_cached_worker(temporary, cache_root, &digest)
295}
296
297pub(super) fn publish_cached_worker(
301 temporary: tempfile::NamedTempFile,
302 cache_root: &Path,
303 digest: &str,
304) -> Result<PathBuf> {
305 verify_worker_build(temporary.path())?;
306 let directory = cache_root.join(digest);
307 std::fs::create_dir_all(&directory)
308 .with_context(|| format!("create pinned worker cache {}", directory.display()))?;
309 let destination = directory.join("hel");
310 if destination.is_file() {
311 verify_cached_worker(&destination, digest)?;
312 return Ok(destination);
313 }
314 match temporary.persist_noclobber(&destination) {
315 Ok(_) => {
316 #[cfg(unix)]
317 File::open(&directory)
318 .and_then(|directory| directory.sync_all())
319 .with_context(|| format!("flush pinned worker cache {}", directory.display()))?;
320 Ok(destination)
321 }
322 Err(error) if error.error.kind() == ErrorKind::AlreadyExists => {
323 if destination.is_file() {
324 verify_cached_worker(&destination, digest)?;
325 Ok(destination)
326 } else {
327 Err(error.error).with_context(|| {
328 format!("publish pinned worker artifact {}", destination.display())
329 })
330 }
331 }
332 Err(error) => Err(error.error)
333 .with_context(|| format!("publish pinned worker artifact {}", destination.display())),
334 }
335}
336
337fn verify_cached_worker(path: &Path, digest: &str) -> Result<()> {
338 verify_worker_build(path)?;
339 ensure!(
340 mj_core::worker_launch::worker_executable_digest(path)? == digest,
341 "content-addressed worker cache {} does not match {digest} checksum",
342 path.display()
343 );
344 Ok(())
345}
346
347pub fn worker_binary_prerequisite_for_arch(arch: &str) -> Result<WorkerBinaryAvailability> {
354 worker_binary_for_arch(arch, WorkerBinaryRequirement::PortableLinux)
355}
356
357pub fn native_worker_binary_prerequisite() -> Result<WorkerBinaryAvailability> {
364 worker_binary_for_arch(std::env::consts::ARCH, WorkerBinaryRequirement::LocalHost)
365}
366
367pub(super) fn worker_binary_for_arch(
368 arch: &str,
369 requirement: WorkerBinaryRequirement,
370) -> Result<WorkerBinaryAvailability> {
371 if let Some(snapshot) = PINNED_WORKER_BINARY_SOURCES.get() {
372 let pinned = snapshot.resolve(arch, requirement);
373 if pinned_source_is_usable(&pinned, &|path| path.is_file()) {
374 return pinned;
375 }
376 return resolve_worker_source_again(arch, requirement, pinned.err());
382 }
383 let current = std::env::current_exe().context("resolve Mjolnir controller binary")?;
384 worker_binary_prerequisite_for_current(arch, requirement, ¤t, &|path| path.is_file())
385}
386
387pub(super) fn pinned_source_is_usable(
393 pinned: &Result<WorkerBinaryAvailability>,
394 is_file: &dyn Fn(&Path) -> bool,
395) -> bool {
396 match pinned {
397 Ok(WorkerBinaryAvailability::Local { path, .. }) => is_file(path),
398 Ok(WorkerBinaryAvailability::Remote { .. }) => true,
399 Err(_) => false,
400 }
401}
402
403fn resolve_worker_source_again(
408 arch: &str,
409 requirement: WorkerBinaryRequirement,
410 pinned_error: Option<anyhow::Error>,
411) -> Result<WorkerBinaryAvailability> {
412 let current = std::env::current_exe().context("resolve Mjolnir controller binary")?;
413 let resolved =
414 worker_binary_prerequisite_for_current(arch, requirement, ¤t, &|path| path.is_file());
415 match resolved {
416 Ok(WorkerBinaryAvailability::Local { path, source }) => {
417 let cache_root = data_dir().join("workers").join("pinned");
418 let path = copy_worker_source_to_cache(&path, &cache_root)
419 .context("pin the re-resolved worker source")?;
420 tracing::info!(
421 arch,
422 requirement = ?requirement,
423 source = %source,
424 "re-resolved a worker source the daemon could not pin at startup"
425 );
426 Ok(WorkerBinaryAvailability::Local { path, source })
427 }
428 Ok(remote) => Ok(remote),
429 Err(error) => Err(match pinned_error {
432 Some(pinned) => error.context(format!("{pinned:#}")),
433 None => error,
434 }),
435 }
436}
437
438pub(super) fn worker_binary_prerequisite_for_current(
441 arch: &str,
442 requirement: WorkerBinaryRequirement,
443 current: &Path,
444 is_file: &dyn Fn(&Path) -> bool,
445) -> Result<WorkerBinaryAvailability> {
446 let triple = format!("{arch}-unknown-linux-musl");
447 let rejected = std::cell::RefCell::new(Vec::new());
448 let matches_build = |path: &Path| {
449 if !is_file(path) {
450 return false;
451 }
452 match verify_worker_build(path) {
453 Ok(()) => true,
454 Err(error) => {
455 tracing::warn!(path = %path.display(), error = %error, "skipping incompatible worker source");
456 rejected.borrow_mut().push(format!("{error:#}"));
457 false
458 }
459 }
460 };
461 if let Some(path) = mj_core::config::env_override_os("WORKER_BINARY").map(PathBuf::from) {
462 if !is_file(&path) {
463 bail!("MJ_WORKER_BINARY is not a file: {}", path.display());
464 }
465 verify_worker_build(&path).context(
468 "MJ_WORKER_BINARY does not match the running mj; rebuild it from the same commit or unset MJ_WORKER_BINARY",
469 )?;
470 return Ok(WorkerBinaryAvailability::Local {
471 path,
472 source: "MJ_WORKER_BINARY".into(),
473 });
474 }
475 let controller_replaced = !is_file(current);
479 let mut candidates = Vec::new();
480 if let Some(directory) = mj_core::config::env_override_os("WORKER_DIR").map(PathBuf::from) {
481 candidates.push((
482 packaged_worker_binary_path(&directory, &triple),
483 "MJ_WORKER_DIR",
484 ));
485 candidates.push((directory.join(&triple).join("hel"), "MJ_WORKER_DIR"));
486 }
487 if let Some((path, source)) = candidates.into_iter().find(|(path, _)| matches_build(path)) {
488 return Ok(WorkerBinaryAvailability::Local {
489 path,
490 source: source.into(),
491 });
492 }
493 if requirement == WorkerBinaryRequirement::LocalHost
494 && let Some((path, source)) = select_native_worker(current, matches_build)
495 {
496 return Ok(WorkerBinaryAvailability::Local {
497 path,
498 source: source.into(),
499 });
500 }
501 if !controller_replaced
502 && let Some((path, source)) = select_sibling_worker(current, &triple, matches_build)
503 {
504 return Ok(WorkerBinaryAvailability::Local {
505 path,
506 source: source.into(),
507 });
508 }
509 if let Some(template) = mj_core::config::env_override("WORKER_URL") {
510 let expected = mj_core::config::env_override("WORKER_SHA256")
511 .context("MJ_WORKER_URL requires MJ_WORKER_SHA256")?;
512 validate_worker_sha256(&expected)?;
513 return Ok(WorkerBinaryAvailability::Remote {
514 url: template.replace("{target}", &triple),
515 sha256: expected,
516 triple,
517 });
518 }
519 let rejected = rejected.into_inner();
520 ensure!(
521 rejected.is_empty(),
522 "no worker matching controller build {BUILD_ID} for {triple}; rejected sources:\n{}\nInstall the worker from the same mj release, or rebuild from the same commit with `cargo build --target {triple} -p brokk-mj-worker --bin mj-worker` and set MJ_WORKER_BINARY to that file (local native: `cargo build -p brokk-mj-worker --bin mj-worker`).",
523 rejected.join("\n")
524 );
525 ensure!(
528 !controller_replaced,
529 "the running mj binary was replaced or removed on disk ({}); restart the Mjolnir daemon so it runs the current build, then retry",
530 display_path(current)
531 );
532 bail!(
533 "no Linux worker for {triple}; install mj-worker-{triple} beside mj, set MJ_WORKER_DIR/MJ_WORKER_BINARY, or configure MJ_WORKER_URL and MJ_WORKER_SHA256"
534 )
535}