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 apply_staged_execution_setting(profile.kind, launch.execution_policy, &profile_stage)?;
86 if profile.kind == mj_core::config::HarnessKind::Claude {
87 if launch.subagent_tools {
88 configure_claude_subagent_mcp(
89 &profile_stage,
90 worker_root,
91 mj_core::subagent::SubagentMcpRole::Parent,
92 )?;
93 } else if launch.handback_tool {
94 configure_claude_subagent_mcp(
95 &profile_stage,
96 worker_root,
97 mj_core::subagent::SubagentMcpRole::Child,
98 )?;
99 }
100 }
101 stage_memory_replica(
102 &project_memory,
103 Path::new(&target_profile_home),
104 &profile_stage,
105 )?;
106 if project_memory.mcp_delivery == ProjectMemoryMcpDelivery::HarnessProfile {
107 configure_kimi_project_memory_mcp(&profile_stage, worker_root, &project_memory)?;
108 }
109 let worker_binary = worker_binary_for(backend, executor)?;
110
111 install_worker_files(
112 executor,
113 backend,
114 session_id,
115 worker_root,
116 &target_profile_home,
117 &worker_binary,
118 &launch_path,
119 &ownership_path,
120 &profile_stage,
121 )?;
122 if session.build_cache.is_some()
125 && let Err(error) = self.install_build_cache_shim(session, backend, executor)
126 {
127 tracing::warn!(
128 session_id,
129 "installing the mbx build cache failed: {error:#}"
130 );
131 }
132 prepare_installed_managed_harness(executor, backend, worker_root, &launch)
133 }
134
135 fn install_build_cache_shim(
140 &self,
141 session: &mj_core::state::SessionRecord,
142 backend: &targets::TargetLocator,
143 executor: &impl CommandExecutor,
144 ) -> Result<()> {
145 let worker_root = targets::worker_root(backend, &session.id)?;
146 let binary =
150 crate::controller::mbx::binary_for(backend, executor).inspect_err(|error| {
151 executor.notify_notice(&format!(
152 "The Rust build cache is unavailable: {error:#}; this session builds without it."
153 ));
154 })?;
155 let configuration = crate::controller::mbx::host_configuration(backend, executor)?;
156 install_mbx_files(
157 executor,
158 backend,
159 &session.id,
160 &worker_root,
161 &binary,
162 configuration.as_deref(),
163 )
164 }
165
166 pub fn diagnose_worker(&self, session_id: &str) -> Option<String> {
170 self.diagnose_worker_controlled(session_id, &crate::targets::ProcessExecutor)
171 }
172
173 pub fn diagnose_worker_controlled(
174 &self,
175 session_id: &str,
176 executor: &impl CommandExecutor,
177 ) -> Option<String> {
178 let session = self.state.sessions.get(session_id)?;
179 let locator = session.target.as_ref()?;
180 let backend = match backend_locator(locator, session, &self.config) {
181 Ok(backend) => backend,
182 Err(error) => {
183 tracing::debug!(
184 session_id,
185 error = format!("{error:#}"),
186 "could not construct a worker diagnostic probe"
187 );
188 return None;
189 }
190 };
191 let worker_root = match targets::worker_root(&backend, session_id) {
192 Ok(root) => root,
193 Err(error) => {
194 tracing::debug!(
195 session_id,
196 error = format!("{error:#}"),
197 "could not derive the worker diagnostic root"
198 );
199 return None;
200 }
201 };
202 let binary_failure = worker_binary_probe_failure(executor, &backend, &worker_root);
203 let last_words = worker_last_words(executor, &backend, &worker_root);
204 match (binary_failure, last_words) {
205 (Some(binary_failure), Some(last_words)) => {
206 Some(format!("{binary_failure}; {last_words}"))
207 }
208 (Some(binary_failure), None) => Some(binary_failure),
209 (None, last_words) => last_words,
210 }
211 }
212
213 pub fn worker_recovery_plan(&self, session_id: &str) -> Result<WorkerRecoveryPlan> {
217 let (backend, worker_root) = self.worker_placement(session_id)?;
218 let launch = self.current_worker_launch_config(session_id, &backend)?;
219 let workspace = worker_workspace_for_recovery(&backend, &launch.cwd);
220 Ok(WorkerRecoveryPlan {
221 source_target: self.state.sessions[session_id]
222 .target
223 .clone()
224 .context("session target is missing")?,
225 target: targets::target_recovery_plan(&backend, session_id)?,
226 workspace,
227 liveness_probe: worker_liveness_command(&backend, &worker_root),
228 binary_refresh: worker_binary_refresh_plan(&backend, session_id)?,
229 launch_refresh: Some(worker_launch_refresh_plan(&backend, session_id, &launch)?),
230 restart: CommandPlan {
231 description: format!("restart Mjolnir worker for session {session_id}"),
232 commands: vec![
233 stop_worker_command(&backend, &worker_root),
234 start_worker_command(&backend, &worker_root),
235 ],
236 },
237 })
238 }
239
240 fn session_launch_config(
245 &self,
246 session_id: &str,
247 backend: &targets::TargetLocator,
248 ) -> Result<(WorkerLaunchConfig, ProjectMemoryLaunchConfig, String)> {
249 let session = self
250 .state
251 .sessions
252 .get(session_id)
253 .with_context(|| format!("unknown session {session_id}"))?;
254 session.validate_configuration(&self.config)?;
255 let profile = self
256 .config
257 .profiles
258 .get(&session.last_profile)
259 .context("session profile is missing")?;
260 let bundle = session
261 .project_directory
262 .is_none()
263 .then(|| self.config.bundles.get(&session.bundle_id))
264 .flatten();
265 let target = session.target_runtime_settings(&self.config)?;
266 let subagent = crate::database::load_subagent(session_id)?;
267 let (workspace_session_id, workspace_container) = match subagent.as_ref() {
270 Some(child) => {
271 let parent = self
272 .state
273 .sessions
274 .get(&child.parent_session_id)
275 .context("sub-agent parent session is missing")?;
276 (parent.id.clone(), parent.container_workspace.clone())
277 }
278 None => (session_id.to_owned(), session.container_workspace.clone()),
279 };
280 let (mut launch, project_memory, target_profile_home) = worker_launch_config(
281 session,
282 profile,
283 bundle,
284 backend,
285 &workspace_session_id,
286 workspace_container.as_deref(),
287 &target,
288 )?;
289 apply_jev_switch(&mut launch, self.config.jev.enabled);
290 launch.subagent_tools = subagent_tools_enabled(session, subagent.is_some());
291 launch.handback_tool = subagent.as_ref().is_some_and(|child| child.handback_tool);
293 launch.review_capture =
299 mj_core::review::settings::can_review(&self.config) && subagent.is_none();
300 if let Some(subagent) = &subagent {
301 let parent = self
302 .state
303 .sessions
304 .get(&subagent.parent_session_id)
305 .context("sub-agent parent session is missing")?;
306 let parent_profile = self
307 .config
308 .profiles
309 .get(&parent.last_profile)
310 .context("sub-agent parent profile is missing")?;
311 let parent_target = parent.target_runtime_settings(&self.config)?;
312 let parent_locator = parent
313 .target
314 .as_ref()
315 .context("sub-agent parent has no live target")?;
316 let parent_backend = backend_locator(parent_locator, parent, &self.config)?;
317 let parent_bundle = parent
318 .project_directory
319 .is_none()
320 .then(|| self.config.bundles.get(&parent.bundle_id))
321 .flatten();
322 let (parent_launch, _, _) = worker_launch_config(
323 parent,
324 parent_profile,
325 parent_bundle,
326 &parent_backend,
327 &parent.id,
328 parent.container_workspace.as_deref(),
329 &parent_target,
330 )?;
331 launch.cwd = if subagent.working_directory.as_os_str().is_empty() {
332 parent_launch.cwd
333 } else {
334 parent_launch.cwd.join(&subagent.working_directory)
335 };
336 launch.additional_directories = parent_launch.additional_directories;
337 }
338 Ok((launch, project_memory, target_profile_home))
339 }
340
341 pub(in crate::controller) fn current_worker_launch_config(
342 &self,
343 session_id: &str,
344 backend: &targets::TargetLocator,
345 ) -> Result<WorkerLaunchConfig> {
346 let session = self
347 .state
348 .sessions
349 .get(session_id)
350 .with_context(|| format!("unknown session {session_id}"))?;
351 let (mut launch, _, _) = self.session_launch_config(session_id, backend)?;
352 if crate::database::load_move_operation(session_id)?.is_some_and(|operation| {
353 operation.source_checkpoint_only
354 && operation.destination_target.is_none()
355 && matches!(
356 operation.phase,
357 mj_core::state::MovePhase::Preparing
358 | mj_core::state::MovePhase::ClosingSource
359 | mj_core::state::MovePhase::Failed
360 | mj_core::state::MovePhase::Cancelled
361 )
362 && session.last_profile == operation.source_profile_id
363 && session.target == operation.source_target
364 && matches!(
365 session.state,
366 mj_core::state::SessionState::Running
367 | mj_core::state::SessionState::Disconnected
368 | mj_core::state::SessionState::Closing
369 )
370 }) {
371 launch.run_mode = mj_core::worker_launch::WorkerRunMode::CheckpointOnly;
372 }
373 Ok(launch)
374 }
375
376 pub fn project_memory_sync_target(&self, session_id: &str) -> Result<ProjectMemorySyncTarget> {
377 let session = self
378 .state
379 .sessions
380 .get(session_id)
381 .with_context(|| format!("unknown session {session_id}"))?;
382 session.validate_configuration(&self.config)?;
383 let locator = session
384 .target
385 .as_ref()
386 .context("session target is missing")?;
387 let backend = backend_locator(locator, session, &self.config)?;
388 let profile = self
389 .config
390 .profiles
391 .get(&session.last_profile)
392 .context("session profile is missing")?;
393 let bundle = session
394 .project_directory
395 .is_none()
396 .then(|| self.config.bundles.get(&session.bundle_id))
397 .flatten();
398 let workspace = if let Some(project_directory) = &session.project_directory {
399 (project_directory.to_string_lossy().into_owned(), Vec::new())
400 } else {
401 workspace_paths(
402 &backend,
403 bundle.context("session bundle is missing")?,
404 session_id,
405 session.container_workspace.as_deref(),
406 )?
407 };
408 let target_home = target_profile_home(&backend, session_id, profile);
409 let launch = project_memory_launch(session, bundle, &workspace, &target_home)?;
410 Ok(ProjectMemorySyncTarget {
411 canonical_root: canonical_memory_root(&launch.project_key),
412 })
413 }
414}
415
416pub(super) fn apply_jev_switch(launch: &mut WorkerLaunchConfig, enabled: bool) {
420 if enabled {
421 return;
422 }
423 for environment in [&mut launch.target_environment, &mut launch.environment] {
424 environment.remove("TYPESAFE_API_KEY");
425 environment.insert(
426 mj_core::jev::DISABLED_ENVIRONMENT.to_owned(),
427 "1".to_owned(),
428 );
429 }
430}
431
432pub(super) fn subagent_tools_enabled(
439 session: &mj_core::state::SessionRecord,
440 is_child: bool,
441) -> bool {
442 session.mjolnir_subagents.unwrap_or(false)
443 && !is_child
444 && matches!(
445 session.harness_kind,
446 mj_core::config::HarnessKind::Claude | mj_core::config::HarnessKind::Codex
447 )
448}
449
450pub(super) fn worker_workspace_for_recovery(
451 backend: &targets::TargetLocator,
452 directory: &Path,
453) -> Option<WorkerWorkspace> {
454 let target = match backend {
455 targets::TargetLocator::LocalBare { .. } => mj_core::state::ManagedWorktreeTarget::Local,
456 targets::TargetLocator::SshBare { ssh, .. } => mj_core::state::ManagedWorktreeTarget::Ssh {
457 destination: ssh.destination.clone(),
458 ssh_args: ssh.ssh_args.clone(),
459 },
460 targets::TargetLocator::LocalPodman { .. }
461 | targets::TargetLocator::LocalDocker { .. }
462 | targets::TargetLocator::AppleContainer { .. }
463 | targets::TargetLocator::AwsEc2 { .. }
464 | targets::TargetLocator::SshPodman { .. }
465 | targets::TargetLocator::SshDocker { .. } => return None,
466 };
467 Some(WorkerWorkspace {
468 target,
469 directory: directory.to_path_buf(),
470 })
471}
472
473pub(super) fn worker_launch_config(
474 session: &mj_core::state::SessionRecord,
475 profile: &mj_core::config::HarnessProfile,
476 bundle: Option<&ProjectBundle>,
477 backend: &targets::TargetLocator,
478 workspace_session_id: &str,
479 workspace_container: Option<&Path>,
480 target: &mj_core::state::TargetRuntimeSettings,
481) -> Result<(WorkerLaunchConfig, ProjectMemoryLaunchConfig, String)> {
482 let session_id = session.id.as_str();
483 let execution_policy = profile
484 .kind
485 .effective_execution_policy(target.execution_policy);
486 let target_profile_home = target_profile_home(backend, session_id, profile);
487 let workspace = if let Some(project_directory) = &session.project_directory {
488 (project_directory.to_string_lossy().into_owned(), Vec::new())
489 } else {
490 workspace_paths(
491 backend,
492 bundle.context("session bundle is missing")?,
493 workspace_session_id,
494 workspace_container,
495 )?
496 };
497 let mut additional_directories = workspace.1.iter().map(PathBuf::from).collect::<Vec<_>>();
498 additional_directories.extend(
499 session
500 .additional_mounts
501 .iter()
502 .map(|resource| resource.destination.clone()),
503 );
504 if profile.kind == mj_core::config::HarnessKind::Muse && !additional_directories.is_empty() {
505 bail!(
506 "{} ACP does not support multiple workspace roots; use a single-repository bundle",
507 profile.kind.display_name()
508 );
509 }
510 let (bridge_command, bridge_args) = bridge_launch(profile.kind, execution_policy);
511 let mut target_environment = target.environment.clone();
512 for name in [
522 "MJ_TURN_STALL_TIMEOUT_MS",
523 "MJ_TURN_TOOL_STALL_TIMEOUT_MS",
524 "RUST_LOG",
525 ] {
526 if let Ok(value) = std::env::var(name) {
527 target_environment.insert(name.to_owned(), value);
528 }
529 }
530 if let Some(key) = mj_core::activity::verdict::api_key() {
533 target_environment.insert("TYPESAFE_API_KEY".to_owned(), key);
534 }
535 target_environment.insert(
540 "MJ_INSTANCE".to_owned(),
541 mj_core::config::instance_identity(),
542 );
543 if let Some(build_cache) = &session.build_cache {
546 target_environment.insert(
547 "MBX_CACHE_DIR".into(),
548 build_cache.directory.to_string_lossy().into_owned(),
549 );
550 target_environment.insert(
553 "MBX_SHIMS_DIR".into(),
554 Path::new(&targets::worker_root(backend, session_id)?)
555 .join("mbx-shims")
556 .to_string_lossy()
557 .into_owned(),
558 );
559 if let Some(max_size) = &build_cache.max_size {
560 target_environment.insert("MBX_GC_MAX_TOTAL_SIZE".into(), max_size.clone());
561 }
562 target_environment.insert("MBX_SUMMARY".into(), "off".into());
565 target_environment.insert("MBX_SAVINGS".into(), "off".into());
566 }
567 let mut environment = target_environment.clone();
568 environment.extend(profile.environment.clone());
569 profile
570 .kind
571 .configure_home_environment(Path::new(&target_profile_home), &mut environment);
572 profile
573 .kind
574 .configure_execution_environment(execution_policy, &mut environment)?;
575 let mut project_memory =
576 project_memory_launch(session, bundle, &workspace, &target_profile_home)?;
577 project_memory.mcp_delivery = project_memory_mcp_delivery(profile.kind, backend);
578 project_memory.history_socket =
579 Some(Path::new(&targets::worker_root(backend, &session.id)?).join("control.sock"));
580 if profile.kind == mj_core::config::HarnessKind::Claude {
581 environment.insert(
582 "CLAUDE_CODE_PROJECT_DIR_NAME".into(),
583 project_memory_replica_slug(&project_memory.project_key, session_id),
584 );
585 }
586 apply_claude_setup_token(
587 &mut environment,
588 profile.kind,
589 &mj_core::credentials::claude_oauth_token_path(&session.last_profile),
590 );
591 let excluded_environment =
592 exclude_harness_environment(&session.last_profile, profile, &mut environment);
593 Ok((
594 WorkerLaunchConfig {
595 goal_resume_request: None,
596 target_environment,
597 seed_image_environment: backend.container_engine().is_some(),
598 run_mode: Default::default(),
599 session_id: session_id.to_string(),
600 subagent_tools: false,
601 handback_tool: false,
602 review_capture: false,
603 harness: profile.kind,
604 harness_home: PathBuf::from(&target_profile_home),
605 authentication_marker: profile
611 .authentication_marker()
612 .strip_prefix(&profile.home)
613 .ok()
614 .map(|name| name.to_string_lossy().into_owned()),
615 bridge_command: PathBuf::from(bridge_command),
616 bridge_args,
617 harness_runtime: harness_runtime_policy(backend),
618 environment,
619 excluded_environment,
620 cwd: PathBuf::from(&workspace.0),
621 additional_directories,
622 native_session_id: session.native_session_id.clone(),
623 project_memory: profile
624 .kind
625 .supports_injected_mcp()
626 .then(|| project_memory.clone()),
627 execution_policy,
628 },
629 project_memory,
630 target_profile_home,
631 ))
632}
633
634pub(super) fn exclude_harness_environment(
641 profile_id: &str,
642 profile: &mj_core::config::HarnessProfile,
643 environment: &mut std::collections::BTreeMap<String, String>,
644) -> Vec<String> {
645 static REPORTED: std::sync::Mutex<std::collections::BTreeSet<String>> =
646 std::sync::Mutex::new(std::collections::BTreeSet::new());
647 let before = environment.clone();
648 let excluded = profile.exclude_harness_environment(environment);
649 let removed = excluded
650 .iter()
651 .filter(|name| before.contains_key(*name))
652 .cloned()
653 .collect::<Vec<_>>();
654 if !removed.is_empty()
655 && REPORTED
656 .lock()
657 .map(|mut reported| reported.insert(profile_id.to_owned()))
658 .unwrap_or(true)
659 {
660 tracing::info!(
661 profile_id,
662 removed = removed.join(", "),
663 "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"
664 );
665 }
666 excluded
667}
668
669pub(super) fn harness_runtime_policy(backend: &targets::TargetLocator) -> HarnessRuntimePolicy {
670 match backend {
671 targets::TargetLocator::LocalBare { .. }
672 | targets::TargetLocator::AwsEc2 { .. }
673 | targets::TargetLocator::SshBare { .. } => HarnessRuntimePolicy::Managed,
674 _ => HarnessRuntimePolicy::Ambient,
675 }
676}