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 triple.ends_with("-apple-darwin") {
109 candidates.push((
110 packaged_worker_binary_path(directory, "universal-apple-darwin"),
111 "universal Darwin worker beside mj",
112 ));
113 }
114 if let (Some(profile), Some(target_dir)) = (directory.file_name(), directory.parent()) {
120 candidates.push((
121 target_dir
122 .join("worker")
123 .join(triple)
124 .join(profile)
125 .join("mj-worker"),
126 "isolated development musl worker",
127 ));
128 candidates.push((
129 target_dir.join(triple).join(profile).join("mj-worker"),
130 "development musl worker",
131 ));
132 for name in &names {
133 candidates.push((
134 target_dir.join(triple).join(profile).join(name),
135 "development musl sibling",
136 ));
137 }
138 }
139 let controller_name = running_executable_file_name(controller);
144 for name in names
145 .iter()
146 .filter(|_| !triple.ends_with("-apple-darwin"))
147 .filter(|name| Some(name.as_os_str()) != controller_name.as_deref())
148 {
149 candidates.push((directory.join(name), "beside the running executable"));
150 }
151 candidates.into_iter().find(|(path, _)| is_file(path))
152}
153
154#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
155pub(super) enum WorkerBinaryRequirement {
156 PortableLinux,
157 Darwin,
158 LocalHost,
159}
160
161impl WorkerBinaryRequirement {
162 pub(super) fn for_os(os: targets::TargetOs) -> Self {
163 match os {
164 targets::TargetOs::Linux => Self::PortableLinux,
165 targets::TargetOs::Darwin => Self::Darwin,
166 }
167 }
168
169 pub(super) fn triple(self, arch: &str) -> String {
170 match self {
171 Self::Darwin => format!("{arch}-apple-darwin"),
172 Self::LocalHost if cfg!(target_os = "macos") => format!("{arch}-apple-darwin"),
173 _ => format!("{arch}-unknown-linux-musl"),
174 }
175 }
176}
177
178impl WorkerBinarySourceSnapshot {
179 pub(super) fn capture<F>(cache_root: &Path, resolve: F) -> Self
185 where
186 F: Fn(&str, WorkerBinaryRequirement) -> Result<WorkerBinaryAvailability> + Sync,
187 {
188 let architectures = [
189 (std::env::consts::ARCH, WorkerBinaryRequirement::LocalHost),
190 ("x86_64", WorkerBinaryRequirement::PortableLinux),
191 ("aarch64", WorkerBinaryRequirement::PortableLinux),
192 ("x86_64", WorkerBinaryRequirement::Darwin),
193 ("aarch64", WorkerBinaryRequirement::Darwin),
194 ];
195 let resolved: Vec<_> = std::thread::scope(|scope| {
196 let resolve = &resolve;
197 let handles: Vec<_> = architectures
198 .iter()
199 .map(|&(arch, requirement)| {
200 scope.spawn(move || {
201 let started = Instant::now();
202 (resolve(arch, requirement), started.elapsed())
203 })
204 })
205 .collect();
206 handles
207 .into_iter()
208 .zip(architectures)
209 .map(|(handle, (arch, requirement))| {
210 handle.join().unwrap_or_else(|_| {
211 (
212 Err(anyhow::anyhow!(
213 "resolving the worker source for {arch} ({requirement:?}) panicked"
214 )),
215 Duration::ZERO,
216 )
217 })
218 })
219 .collect()
220 });
221
222 let mut unique_paths: Vec<&Path> = Vec::new();
223 for (result, _) in &resolved {
224 if let Ok(WorkerBinaryAvailability::Local { path, .. }) = result
225 && !unique_paths.contains(&path.as_path())
226 {
227 unique_paths.push(path);
228 }
229 }
230 let pinned_paths: HashMap<PathBuf, (Result<PathBuf>, Duration)> =
231 std::thread::scope(|scope| {
232 let handles: Vec<_> = unique_paths
233 .iter()
234 .map(|&path| {
235 scope.spawn(move || {
236 let started = Instant::now();
237 (
238 copy_worker_source_to_cache(path, cache_root),
239 started.elapsed(),
240 )
241 })
242 })
243 .collect();
244 unique_paths
245 .iter()
246 .zip(handles)
247 .map(|(&path, handle)| {
248 let outcome = handle.join().unwrap_or_else(|_| {
249 (
250 Err(anyhow::anyhow!("pinning {} panicked", path.display())),
251 Duration::ZERO,
252 )
253 });
254 (path.to_path_buf(), outcome)
255 })
256 .collect()
257 });
258
259 let mut entries = HashMap::new();
260 for ((arch, requirement), (result, resolve_elapsed)) in
261 architectures.into_iter().zip(resolved)
262 {
263 let pinned = match result {
264 Ok(WorkerBinaryAvailability::Local { path, source }) => {
265 let (outcome, pin_elapsed) = &pinned_paths[&path];
266 match outcome {
267 Ok(cached) => {
268 tracing::info!(
269 triple = requirement.triple(arch),
270 source = source.as_str(),
271 path = %path.display(),
272 pinned = %cached.display(),
273 build = BUILD_ID,
274 resolve_ms = resolve_elapsed.as_millis(),
275 pin_ms = pin_elapsed.as_millis(),
276 "worker source selected"
277 );
278 Ok(WorkerBinaryAvailability::Local {
279 path: cached.clone(),
280 source,
281 })
282 }
283 Err(error) => {
284 let error = format!(
285 "pin worker source {} for {arch} ({requirement:?}): {error:#}",
286 path.display()
287 );
288 tracing::warn!(arch, requirement = ?requirement, error = %error);
289 Err(error)
290 }
291 }
292 }
293 Ok(WorkerBinaryAvailability::Remote {
294 url,
295 sha256,
296 triple,
297 }) => {
298 tracing::info!(
299 triple = triple.as_str(),
300 source = "MJ_WORKER_URL",
301 url = url.as_str(),
302 sha256 = sha256.as_str(),
303 build = BUILD_ID,
304 "worker source selected"
305 );
306 Ok(WorkerBinaryAvailability::Remote {
307 url,
308 sha256,
309 triple,
310 })
311 }
312 Err(error) => {
313 let error = format!("{error:#}");
314 tracing::debug!(
315 arch,
316 requirement = ?requirement,
317 error = %error,
318 "worker source was unavailable when the daemon started"
319 );
320 Err(error)
321 }
322 };
323 entries.insert((arch.to_owned(), requirement), pinned);
324 }
325
326 Self { entries }
327 }
328
329 pub(super) fn resolve(
330 &self,
331 arch: &str,
332 requirement: WorkerBinaryRequirement,
333 ) -> Result<WorkerBinaryAvailability> {
334 let Some(source) = self.entries.get(&(arch.to_owned(), requirement)) else {
335 bail!(
336 "worker source for {arch} ({requirement:?}) was not captured when the daemon started"
337 );
338 };
339 match source {
340 Ok(availability) => {
341 if let WorkerBinaryAvailability::Local { path, .. } = availability {
342 verify_worker_build(path).inspect_err(|error| {
343 tracing::warn!(path = %path.display(), error = %error, "rejecting pinned worker source");
344 })?;
345 }
346 Ok(availability.clone())
347 }
348 Err(error) => bail!(
349 "worker source for {arch} ({requirement:?}) was unavailable when the daemon started: {error}"
350 ),
351 }
352 }
353}
354
355pub fn pin_worker_binary_sources() -> Result<()> {
359 if PINNED_WORKER_BINARY_SOURCES.get().is_some() {
360 return Ok(());
361 }
362 let current = std::env::current_exe().context("resolve Mjolnir controller binary")?;
363 let cache_root = data_dir().join("workers").join("pinned");
364 let started = std::time::Instant::now();
365 let snapshot = WorkerBinarySourceSnapshot::capture(&cache_root, |arch, requirement| {
366 worker_binary_prerequisite_with_verifier(
367 arch,
368 requirement,
369 ¤t,
370 &|path| path.is_file(),
371 &|path| verify_worker_build_indexed(&cache_root, path),
372 )
373 });
374 tracing::info!(
375 elapsed_ms = started.elapsed().as_millis(),
376 "worker sources pinned"
377 );
378 let _ = PINNED_WORKER_BINARY_SOURCES.set(snapshot);
381 Ok(())
382}
383
384fn index_entry_path(cache_root: &Path, source: &Path) -> Option<PathBuf> {
391 let metadata = std::fs::metadata(source).ok()?;
392 let mut key = Sha256::new();
393 key.update(BUILD_ID.as_bytes());
394 key.update([0]);
395 key.update(source.as_os_str().as_encoded_bytes());
396 key.update(metadata.len().to_le_bytes());
397 #[cfg(unix)]
398 {
399 use std::os::unix::fs::MetadataExt;
400 for value in [
401 metadata.mtime(),
402 metadata.mtime_nsec(),
403 metadata.ctime(),
404 metadata.ctime_nsec(),
405 ] {
406 key.update(value.to_le_bytes());
407 }
408 key.update(metadata.ino().to_le_bytes());
409 }
410 #[cfg(not(unix))]
411 {
412 let modified = metadata.modified().ok()?;
413 key.update(format!("{modified:?}").as_bytes());
414 }
415 Some(cache_root.join("index").join(lower_hex(key.finalize())))
416}
417
418fn indexed_pinned_worker(cache_root: &Path, source: &Path) -> Option<PathBuf> {
421 let entry = index_entry_path(cache_root, source)?;
422 let digest = std::fs::read_to_string(entry).ok()?;
423 let digest = digest.trim();
424 if digest.len() != 64 || !digest.bytes().all(|byte| byte.is_ascii_hexdigit()) {
425 return None;
426 }
427 let pinned = cache_root.join(digest).join("hel");
428 let source_len = std::fs::metadata(source).ok()?.len();
429 (std::fs::metadata(&pinned).ok()?.len() == source_len).then_some(pinned)
430}
431
432fn record_indexed_pin(cache_root: &Path, source: &Path, digest: &str) {
433 let Some(entry) = index_entry_path(cache_root, source) else {
434 return;
435 };
436 let write = || -> std::io::Result<()> {
438 let directory = entry.parent().expect("index entries have a parent");
439 std::fs::create_dir_all(directory)?;
440 let mut staged = tempfile::NamedTempFile::new_in(directory)?;
441 staged.write_all(digest.as_bytes())?;
442 staged
443 .persist(&entry)
444 .map(drop)
445 .map_err(|error| error.error)
446 };
447 if let Err(error) = write() {
448 tracing::debug!(source = %source.display(), %error, "could not index a pinned worker");
449 }
450}
451
452pub(super) fn verify_worker_build_indexed(cache_root: &Path, path: &Path) -> Result<()> {
456 if indexed_pinned_worker(cache_root, path).is_some() {
457 return Ok(());
458 }
459 verify_worker_build(path)
460}
461
462pub(super) fn copy_worker_source_to_cache(source: &Path, cache_root: &Path) -> Result<PathBuf> {
463 if let Some(pinned) = indexed_pinned_worker(cache_root, source) {
464 return Ok(pinned);
465 }
466 std::fs::create_dir_all(cache_root)
467 .with_context(|| format!("create pinned worker cache {}", cache_root.display()))?;
468 let mut input =
469 File::open(source).with_context(|| format!("open worker source {}", source.display()))?;
470 let metadata = input
471 .metadata()
472 .with_context(|| format!("stat worker source {}", source.display()))?;
473 let mut temporary = tempfile::NamedTempFile::new_in(cache_root)
474 .with_context(|| format!("create pinned worker staging file {}", cache_root.display()))?;
475 let mut digest = Sha256::new();
476 let mut buffer = [0_u8; 128 * 1024];
477 loop {
478 let count = input
479 .read(&mut buffer)
480 .with_context(|| format!("read worker source {}", source.display()))?;
481 if count == 0 {
482 break;
483 }
484 temporary
485 .write_all(&buffer[..count])
486 .with_context(|| format!("copy worker source {}", source.display()))?;
487 digest.update(&buffer[..count]);
488 }
489 temporary
490 .as_file_mut()
491 .sync_all()
492 .with_context(|| format!("flush pinned worker source {}", source.display()))?;
493 std::fs::set_permissions(temporary.path(), metadata.permissions())
494 .with_context(|| format!("preserve permissions for {}", source.display()))?;
495 let digest = lower_hex(digest.finalize());
496 let pinned = publish_cached_worker(temporary, cache_root, &digest)?;
497 record_indexed_pin(cache_root, source, &digest);
498 Ok(pinned)
499}
500
501pub(super) fn publish_cached_worker(
505 temporary: tempfile::NamedTempFile,
506 cache_root: &Path,
507 digest: &str,
508) -> Result<PathBuf> {
509 verify_worker_build(temporary.path())?;
510 let directory = cache_root.join(digest);
511 std::fs::create_dir_all(&directory)
512 .with_context(|| format!("create pinned worker cache {}", directory.display()))?;
513 let destination = directory.join("hel");
514 if destination.is_file() {
515 verify_cached_worker(&destination, digest)?;
516 return Ok(destination);
517 }
518 match temporary.persist_noclobber(&destination) {
519 Ok(_) => {
520 #[cfg(unix)]
521 File::open(&directory)
522 .and_then(|directory| directory.sync_all())
523 .with_context(|| format!("flush pinned worker cache {}", directory.display()))?;
524 Ok(destination)
525 }
526 Err(error) if error.error.kind() == ErrorKind::AlreadyExists => {
527 if destination.is_file() {
528 verify_cached_worker(&destination, digest)?;
529 Ok(destination)
530 } else {
531 Err(error.error).with_context(|| {
532 format!("publish pinned worker artifact {}", destination.display())
533 })
534 }
535 }
536 Err(error) => Err(error.error)
537 .with_context(|| format!("publish pinned worker artifact {}", destination.display())),
538 }
539}
540
541fn verify_cached_worker(path: &Path, digest: &str) -> Result<()> {
542 verify_worker_build(path)?;
543 ensure!(
544 mj_core::worker_launch::worker_executable_digest(path)? == digest,
545 "content-addressed worker cache {} does not match {digest} checksum",
546 path.display()
547 );
548 Ok(())
549}
550
551pub fn worker_binary_prerequisite_for_arch(arch: &str) -> Result<WorkerBinaryAvailability> {
558 worker_binary_for_arch(arch, WorkerBinaryRequirement::PortableLinux)
559}
560
561pub fn native_worker_binary_prerequisite() -> Result<WorkerBinaryAvailability> {
568 worker_binary_for_arch(std::env::consts::ARCH, WorkerBinaryRequirement::LocalHost)
569}
570
571pub(super) fn worker_binary_for_arch(
572 arch: &str,
573 requirement: WorkerBinaryRequirement,
574) -> Result<WorkerBinaryAvailability> {
575 if let Some(snapshot) = PINNED_WORKER_BINARY_SOURCES.get() {
576 let pinned = snapshot.resolve(arch, requirement);
577 if pinned_source_is_usable(&pinned, &|path| path.is_file()) {
578 return pinned;
579 }
580 return resolve_worker_source_again(arch, requirement, pinned.err());
586 }
587 let current = std::env::current_exe().context("resolve Mjolnir controller binary")?;
588 worker_binary_prerequisite_for_current(arch, requirement, ¤t, &|path| path.is_file())
589}
590
591pub(super) fn pinned_source_is_usable(
597 pinned: &Result<WorkerBinaryAvailability>,
598 is_file: &dyn Fn(&Path) -> bool,
599) -> bool {
600 match pinned {
601 Ok(WorkerBinaryAvailability::Local { path, .. }) => is_file(path),
602 Ok(WorkerBinaryAvailability::Remote { .. }) => true,
603 Err(_) => false,
604 }
605}
606
607fn resolve_worker_source_again(
612 arch: &str,
613 requirement: WorkerBinaryRequirement,
614 pinned_error: Option<anyhow::Error>,
615) -> Result<WorkerBinaryAvailability> {
616 let current = std::env::current_exe().context("resolve Mjolnir controller binary")?;
617 let resolved =
618 worker_binary_prerequisite_for_current(arch, requirement, ¤t, &|path| path.is_file());
619 match resolved {
620 Ok(WorkerBinaryAvailability::Local { path, source }) => {
621 let cache_root = data_dir().join("workers").join("pinned");
622 let path = copy_worker_source_to_cache(&path, &cache_root)
623 .context("pin the re-resolved worker source")?;
624 tracing::info!(
625 arch,
626 requirement = ?requirement,
627 source = %source,
628 "re-resolved a worker source the daemon could not pin at startup"
629 );
630 Ok(WorkerBinaryAvailability::Local { path, source })
631 }
632 Ok(remote) => Ok(remote),
633 Err(error) => Err(match pinned_error {
636 Some(pinned) => error.context(format!("{pinned:#}")),
637 None => error,
638 }),
639 }
640}
641
642pub(super) fn worker_binary_prerequisite_for_current(
645 arch: &str,
646 requirement: WorkerBinaryRequirement,
647 current: &Path,
648 is_file: &dyn Fn(&Path) -> bool,
649) -> Result<WorkerBinaryAvailability> {
650 worker_binary_prerequisite_with_verifier(
651 arch,
652 requirement,
653 current,
654 is_file,
655 &verify_worker_build,
656 )
657}
658
659pub(super) fn worker_binary_prerequisite_with_verifier(
663 arch: &str,
664 requirement: WorkerBinaryRequirement,
665 current: &Path,
666 is_file: &dyn Fn(&Path) -> bool,
667 verify: &dyn Fn(&Path) -> Result<()>,
668) -> Result<WorkerBinaryAvailability> {
669 let triple = requirement.triple(arch);
670 let rejected = std::cell::RefCell::new(Vec::new());
671 let matches_build = |path: &Path| {
672 if !is_file(path) {
673 return false;
674 }
675 match verify(path) {
676 Ok(()) => true,
677 Err(error) => {
678 tracing::warn!(path = %path.display(), error = %error, "skipping incompatible worker source");
679 rejected.borrow_mut().push(format!("{error:#}"));
680 false
681 }
682 }
683 };
684 if let Some(path) = mj_core::config::env_override_os("WORKER_BINARY").map(PathBuf::from) {
685 if !is_file(&path) {
686 bail!("MJ_WORKER_BINARY is not a file: {}", path.display());
687 }
688 verify(&path).context(
691 "MJ_WORKER_BINARY does not match the running mj; rebuild it from the same commit or unset MJ_WORKER_BINARY",
692 )?;
693 return Ok(WorkerBinaryAvailability::Local {
694 path,
695 source: "MJ_WORKER_BINARY".into(),
696 });
697 }
698 let controller_replaced = !is_file(current);
702 let mut candidates = Vec::new();
703 if let Some(directory) = mj_core::config::env_override_os("WORKER_DIR").map(PathBuf::from) {
704 candidates.push((
705 packaged_worker_binary_path(&directory, &triple),
706 "MJ_WORKER_DIR",
707 ));
708 candidates.push((directory.join(&triple).join("hel"), "MJ_WORKER_DIR"));
709 if triple.ends_with("-apple-darwin") {
710 candidates.push((
711 packaged_worker_binary_path(&directory, "universal-apple-darwin"),
712 "MJ_WORKER_DIR",
713 ));
714 }
715 }
716 if let Some((path, source)) = candidates.into_iter().find(|(path, _)| matches_build(path)) {
717 return Ok(WorkerBinaryAvailability::Local {
718 path,
719 source: source.into(),
720 });
721 }
722 if requirement == WorkerBinaryRequirement::LocalHost
723 && let Some((path, source)) = select_native_worker(current, matches_build)
724 {
725 return Ok(WorkerBinaryAvailability::Local {
726 path,
727 source: source.into(),
728 });
729 }
730 if !controller_replaced
731 && let Some((path, source)) = select_sibling_worker(current, &triple, matches_build)
732 {
733 return Ok(WorkerBinaryAvailability::Local {
734 path,
735 source: source.into(),
736 });
737 }
738 if let Some(template) = mj_core::config::env_override("WORKER_URL") {
739 let expected = mj_core::config::env_override("WORKER_SHA256")
740 .context("MJ_WORKER_URL requires MJ_WORKER_SHA256")?;
741 validate_worker_sha256(&expected)?;
742 return Ok(WorkerBinaryAvailability::Remote {
743 url: template.replace("{target}", &triple),
744 sha256: expected,
745 triple,
746 });
747 }
748 let rejected = rejected.into_inner();
749 ensure!(
750 rejected.is_empty(),
751 "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`).",
752 rejected.join("\n")
753 );
754 ensure!(
757 !controller_replaced,
758 "the running mj binary was replaced or removed on disk ({}); restart the Mjolnir daemon so it runs the current build, then retry",
759 display_path(current)
760 );
761 bail!(
762 "no 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"
763 )
764}