1use super::*;
2
3impl Controller {
4 pub(in crate::controller) fn worker_placement(
8 &self,
9 session_id: &str,
10 ) -> Result<(targets::TargetLocator, String)> {
11 let session = self
12 .state
13 .sessions
14 .get(session_id)
15 .with_context(|| format!("unknown session {session_id}"))?;
16 let locator = session
17 .target
18 .as_ref()
19 .context("session target is missing")?;
20 let backend = backend_locator(locator, session, &self.config)?;
21 let worker_root = targets::worker_root(&backend, session_id)?;
22 Ok((backend, worker_root))
23 }
24
25 pub(in crate::controller) fn prepare_worker_files(
26 &self,
27 session_id: &str,
28 backend: &targets::TargetLocator,
29 worker_root: &str,
30 executor: &impl CommandExecutor,
31 ) -> Result<()> {
32 let session = self
33 .state
34 .sessions
35 .get(session_id)
36 .with_context(|| format!("unknown session {session_id}"))?;
37 let profile = self
38 .config
39 .profiles
40 .get(&session.last_profile)
41 .context("session profile is missing")?;
42 let (mut launch, project_memory, target_profile_home) =
43 self.session_launch_config(session_id, backend)?;
44
45 if session.native_session_id.is_some()
46 && profile.kind == mj_core::config::HarnessKind::Codex
47 {
48 launch.goal_resume_request = Some(mj_core::state::new_session_id()?);
49 }
50 let staging = tempfile::tempdir().context("create worker staging directory")?;
51 let launch_path = staging.path().join("launch.json");
52 launch.write(&launch_path)?;
53 let ownership_path = staging.path().join("ownership.json");
54 WorkerOwnership {
55 version: WorkerOwnership::VERSION,
56 workspace_id: session.workspace_id.clone(),
57 session_id: session_id.to_string(),
58 profile_id: session.last_profile.clone(),
59 bundle_id: session.bundle_id.clone(),
60 target_template_id: session.target_template_id.clone(),
61 instance_id: Some(mj_core::config::instance_identity()),
62 }
63 .write(&ownership_path)?;
64 let profile_stage = staging.path().join("profile");
68 let started = Instant::now();
69 let result = stage_profile(profile, &profile_stage);
70 tracing::debug!(
71 session_id,
72 elapsed_ms = started.elapsed().as_millis(),
73 "profile staging completed"
74 );
75 result?;
76 stage_managed_skills(profile.kind, &profile_stage)?;
77 stage_codex_catalog(
78 &session.last_profile,
79 profile,
80 &profile_stage,
81 &fetch_catalog_over_https,
82 &SharedCatalogCache,
83 )?;
84 append_hel_target_environment(profile.kind, &profile_stage, backend)?;
85 append_subagent_policy(
86 profile.kind,
87 &profile_stage,
88 &launch.subagents,
89 self.config.subagents.max_concurrent,
90 )?;
91 apply_staged_execution_setting(profile.kind, launch.execution_policy, &profile_stage)?;
92 if profile.kind == mj_core::config::HarnessKind::Claude {
93 if let Some(role) = launch.subagents.parent_role() {
94 configure_claude_subagent_mcp(&profile_stage, worker_root, role)?;
95 } else if launch.handback_tool {
96 configure_claude_subagent_mcp(
97 &profile_stage,
98 worker_root,
99 mj_core::subagent::SubagentMcpRole::Child,
100 )?;
101 }
102 }
103 stage_memory_replica(
104 &project_memory,
105 Path::new(&target_profile_home),
106 &profile_stage,
107 )?;
108 if project_memory.mcp_delivery == ProjectMemoryMcpDelivery::HarnessProfile {
109 configure_kimi_project_memory_mcp(&profile_stage, worker_root, &project_memory)?;
110 }
111 let worker_binary = worker_binary_for(backend, executor)?;
112
113 install_worker_files(
114 executor,
115 backend,
116 session_id,
117 worker_root,
118 &target_profile_home,
119 &worker_binary,
120 &launch_path,
121 &ownership_path,
122 &profile_stage,
123 )?;
124 if session.build_cache.is_some() {
125 self.install_build_cache_shim(session, backend, &launch, executor)
126 .context("install shared machine build cache configuration")?;
127 }
128 prepare_installed_managed_harness(executor, backend, worker_root, &launch)
129 }
130
131 pub(in crate::controller) fn install_build_cache_shim(
136 &self,
137 session: &mj_core::state::SessionRecord,
138 backend: &targets::TargetLocator,
139 launch: &WorkerLaunchConfig,
140 executor: &impl CommandExecutor,
141 ) -> Result<()> {
142 let worker_root = targets::worker_root(backend, &session.id)?;
143 let binary =
147 crate::controller::mbx::binary_for(backend, executor).inspect_err(|error| {
148 executor.notify_notice(&format!(
149 "The Rust build cache could not be prepared: {error:#}."
150 ));
151 })?;
152 let (configuration, config_roots) =
153 self.build_cache_configuration(session, backend, launch, executor)?;
154 install_mbx_files(
155 executor,
156 backend,
157 &session.id,
158 &worker_root,
159 &binary,
160 &configuration,
161 &config_roots,
162 )
163 }
164
165 pub(in crate::controller) fn prepare_build_cache_links(
166 &self,
167 session: &mj_core::state::SessionRecord,
168 backend: &targets::TargetLocator,
169 launch: &WorkerLaunchConfig,
170 executor: &impl CommandExecutor,
171 ) -> Result<()> {
172 let (configuration, config_roots) =
173 self.build_cache_configuration(session, backend, launch, executor)?;
174 link_mbx_configuration(executor, backend, &configuration, &config_roots)
175 }
176
177 fn build_cache_configuration(
178 &self,
179 session: &mj_core::state::SessionRecord,
180 backend: &targets::TargetLocator,
181 launch: &WorkerLaunchConfig,
182 executor: &impl CommandExecutor,
183 ) -> Result<(PathBuf, Vec<PathBuf>)> {
184 let configuration = crate::controller::mbx::prepare_session_configuration(
185 &self.config,
186 backend,
187 session
188 .build_cache
189 .as_ref()
190 .context("session has no build cache")?,
191 executor,
192 )?;
193 let config_roots = [&launch.target_environment, &launch.environment]
194 .into_iter()
195 .filter_map(|environment| environment.get("XDG_CONFIG_HOME").map(PathBuf::from))
196 .collect::<std::collections::BTreeSet<_>>();
197 Ok((configuration, config_roots.into_iter().collect()))
198 }
199
200 pub fn diagnose_worker(&self, session_id: &str) -> Option<String> {
204 self.diagnose_worker_controlled(session_id, &crate::targets::ProcessExecutor)
205 }
206
207 pub fn diagnose_worker_controlled(
208 &self,
209 session_id: &str,
210 executor: &impl CommandExecutor,
211 ) -> Option<String> {
212 let session = self.state.sessions.get(session_id)?;
213 let locator = session.target.as_ref()?;
214 let backend = match backend_locator(locator, session, &self.config) {
215 Ok(backend) => backend,
216 Err(error) => {
217 tracing::debug!(
218 session_id,
219 error = format!("{error:#}"),
220 "could not construct a worker diagnostic probe"
221 );
222 return None;
223 }
224 };
225 let worker_root = match targets::worker_root(&backend, session_id) {
226 Ok(root) => root,
227 Err(error) => {
228 tracing::debug!(
229 session_id,
230 error = format!("{error:#}"),
231 "could not derive the worker diagnostic root"
232 );
233 return None;
234 }
235 };
236 let binary_failure = worker_binary_probe_failure(executor, &backend, &worker_root);
237 let probe = match probe_worker(executor, &backend, &worker_root) {
238 Ok(probe) => probe.to_string(),
239 Err(error) => format!("the worker could not be probed: {error:#}"),
240 };
241 Some(match binary_failure {
242 Some(binary_failure) => format!("{binary_failure}; {probe}"),
243 None => probe,
244 })
245 }
246
247 pub fn worker_recovery_plan(
251 &self,
252 session_id: &str,
253 operation: Option<&mj_core::state::MoveOperation>,
254 ) -> Result<WorkerRecoveryPlan> {
255 let (backend, worker_root) = self.worker_placement(session_id)?;
256 let launch = self.worker_launch_config_for_move(session_id, &backend, operation)?;
257 let workspace = worker_workspace_for_recovery(&backend, &launch.cwd);
258 Ok(WorkerRecoveryPlan {
259 source_target: self.state.sessions[session_id]
260 .target
261 .clone()
262 .context("session target is missing")?,
263 target: targets::target_recovery_plan(&backend, session_id)?,
264 workspace,
265 liveness_probe: worker_liveness_command(&backend, &worker_root),
266 binary_refresh: worker_binary_refresh_plan(&backend, session_id)?,
267 launch_refresh: Some(worker_launch_refresh_plan(&backend, session_id, &launch)?),
268 restart: CommandPlan {
269 description: format!("restart Mjolnir worker for session {session_id}"),
270 commands: vec![
271 stop_worker_command(&backend, &worker_root),
272 start_worker_command(&backend, &worker_root),
273 ],
274 },
275 })
276 }
277
278 fn session_launch_config(
283 &self,
284 session_id: &str,
285 backend: &targets::TargetLocator,
286 ) -> Result<(WorkerLaunchConfig, ProjectMemoryLaunchConfig, String)> {
287 let session = self
288 .state
289 .sessions
290 .get(session_id)
291 .with_context(|| format!("unknown session {session_id}"))?;
292 session.validate_configuration(&self.config)?;
293 let profile = self
294 .config
295 .profiles
296 .get(&session.last_profile)
297 .context("session profile is missing")?;
298 let bundle = session
299 .project_directory
300 .is_none()
301 .then(|| self.config.bundles.get(&session.bundle_id))
302 .flatten();
303 let target = session.target_runtime_settings(&self.config)?;
304 let subagent = self.state.subagents.get(session_id);
305 let (workspace_session_id, workspace_container) = match subagent.as_ref() {
308 Some(child) => {
309 let parent = self
310 .state
311 .sessions
312 .get(&child.parent_session_id)
313 .context("sub-agent parent session is missing")?;
314 (parent.id.clone(), parent.container_workspace.clone())
315 }
316 None => (session_id.to_owned(), session.container_workspace.clone()),
317 };
318 let (mut launch, project_memory, target_profile_home) = worker_launch_config(
319 session,
320 profile,
321 bundle,
322 backend,
323 LaunchWorkspace {
324 session_id: &workspace_session_id,
325 container: workspace_container.as_deref(),
326 parent_worktree: self.subagent_parent_worktree(session_id),
327 },
328 &target,
329 )?;
330 apply_jev_switch(&mut launch, self.config.jev.enabled);
331 apply_continuation_switch(&mut launch, self.config.automatic_continuation_enabled());
332 launch.subagents = session
333 .subagents
334 .clone()
335 .unwrap_or_default()
336 .for_launch(profile.kind, subagent.is_some());
337 launch.handback_tool = subagent.as_ref().is_some_and(|child| child.handback_tool);
339 launch.review_capture =
345 mj_core::review::settings::can_review(&self.config) && subagent.is_none();
346 launch.bifrost_binary = mj_review::bifrost::configured_bifrost_binary();
347 if let Some(subagent) = &subagent {
348 let parent = self
349 .state
350 .sessions
351 .get(&subagent.parent_session_id)
352 .context("sub-agent parent session is missing")?;
353 let parent_profile = self
354 .config
355 .profiles
356 .get(&parent.last_profile)
357 .context("sub-agent parent profile is missing")?;
358 let parent_target = parent.target_runtime_settings(&self.config)?;
359 let parent_locator = parent
360 .target
361 .as_ref()
362 .context("sub-agent parent has no live target")?;
363 let parent_backend = backend_locator(parent_locator, parent, &self.config)?;
364 let parent_bundle = parent
365 .project_directory
366 .is_none()
367 .then(|| self.config.bundles.get(&parent.bundle_id))
368 .flatten();
369 let (parent_launch, _, _) = worker_launch_config(
370 parent,
371 parent_profile,
372 parent_bundle,
373 &parent_backend,
374 LaunchWorkspace {
375 session_id: &parent.id,
376 container: parent.container_workspace.as_deref(),
377 parent_worktree: None,
378 },
379 &parent_target,
380 )?;
381 launch.cwd = if subagent.working_directory.as_os_str().is_empty() {
382 parent_launch.cwd
383 } else {
384 parent_launch.cwd.join(&subagent.working_directory)
385 };
386 launch.additional_directories = parent_launch.additional_directories;
387 }
388 Ok((launch, project_memory, target_profile_home))
389 }
390
391 pub(in crate::controller) fn current_worker_launch_config(
392 &self,
393 session_id: &str,
394 backend: &targets::TargetLocator,
395 ) -> Result<WorkerLaunchConfig> {
396 let operation = crate::database::load_move_operation(session_id)?;
397 self.worker_launch_config_for_move(session_id, backend, operation.as_ref())
398 }
399
400 fn worker_launch_config_for_move(
401 &self,
402 session_id: &str,
403 backend: &targets::TargetLocator,
404 operation: Option<&mj_core::state::MoveOperation>,
405 ) -> Result<WorkerLaunchConfig> {
406 let session = self
407 .state
408 .sessions
409 .get(session_id)
410 .with_context(|| format!("unknown session {session_id}"))?;
411 let (mut launch, _, _) = self.session_launch_config(session_id, backend)?;
412 if operation.is_some_and(|operation| {
413 operation.source_checkpoint_only
414 && operation.destination_target.is_none()
415 && matches!(
416 operation.phase,
417 mj_core::state::MovePhase::Preparing
418 | mj_core::state::MovePhase::ClosingSource
419 | mj_core::state::MovePhase::Failed
420 | mj_core::state::MovePhase::Cancelled
421 )
422 && session.last_profile == operation.source_profile_id
423 && session.target == operation.source_target
424 && matches!(
425 session.state,
426 mj_core::state::SessionState::Running
427 | mj_core::state::SessionState::Disconnected
428 | mj_core::state::SessionState::Closing
429 )
430 }) {
431 launch.run_mode = mj_core::worker_launch::WorkerRunMode::CheckpointOnly;
432 }
433 Ok(launch)
434 }
435
436 pub fn project_memory_sync_target(&self, session_id: &str) -> Result<ProjectMemorySyncTarget> {
437 let session = self
438 .state
439 .sessions
440 .get(session_id)
441 .with_context(|| format!("unknown session {session_id}"))?;
442 session.validate_configuration(&self.config)?;
443 let locator = session
444 .target
445 .as_ref()
446 .context("session target is missing")?;
447 let backend = backend_locator(locator, session, &self.config)?;
448 let profile = self
449 .config
450 .profiles
451 .get(&session.last_profile)
452 .context("session profile is missing")?;
453 let bundle = session
454 .project_directory
455 .is_none()
456 .then(|| self.config.bundles.get(&session.bundle_id))
457 .flatten();
458 let workspace = if let Some(project_directory) = &session.project_directory {
459 (project_directory.to_string_lossy().into_owned(), Vec::new())
460 } else {
461 workspace_paths(
462 &backend,
463 bundle.context("session bundle is missing")?,
464 session_id,
465 session.container_workspace.as_deref(),
466 )?
467 };
468 let target_home = target_profile_home(&backend, session_id, profile);
469 let launch = project_memory_launch(
470 session,
471 bundle,
472 &workspace,
473 &target_home,
474 self.subagent_parent_worktree(session_id),
475 )?;
476 Ok(ProjectMemorySyncTarget {
477 canonical_root: canonical_memory_root(&launch.project_key),
478 })
479 }
480}
481
482pub(super) fn apply_jev_switch(launch: &mut WorkerLaunchConfig, enabled: bool) {
486 if enabled {
487 return;
488 }
489 for environment in [&mut launch.target_environment, &mut launch.environment] {
490 environment.remove("TYPESAFE_API_KEY");
491 environment.insert(
492 mj_core::jev::DISABLED_ENVIRONMENT.to_owned(),
493 "1".to_owned(),
494 );
495 }
496}
497
498pub(super) fn apply_continuation_switch(launch: &mut WorkerLaunchConfig, enabled: bool) {
501 if enabled {
502 return;
503 }
504 for environment in [&mut launch.target_environment, &mut launch.environment] {
505 environment.insert(
506 mj_core::jev::CONTINUATION_DISABLED_ENVIRONMENT.to_owned(),
507 "1".to_owned(),
508 );
509 }
510}
511
512#[cfg(test)]
519pub(super) fn subagent_tools_enabled(
520 session: &mj_core::state::SessionRecord,
521 is_child: bool,
522) -> bool {
523 session
524 .subagents
525 .clone()
526 .unwrap_or_default()
527 .for_launch(session.harness_kind, is_child)
528 .uses_mjolnir()
529}
530
531pub(super) fn worker_workspace_for_recovery(
532 backend: &targets::TargetLocator,
533 directory: &Path,
534) -> Option<WorkerWorkspace> {
535 let target = match backend {
536 targets::TargetLocator::LocalBare { .. } => mj_core::state::ManagedWorktreeTarget::Local,
537 targets::TargetLocator::SshBare { ssh, .. } => mj_core::state::ManagedWorktreeTarget::Ssh {
538 destination: ssh.destination.clone(),
539 ssh_args: ssh.ssh_args.clone(),
540 },
541 targets::TargetLocator::LocalPodman { .. }
542 | targets::TargetLocator::LocalDocker { .. }
543 | targets::TargetLocator::AppleContainer { .. }
544 | targets::TargetLocator::AwsEc2 { .. }
545 | targets::TargetLocator::SshPodman { .. }
546 | targets::TargetLocator::SshDocker { .. } => return None,
547 };
548 Some(WorkerWorkspace {
549 target,
550 directory: directory.to_path_buf(),
551 })
552}
553
554pub(super) struct LaunchWorkspace<'a> {
555 pub session_id: &'a str,
556 pub container: Option<&'a Path>,
557 pub parent_worktree: Option<&'a mj_core::state::ManagedWorktree>,
558}
559
560pub(super) fn worker_launch_config(
561 session: &mj_core::state::SessionRecord,
562 profile: &mj_core::config::HarnessProfile,
563 bundle: Option<&ProjectBundle>,
564 backend: &targets::TargetLocator,
565 worker_workspace: LaunchWorkspace<'_>,
566 target: &mj_core::state::TargetRuntimeSettings,
567) -> Result<(WorkerLaunchConfig, ProjectMemoryLaunchConfig, String)> {
568 let session_id = session.id.as_str();
569 let execution_policy = profile
570 .kind
571 .effective_execution_policy(target.execution_policy);
572 let target_profile_home = target_profile_home(backend, session_id, profile);
573 let workspace = if let Some(project_directory) = &session.project_directory {
574 (project_directory.to_string_lossy().into_owned(), Vec::new())
575 } else {
576 workspace_paths(
577 backend,
578 bundle.context("session bundle is missing")?,
579 worker_workspace.session_id,
580 worker_workspace.container,
581 )?
582 };
583 let mut additional_directories = workspace.1.iter().map(PathBuf::from).collect::<Vec<_>>();
584 additional_directories.extend(
585 session
586 .additional_mounts
587 .iter()
588 .map(|resource| resource.destination.clone()),
589 );
590 if profile.kind == mj_core::config::HarnessKind::Muse && !additional_directories.is_empty() {
591 bail!(
592 "{} ACP does not support multiple workspace roots; use a single-repository bundle",
593 profile.kind.display_name()
594 );
595 }
596 let (bridge_command, bridge_args) = bridge_launch(profile.kind, execution_policy);
597 let mut target_environment = target.environment.clone();
598 for name in [
608 "MJ_TURN_STALL_TIMEOUT_MS",
609 "MJ_TURN_TOOL_STALL_TIMEOUT_MS",
610 "RUST_LOG",
611 ] {
612 if let Ok(value) = std::env::var(name) {
613 target_environment.insert(name.to_owned(), value);
614 }
615 }
616 if let Some(key) = mj_core::activity::verdict::api_key() {
619 target_environment.insert("TYPESAFE_API_KEY".to_owned(), key);
620 }
621 if let Some(build_cache) = &session.build_cache {
624 target_environment.insert(
625 "MBX_CACHE_DIR".into(),
626 build_cache.directory.to_string_lossy().into_owned(),
627 );
628 target_environment.insert(
629 "MJ_MBX_CONFIG_DIR".into(),
630 mj_core::config::build_cache_configuration_directory(&build_cache.directory)
631 .to_string_lossy()
632 .into_owned(),
633 );
634 target_environment.insert(
637 "MBX_SHIMS_DIR".into(),
638 Path::new(&targets::worker_root(backend, session_id)?)
639 .join("mbx-shims")
640 .to_string_lossy()
641 .into_owned(),
642 );
643 target_environment.insert("MBX_SUMMARY".into(), "off".into());
646 target_environment.insert("MBX_SAVINGS".into(), "off".into());
647 }
648 let mut environment = target_environment.clone();
649 environment.extend(profile.environment.resolved().clone());
650 profile
651 .kind
652 .configure_home_environment(Path::new(&target_profile_home), &mut environment);
653 profile
654 .kind
655 .configure_execution_environment(execution_policy, &mut environment)?;
656 let mut project_memory = project_memory_launch(
657 session,
658 bundle,
659 &workspace,
660 &target_profile_home,
661 worker_workspace.parent_worktree,
662 )?;
663 project_memory.mcp_delivery = project_memory_mcp_delivery(profile.kind, backend);
664 project_memory.history_socket =
665 Some(Path::new(&targets::worker_root(backend, &session.id)?).join("control.sock"));
666 if profile.kind == mj_core::config::HarnessKind::Claude {
667 environment.insert(
668 "CLAUDE_CODE_PROJECT_DIR_NAME".into(),
669 project_memory_replica_slug(&project_memory.project_key, session_id),
670 );
671 }
672 apply_claude_setup_token(
673 &mut environment,
674 profile.kind,
675 &mj_core::credentials::claude_oauth_token_path(&session.last_profile),
676 );
677 let excluded_environment =
678 exclude_harness_environment(&session.last_profile, profile, &mut environment);
679 Ok((
680 WorkerLaunchConfig {
681 goal_resume_request: None,
682 target_environment,
683 seed_image_environment: backend.container_engine().is_some(),
684 run_mode: Default::default(),
685 expected_runtime_identity: session.expected_runtime_identity.clone(),
686 session_id: session_id.to_string(),
687 subagents: mj_core::subagent::SubagentPolicy::Native,
688 handback_tool: false,
689 review_capture: false,
690 bifrost_binary: None,
691 harness: profile.kind,
692 harness_home: PathBuf::from(&target_profile_home),
693 authentication_marker: profile
699 .authentication_marker()
700 .strip_prefix(&profile.home)
701 .ok()
702 .map(|name| name.to_string_lossy().into_owned()),
703 bridge_command: PathBuf::from(bridge_command),
704 bridge_args,
705 harness_runtime: harness_runtime_policy(backend),
706 environment,
707 excluded_environment,
708 cwd: PathBuf::from(&workspace.0),
709 additional_directories,
710 native_session_id: session.native_session_id.clone(),
711 project_memory: Some(project_memory.clone()),
712 execution_policy,
713 },
714 project_memory,
715 target_profile_home,
716 ))
717}
718
719pub(super) fn exclude_harness_environment(
726 profile_id: &str,
727 profile: &mj_core::config::HarnessProfile,
728 environment: &mut std::collections::BTreeMap<String, String>,
729) -> Vec<String> {
730 static REPORTED: std::sync::Mutex<std::collections::BTreeSet<String>> =
731 std::sync::Mutex::new(std::collections::BTreeSet::new());
732 let before = environment.clone();
733 let excluded = profile.exclude_harness_environment(environment);
734 let removed = excluded
735 .iter()
736 .filter(|name| before.contains_key(*name))
737 .cloned()
738 .collect::<Vec<_>>();
739 if !removed.is_empty()
740 && REPORTED
741 .lock()
742 .map(|mut reported| reported.insert(profile_id.to_owned()))
743 .unwrap_or(true)
744 {
745 tracing::info!(
746 profile_id,
747 removed = removed.join(", "),
748 "left API key settings out of the harness environment: this Codex profile signs in with ChatGPT and must not fall back to an API key"
749 );
750 }
751 excluded
752}
753
754pub(super) fn harness_runtime_policy(backend: &targets::TargetLocator) -> HarnessRuntimePolicy {
755 match backend {
756 targets::TargetLocator::LocalBare { .. }
757 | targets::TargetLocator::AwsEc2 { .. }
758 | targets::TargetLocator::SshBare { .. } => HarnessRuntimePolicy::Managed,
759 _ => HarnessRuntimePolicy::Ambient,
760 }
761}