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 started = std::time::Instant::now();
363 let snapshot = capture_worker_binary_sources()?;
364 tracing::info!(
365 elapsed_ms = started.elapsed().as_millis(),
366 "worker sources pinned"
367 );
368 let _ = PINNED_WORKER_BINARY_SOURCES.set(snapshot);
371 Ok(())
372}
373
374pub fn warm_worker_binary_sources() -> Result<()> {
385 capture_worker_binary_sources().map(drop)
386}
387
388fn capture_worker_binary_sources() -> Result<WorkerBinarySourceSnapshot> {
389 let current = std::env::current_exe().context("resolve Mjolnir controller binary")?;
390 let cache_root = data_dir().join("workers").join("pinned");
391 Ok(WorkerBinarySourceSnapshot::capture(
392 &cache_root,
393 |arch, requirement| {
394 worker_binary_prerequisite_with_verifier(
395 arch,
396 requirement,
397 ¤t,
398 &|path| path.is_file(),
399 &|path| verify_worker_build_indexed(&cache_root, path),
400 )
401 },
402 ))
403}
404
405fn index_entry_path(cache_root: &Path, source: &Path) -> Option<PathBuf> {
412 let metadata = std::fs::metadata(source).ok()?;
413 let mut key = Sha256::new();
414 key.update(BUILD_ID.as_bytes());
415 key.update([0]);
416 key.update(source.as_os_str().as_encoded_bytes());
417 key.update(metadata.len().to_le_bytes());
418 #[cfg(unix)]
419 {
420 use std::os::unix::fs::MetadataExt;
421 for value in [
422 metadata.mtime(),
423 metadata.mtime_nsec(),
424 metadata.ctime(),
425 metadata.ctime_nsec(),
426 ] {
427 key.update(value.to_le_bytes());
428 }
429 key.update(metadata.ino().to_le_bytes());
430 }
431 #[cfg(not(unix))]
432 {
433 let modified = metadata.modified().ok()?;
434 key.update(format!("{modified:?}").as_bytes());
435 }
436 Some(cache_root.join("index").join(lower_hex(key.finalize())))
437}
438
439pub(super) fn indexed_pinned_worker(cache_root: &Path, source: &Path) -> Option<PathBuf> {
442 let entry = index_entry_path(cache_root, source)?;
443 let digest = std::fs::read_to_string(entry).ok()?;
444 let digest = digest.trim();
445 if digest.len() != 64 || !digest.bytes().all(|byte| byte.is_ascii_hexdigit()) {
446 return None;
447 }
448 let pinned = cache_root.join(digest).join("hel");
449 let source_len = std::fs::metadata(source).ok()?.len();
450 (std::fs::metadata(&pinned).ok()?.len() == source_len).then_some(pinned)
451}
452
453fn record_indexed_pin(cache_root: &Path, source: &Path, digest: &str) {
454 let Some(entry) = index_entry_path(cache_root, source) else {
455 return;
456 };
457 let write = || -> std::io::Result<()> {
459 let directory = entry.parent().expect("index entries have a parent");
460 std::fs::create_dir_all(directory)?;
461 let mut staged = tempfile::NamedTempFile::new_in(directory)?;
462 staged.write_all(digest.as_bytes())?;
463 staged
464 .persist(&entry)
465 .map(drop)
466 .map_err(|error| error.error)
467 };
468 if let Err(error) = write() {
469 tracing::debug!(source = %source.display(), %error, "could not index a pinned worker");
470 }
471}
472
473pub(super) fn verify_worker_build_indexed(cache_root: &Path, path: &Path) -> Result<()> {
477 if indexed_pinned_worker(cache_root, path).is_some() {
478 return Ok(());
479 }
480 verify_worker_build(path)
481}
482
483pub(super) fn copy_worker_source_to_cache(source: &Path, cache_root: &Path) -> Result<PathBuf> {
484 if let Some(pinned) = indexed_pinned_worker(cache_root, source) {
485 return Ok(pinned);
486 }
487 std::fs::create_dir_all(cache_root)
488 .with_context(|| format!("create pinned worker cache {}", cache_root.display()))?;
489 let mut input =
490 File::open(source).with_context(|| format!("open worker source {}", source.display()))?;
491 let metadata = input
492 .metadata()
493 .with_context(|| format!("stat worker source {}", source.display()))?;
494 let mut temporary = tempfile::NamedTempFile::new_in(cache_root)
495 .with_context(|| format!("create pinned worker staging file {}", cache_root.display()))?;
496 let mut digest = Sha256::new();
497 let mut buffer = [0_u8; 128 * 1024];
498 loop {
499 let count = input
500 .read(&mut buffer)
501 .with_context(|| format!("read worker source {}", source.display()))?;
502 if count == 0 {
503 break;
504 }
505 temporary
506 .write_all(&buffer[..count])
507 .with_context(|| format!("copy worker source {}", source.display()))?;
508 digest.update(&buffer[..count]);
509 }
510 temporary
511 .as_file_mut()
512 .sync_all()
513 .with_context(|| format!("flush pinned worker source {}", source.display()))?;
514 std::fs::set_permissions(temporary.path(), metadata.permissions())
515 .with_context(|| format!("preserve permissions for {}", source.display()))?;
516 let digest = lower_hex(digest.finalize());
517 let pinned = publish_cached_worker(temporary, cache_root, &digest)?;
518 record_indexed_pin(cache_root, source, &digest);
519 Ok(pinned)
520}
521
522pub(super) fn publish_cached_worker(
526 temporary: tempfile::NamedTempFile,
527 cache_root: &Path,
528 digest: &str,
529) -> Result<PathBuf> {
530 verify_worker_build(temporary.path())?;
531 let directory = cache_root.join(digest);
532 std::fs::create_dir_all(&directory)
533 .with_context(|| format!("create pinned worker cache {}", directory.display()))?;
534 let destination = directory.join("hel");
535 if destination.is_file() {
536 verify_cached_worker(&destination, digest)?;
537 return Ok(destination);
538 }
539 match temporary.persist_noclobber(&destination) {
540 Ok(_) => {
541 #[cfg(unix)]
542 File::open(&directory)
543 .and_then(|directory| directory.sync_all())
544 .with_context(|| format!("flush pinned worker cache {}", directory.display()))?;
545 Ok(destination)
546 }
547 Err(error) if error.error.kind() == ErrorKind::AlreadyExists => {
548 if destination.is_file() {
549 verify_cached_worker(&destination, digest)?;
550 Ok(destination)
551 } else {
552 Err(error.error).with_context(|| {
553 format!("publish pinned worker artifact {}", destination.display())
554 })
555 }
556 }
557 Err(error) => Err(error.error)
558 .with_context(|| format!("publish pinned worker artifact {}", destination.display())),
559 }
560}
561
562fn verify_cached_worker(path: &Path, digest: &str) -> Result<()> {
563 verify_worker_build(path)?;
564 ensure!(
565 mj_core::worker_launch::worker_executable_digest(path)? == digest,
566 "content-addressed worker cache {} does not match {digest} checksum",
567 path.display()
568 );
569 Ok(())
570}
571
572pub fn worker_binary_prerequisite_for_arch(arch: &str) -> Result<WorkerBinaryAvailability> {
579 worker_binary_for_arch(arch, WorkerBinaryRequirement::PortableLinux)
580}
581
582pub fn native_worker_binary_prerequisite() -> Result<WorkerBinaryAvailability> {
589 worker_binary_for_arch(std::env::consts::ARCH, WorkerBinaryRequirement::LocalHost)
590}
591
592pub(super) fn worker_binary_for_arch(
593 arch: &str,
594 requirement: WorkerBinaryRequirement,
595) -> Result<WorkerBinaryAvailability> {
596 if let Some(snapshot) = PINNED_WORKER_BINARY_SOURCES.get() {
597 let pinned = snapshot.resolve(arch, requirement);
598 if pinned_source_is_usable(&pinned, &|path| path.is_file()) {
599 return pinned;
600 }
601 return resolve_worker_source_again(arch, requirement, pinned.err());
607 }
608 let current = std::env::current_exe().context("resolve Mjolnir controller binary")?;
609 worker_binary_prerequisite_for_current(arch, requirement, ¤t, &|path| path.is_file())
610}
611
612pub(super) fn pinned_source_is_usable(
618 pinned: &Result<WorkerBinaryAvailability>,
619 is_file: &dyn Fn(&Path) -> bool,
620) -> bool {
621 match pinned {
622 Ok(WorkerBinaryAvailability::Local { path, .. }) => is_file(path),
623 Ok(WorkerBinaryAvailability::Remote { .. }) => true,
624 Err(_) => false,
625 }
626}
627
628fn resolve_worker_source_again(
633 arch: &str,
634 requirement: WorkerBinaryRequirement,
635 pinned_error: Option<anyhow::Error>,
636) -> Result<WorkerBinaryAvailability> {
637 let current = std::env::current_exe().context("resolve Mjolnir controller binary")?;
638 let resolved =
639 worker_binary_prerequisite_for_current(arch, requirement, ¤t, &|path| path.is_file());
640 match resolved {
641 Ok(WorkerBinaryAvailability::Local { path, source }) => {
642 let cache_root = data_dir().join("workers").join("pinned");
643 let path = copy_worker_source_to_cache(&path, &cache_root)
644 .context("pin the re-resolved worker source")?;
645 tracing::info!(
646 arch,
647 requirement = ?requirement,
648 source = %source,
649 "re-resolved a worker source the daemon could not pin at startup"
650 );
651 Ok(WorkerBinaryAvailability::Local { path, source })
652 }
653 Ok(remote) => Ok(remote),
654 Err(error) => Err(match pinned_error {
657 Some(pinned) => error.context(format!("{pinned:#}")),
658 None => error,
659 }),
660 }
661}
662
663pub(super) fn worker_binary_prerequisite_for_current(
666 arch: &str,
667 requirement: WorkerBinaryRequirement,
668 current: &Path,
669 is_file: &dyn Fn(&Path) -> bool,
670) -> Result<WorkerBinaryAvailability> {
671 worker_binary_prerequisite_with_verifier(
672 arch,
673 requirement,
674 current,
675 is_file,
676 &verify_worker_build,
677 )
678}
679
680pub(super) fn worker_binary_prerequisite_with_verifier(
684 arch: &str,
685 requirement: WorkerBinaryRequirement,
686 current: &Path,
687 is_file: &dyn Fn(&Path) -> bool,
688 verify: &dyn Fn(&Path) -> Result<()>,
689) -> Result<WorkerBinaryAvailability> {
690 let triple = requirement.triple(arch);
691 let rejected = std::cell::RefCell::new(Vec::new());
692 let matches_build = |path: &Path| {
693 if !is_file(path) {
694 return false;
695 }
696 match verify(path) {
697 Ok(()) => true,
698 Err(error) => {
699 tracing::warn!(path = %path.display(), error = %error, "skipping incompatible worker source");
700 rejected.borrow_mut().push(format!("{error:#}"));
701 false
702 }
703 }
704 };
705 if let Some(path) = mj_core::config::env_override_os("WORKER_BINARY").map(PathBuf::from) {
706 if !is_file(&path) {
707 bail!("MJ_WORKER_BINARY is not a file: {}", path.display());
708 }
709 verify(&path).context(
712 "MJ_WORKER_BINARY does not match the running mj; rebuild it from the same commit or unset MJ_WORKER_BINARY",
713 )?;
714 return Ok(WorkerBinaryAvailability::Local {
715 path,
716 source: "MJ_WORKER_BINARY".into(),
717 });
718 }
719 let controller_replaced = !is_file(current);
723 let mut candidates = Vec::new();
724 if let Some(directory) = mj_core::config::env_override_os("WORKER_DIR").map(PathBuf::from) {
725 candidates.push((
726 packaged_worker_binary_path(&directory, &triple),
727 "MJ_WORKER_DIR",
728 ));
729 candidates.push((directory.join(&triple).join("hel"), "MJ_WORKER_DIR"));
730 if triple.ends_with("-apple-darwin") {
731 candidates.push((
732 packaged_worker_binary_path(&directory, "universal-apple-darwin"),
733 "MJ_WORKER_DIR",
734 ));
735 }
736 }
737 if let Some((path, source)) = candidates.into_iter().find(|(path, _)| matches_build(path)) {
738 return Ok(WorkerBinaryAvailability::Local {
739 path,
740 source: source.into(),
741 });
742 }
743 if requirement == WorkerBinaryRequirement::LocalHost
744 && let Some((path, source)) = select_native_worker(current, matches_build)
745 {
746 return Ok(WorkerBinaryAvailability::Local {
747 path,
748 source: source.into(),
749 });
750 }
751 if !controller_replaced
752 && let Some((path, source)) = select_sibling_worker(current, &triple, matches_build)
753 {
754 return Ok(WorkerBinaryAvailability::Local {
755 path,
756 source: source.into(),
757 });
758 }
759 if let Some(template) = mj_core::config::env_override("WORKER_URL") {
760 let expected = mj_core::config::env_override("WORKER_SHA256")
761 .context("MJ_WORKER_URL requires MJ_WORKER_SHA256")?;
762 validate_worker_sha256(&expected)?;
763 return Ok(WorkerBinaryAvailability::Remote {
764 url: template.replace("{target}", &triple),
765 sha256: expected,
766 triple,
767 });
768 }
769 let rejected = rejected.into_inner();
770 ensure!(
771 rejected.is_empty(),
772 "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`).",
773 rejected.join("\n")
774 );
775 ensure!(
778 !controller_replaced,
779 "the running mj binary was replaced or removed on disk ({}); restart the Mjolnir daemon so it runs the current build, then retry",
780 display_path(current)
781 );
782 bail!(
783 "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"
784 )
785}