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");
65 if crate::controller::session_owns_profile_home(backend, session_id, profile) {
71 let started = Instant::now();
72 let result = stage_profile(profile, &profile_stage);
73 tracing::debug!(
74 session_id,
75 elapsed_ms = started.elapsed().as_millis(),
76 "profile staging completed"
77 );
78 result?;
79 stage_managed_skills(profile.kind, &profile_stage)?;
80 stage_codex_catalog(
81 &session.last_profile,
82 profile,
83 &profile_stage,
84 &fetch_catalog_over_https,
85 &SharedCatalogCache,
86 )?;
87 append_hel_target_environment(profile.kind, &profile_stage, backend)?;
88 apply_staged_execution_setting(profile.kind, launch.execution_policy, &profile_stage)?;
89 if profile.kind == mj_core::config::HarnessKind::Claude {
90 if launch.subagent_tools {
91 configure_claude_subagent_mcp(
92 &profile_stage,
93 worker_root,
94 mj_core::subagent::SubagentMcpRole::Parent,
95 )?;
96 } else if launch.handback_tool {
97 configure_claude_subagent_mcp(
98 &profile_stage,
99 worker_root,
100 mj_core::subagent::SubagentMcpRole::Child,
101 )?;
102 }
103 }
104 stage_memory_replica(
105 &project_memory,
106 Path::new(&target_profile_home),
107 &profile_stage,
108 )?;
109 if project_memory.mcp_delivery == ProjectMemoryMcpDelivery::HarnessProfile {
110 configure_kimi_project_memory_mcp(&profile_stage, worker_root, &project_memory)?;
111 }
112 } else {
113 seed_local_memory_replica(&project_memory)?;
114 }
115 let worker_binary = worker_binary_for(backend, executor)?;
116
117 install_worker_files(
118 executor,
119 backend,
120 session_id,
121 worker_root,
122 &target_profile_home,
123 &worker_binary,
124 &launch_path,
125 &ownership_path,
126 &profile_stage,
127 )?;
128 if session.build_cache.is_some()
131 && let Err(error) = self.install_build_cache_shim(session, backend, executor)
132 {
133 tracing::warn!(
134 session_id,
135 "installing the mbx build cache failed: {error:#}"
136 );
137 }
138 prepare_installed_managed_harness(executor, backend, worker_root, &launch)
139 }
140
141 fn install_build_cache_shim(
146 &self,
147 session: &mj_core::state::SessionRecord,
148 backend: &targets::TargetLocator,
149 executor: &impl CommandExecutor,
150 ) -> Result<()> {
151 let worker_root = targets::worker_root(backend, &session.id)?;
152 let binary =
156 crate::controller::mbx::binary_for(backend, executor).inspect_err(|error| {
157 executor.notify_notice(&format!(
158 "The Rust build cache is unavailable: {error:#}; this session builds without it."
159 ));
160 })?;
161 let configuration = crate::controller::mbx::host_configuration(backend, executor)?;
162 install_mbx_files(
163 executor,
164 backend,
165 &session.id,
166 &worker_root,
167 &binary,
168 configuration.as_deref(),
169 )
170 }
171
172 pub fn diagnose_worker(&self, session_id: &str) -> Option<String> {
176 self.diagnose_worker_controlled(session_id, &crate::targets::ProcessExecutor)
177 }
178
179 pub fn diagnose_worker_controlled(
180 &self,
181 session_id: &str,
182 executor: &impl CommandExecutor,
183 ) -> Option<String> {
184 let session = self.state.sessions.get(session_id)?;
185 let locator = session.target.as_ref()?;
186 let backend = match backend_locator(locator, session, &self.config) {
187 Ok(backend) => backend,
188 Err(error) => {
189 tracing::debug!(
190 session_id,
191 error = format!("{error:#}"),
192 "could not construct a worker diagnostic probe"
193 );
194 return None;
195 }
196 };
197 let worker_root = match targets::worker_root(&backend, session_id) {
198 Ok(root) => root,
199 Err(error) => {
200 tracing::debug!(
201 session_id,
202 error = format!("{error:#}"),
203 "could not derive the worker diagnostic root"
204 );
205 return None;
206 }
207 };
208 let binary_failure = worker_binary_probe_failure(executor, &backend, &worker_root);
209 let last_words = worker_last_words(executor, &backend, &worker_root);
210 match (binary_failure, last_words) {
211 (Some(binary_failure), Some(last_words)) => {
212 Some(format!("{binary_failure}; {last_words}"))
213 }
214 (Some(binary_failure), None) => Some(binary_failure),
215 (None, last_words) => last_words,
216 }
217 }
218
219 pub fn worker_recovery_plan(&self, session_id: &str) -> Result<WorkerRecoveryPlan> {
223 let (backend, worker_root) = self.worker_placement(session_id)?;
224 let launch = self.current_worker_launch_config(session_id, &backend)?;
225 let workspace = worker_workspace_for_recovery(&backend, &launch.cwd);
226 Ok(WorkerRecoveryPlan {
227 source_target: self.state.sessions[session_id]
228 .target
229 .clone()
230 .context("session target is missing")?,
231 target: targets::target_recovery_plan(&backend, session_id)?,
232 workspace,
233 liveness_probe: worker_liveness_command(&backend, &worker_root),
234 binary_refresh: worker_binary_refresh_plan(&backend, session_id)?,
235 launch_refresh: Some(worker_launch_refresh_plan(&backend, session_id, &launch)?),
236 restart: CommandPlan {
237 description: format!("restart Mjolnir worker for session {session_id}"),
238 commands: vec![
239 stop_worker_command(&backend, &worker_root),
240 start_worker_command(&backend, &worker_root),
241 ],
242 },
243 })
244 }
245
246 fn session_launch_config(
251 &self,
252 session_id: &str,
253 backend: &targets::TargetLocator,
254 ) -> Result<(WorkerLaunchConfig, ProjectMemoryLaunchConfig, String)> {
255 let session = self
256 .state
257 .sessions
258 .get(session_id)
259 .with_context(|| format!("unknown session {session_id}"))?;
260 session.validate_configuration(&self.config)?;
261 let profile = self
262 .config
263 .profiles
264 .get(&session.last_profile)
265 .context("session profile is missing")?;
266 let bundle = session
267 .project_directory
268 .is_none()
269 .then(|| self.config.bundles.get(&session.bundle_id))
270 .flatten();
271 let target = session.target_runtime_settings(&self.config)?;
272 let subagent = crate::database::load_subagent(session_id)?;
273 let (workspace_session_id, workspace_container) = match subagent.as_ref() {
276 Some(child) => {
277 let parent = self
278 .state
279 .sessions
280 .get(&child.parent_session_id)
281 .context("sub-agent parent session is missing")?;
282 (parent.id.clone(), parent.container_workspace.clone())
283 }
284 None => (session_id.to_owned(), session.container_workspace.clone()),
285 };
286 let (mut launch, project_memory, target_profile_home) = worker_launch_config(
287 session,
288 profile,
289 bundle,
290 backend,
291 &workspace_session_id,
292 workspace_container.as_deref(),
293 &target,
294 )?;
295 apply_jev_switch(&mut launch, self.config.jev.enabled);
296 launch.subagent_tools = subagent_tools_enabled(session, subagent.is_some());
297 launch.handback_tool = subagent.as_ref().is_some_and(|child| child.handback_tool);
299 if launch.subagent_tools
305 && profile.kind == mj_core::config::HarnessKind::Claude
306 && !crate::controller::session_owns_profile_home(backend, session_id, profile)
307 {
308 tracing::info!(
309 session_id,
310 "Mjolnir sub-agents need a harness home of their own; this Claude session runs \
311 out of the user's own home and keeps Claude's Agent and Task tools instead"
312 );
313 launch.subagent_tools = false;
314 }
315 launch.review_capture =
321 mj_core::review::settings::can_review(&self.config) && subagent.is_none();
322 if let Some(subagent) = &subagent {
323 let parent = self
324 .state
325 .sessions
326 .get(&subagent.parent_session_id)
327 .context("sub-agent parent session is missing")?;
328 let parent_profile = self
329 .config
330 .profiles
331 .get(&parent.last_profile)
332 .context("sub-agent parent profile is missing")?;
333 let parent_target = parent.target_runtime_settings(&self.config)?;
334 let parent_locator = parent
335 .target
336 .as_ref()
337 .context("sub-agent parent has no live target")?;
338 let parent_backend = backend_locator(parent_locator, parent, &self.config)?;
339 let parent_bundle = parent
340 .project_directory
341 .is_none()
342 .then(|| self.config.bundles.get(&parent.bundle_id))
343 .flatten();
344 let (parent_launch, _, _) = worker_launch_config(
345 parent,
346 parent_profile,
347 parent_bundle,
348 &parent_backend,
349 &parent.id,
350 parent.container_workspace.as_deref(),
351 &parent_target,
352 )?;
353 launch.cwd = if subagent.working_directory.as_os_str().is_empty() {
354 parent_launch.cwd
355 } else {
356 parent_launch.cwd.join(&subagent.working_directory)
357 };
358 launch.additional_directories = parent_launch.additional_directories;
359 }
360 Ok((launch, project_memory, target_profile_home))
361 }
362
363 pub(in crate::controller) fn current_worker_launch_config(
364 &self,
365 session_id: &str,
366 backend: &targets::TargetLocator,
367 ) -> Result<WorkerLaunchConfig> {
368 let session = self
369 .state
370 .sessions
371 .get(session_id)
372 .with_context(|| format!("unknown session {session_id}"))?;
373 let (mut launch, _, _) = self.session_launch_config(session_id, backend)?;
374 if crate::database::load_move_operation(session_id)?.is_some_and(|operation| {
375 operation.source_checkpoint_only
376 && operation.destination_target.is_none()
377 && matches!(
378 operation.phase,
379 mj_core::state::MovePhase::Preparing
380 | mj_core::state::MovePhase::ClosingSource
381 | mj_core::state::MovePhase::Failed
382 | mj_core::state::MovePhase::Cancelled
383 )
384 && session.last_profile == operation.source_profile_id
385 && session.target == operation.source_target
386 && matches!(
387 session.state,
388 mj_core::state::SessionState::Running
389 | mj_core::state::SessionState::Disconnected
390 | mj_core::state::SessionState::Closing
391 )
392 }) {
393 launch.run_mode = mj_core::worker_launch::WorkerRunMode::CheckpointOnly;
394 }
395 Ok(launch)
396 }
397
398 pub fn project_memory_sync_target(&self, session_id: &str) -> Result<ProjectMemorySyncTarget> {
399 let session = self
400 .state
401 .sessions
402 .get(session_id)
403 .with_context(|| format!("unknown session {session_id}"))?;
404 session.validate_configuration(&self.config)?;
405 let locator = session
406 .target
407 .as_ref()
408 .context("session target is missing")?;
409 let backend = backend_locator(locator, session, &self.config)?;
410 let profile = self
411 .config
412 .profiles
413 .get(&session.last_profile)
414 .context("session profile is missing")?;
415 let bundle = session
416 .project_directory
417 .is_none()
418 .then(|| self.config.bundles.get(&session.bundle_id))
419 .flatten();
420 let workspace = if let Some(project_directory) = &session.project_directory {
421 (project_directory.to_string_lossy().into_owned(), Vec::new())
422 } else {
423 workspace_paths(
424 &backend,
425 bundle.context("session bundle is missing")?,
426 session_id,
427 session.container_workspace.as_deref(),
428 )?
429 };
430 let target_home = target_profile_home(&backend, session_id, profile);
431 let launch = project_memory_launch(session, bundle, &workspace, &target_home)?;
432 Ok(ProjectMemorySyncTarget {
433 canonical_root: canonical_memory_root(&launch.project_key),
434 })
435 }
436}
437
438pub(super) fn apply_jev_switch(launch: &mut WorkerLaunchConfig, enabled: bool) {
442 if enabled {
443 return;
444 }
445 for environment in [&mut launch.target_environment, &mut launch.environment] {
446 environment.remove("TYPESAFE_API_KEY");
447 environment.insert(
448 mj_core::jev::DISABLED_ENVIRONMENT.to_owned(),
449 "1".to_owned(),
450 );
451 }
452}
453
454pub(super) fn subagent_tools_enabled(
461 session: &mj_core::state::SessionRecord,
462 is_child: bool,
463) -> bool {
464 session.mjolnir_subagents.unwrap_or(false)
465 && !is_child
466 && matches!(
467 session.harness_kind,
468 mj_core::config::HarnessKind::Claude | mj_core::config::HarnessKind::Codex
469 )
470}
471
472pub(super) fn worker_workspace_for_recovery(
473 backend: &targets::TargetLocator,
474 directory: &Path,
475) -> Option<WorkerWorkspace> {
476 let target = match backend {
477 targets::TargetLocator::LocalBare { .. } => mj_core::state::ManagedWorktreeTarget::Local,
478 targets::TargetLocator::SshBare { ssh, .. } => mj_core::state::ManagedWorktreeTarget::Ssh {
479 destination: ssh.destination.clone(),
480 ssh_args: ssh.ssh_args.clone(),
481 },
482 targets::TargetLocator::LocalPodman { .. }
483 | targets::TargetLocator::LocalDocker { .. }
484 | targets::TargetLocator::AppleContainer { .. }
485 | targets::TargetLocator::AwsEc2 { .. }
486 | targets::TargetLocator::SshPodman { .. }
487 | targets::TargetLocator::SshDocker { .. } => return None,
488 };
489 Some(WorkerWorkspace {
490 target,
491 directory: directory.to_path_buf(),
492 })
493}
494
495pub(super) fn worker_launch_config(
496 session: &mj_core::state::SessionRecord,
497 profile: &mj_core::config::HarnessProfile,
498 bundle: Option<&ProjectBundle>,
499 backend: &targets::TargetLocator,
500 workspace_session_id: &str,
501 workspace_container: Option<&Path>,
502 target: &mj_core::state::TargetRuntimeSettings,
503) -> Result<(WorkerLaunchConfig, ProjectMemoryLaunchConfig, String)> {
504 let session_id = session.id.as_str();
505 let execution_policy = profile
506 .kind
507 .effective_execution_policy(target.execution_policy);
508 let target_profile_home = target_profile_home(backend, session_id, profile);
509 let workspace = if let Some(project_directory) = &session.project_directory {
510 (project_directory.to_string_lossy().into_owned(), Vec::new())
511 } else {
512 workspace_paths(
513 backend,
514 bundle.context("session bundle is missing")?,
515 workspace_session_id,
516 workspace_container,
517 )?
518 };
519 let mut additional_directories = workspace.1.iter().map(PathBuf::from).collect::<Vec<_>>();
520 additional_directories.extend(
521 session
522 .additional_mounts
523 .iter()
524 .map(|resource| resource.destination.clone()),
525 );
526 if profile.kind == mj_core::config::HarnessKind::Muse && !additional_directories.is_empty() {
527 bail!(
528 "{} ACP does not support multiple workspace roots; use a single-repository bundle",
529 profile.kind.display_name()
530 );
531 }
532 let (bridge_command, bridge_args) = bridge_launch(profile.kind, execution_policy);
533 let mut target_environment = target.environment.clone();
534 for name in [
544 "MJ_TURN_STALL_TIMEOUT_MS",
545 "MJ_TURN_TOOL_STALL_TIMEOUT_MS",
546 "RUST_LOG",
547 ] {
548 if let Ok(value) = std::env::var(name) {
549 target_environment.insert(name.to_owned(), value);
550 }
551 }
552 if let Some(key) = mj_core::activity::verdict::api_key() {
555 target_environment.insert("TYPESAFE_API_KEY".to_owned(), key);
556 }
557 target_environment.insert(
562 "MJ_INSTANCE".to_owned(),
563 mj_core::config::instance_identity(),
564 );
565 if let Some(build_cache) = &session.build_cache {
568 target_environment.insert(
569 "MBX_CACHE_DIR".into(),
570 build_cache.directory.to_string_lossy().into_owned(),
571 );
572 target_environment.insert(
575 "MBX_SHIMS_DIR".into(),
576 Path::new(&targets::worker_root(backend, session_id)?)
577 .join("mbx-shims")
578 .to_string_lossy()
579 .into_owned(),
580 );
581 if let Some(max_size) = &build_cache.max_size {
582 target_environment.insert("MBX_GC_MAX_TOTAL_SIZE".into(), max_size.clone());
583 }
584 target_environment.insert("MBX_SUMMARY".into(), "off".into());
587 target_environment.insert("MBX_SAVINGS".into(), "off".into());
588 }
589 let mut environment = target_environment.clone();
590 environment.extend(profile.environment.clone());
591 profile.kind.configure_home_environment(
592 Path::new(&target_profile_home),
593 backend.harness_host(),
594 &mut environment,
595 );
596 profile
597 .kind
598 .configure_execution_environment(execution_policy, &mut environment)?;
599 let mut project_memory =
600 project_memory_launch(session, bundle, &workspace, &target_profile_home)?;
601 project_memory.mcp_delivery = project_memory_mcp_delivery(profile.kind, backend);
602 project_memory.history_socket =
603 Some(Path::new(&targets::worker_root(backend, &session.id)?).join("control.sock"));
604 if profile.kind == mj_core::config::HarnessKind::Claude {
605 environment.insert(
606 "CLAUDE_CODE_PROJECT_DIR_NAME".into(),
607 project_memory_replica_slug(&project_memory.project_key, session_id),
608 );
609 }
610 apply_claude_setup_token(
611 &mut environment,
612 profile.kind,
613 &mj_core::credentials::claude_oauth_token_path(&session.last_profile),
614 );
615 Ok((
616 WorkerLaunchConfig {
617 goal_resume_request: None,
618 target_environment,
619 seed_image_environment: backend.container_engine().is_some(),
620 run_mode: Default::default(),
621 session_id: session_id.to_string(),
622 subagent_tools: false,
623 handback_tool: false,
624 review_capture: false,
625 harness: profile.kind,
626 harness_home: PathBuf::from(&target_profile_home),
627 authentication_marker: profile
630 .authentication_marker()
631 .file_name()
632 .map(|name| name.to_string_lossy().into_owned()),
633 bridge_command: PathBuf::from(bridge_command),
634 bridge_args,
635 harness_runtime: harness_runtime_policy(backend),
636 environment,
637 cwd: PathBuf::from(&workspace.0),
638 additional_directories,
639 native_session_id: session.native_session_id.clone(),
640 project_memory: profile
641 .kind
642 .supports_injected_mcp()
643 .then(|| project_memory.clone()),
644 execution_policy,
645 },
646 project_memory,
647 target_profile_home,
648 ))
649}
650
651pub(super) fn harness_runtime_policy(backend: &targets::TargetLocator) -> HarnessRuntimePolicy {
652 match backend {
653 targets::TargetLocator::LocalBare { .. }
654 | targets::TargetLocator::AwsEc2 { .. }
655 | targets::TargetLocator::SshBare { .. } => HarnessRuntimePolicy::Managed,
656 _ => HarnessRuntimePolicy::Ambient,
657 }
658}