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