1use std::collections::{BTreeMap, HashMap};
4use std::path::{Path, PathBuf};
5use std::process::{Command, Stdio};
6use std::sync::{Arc, Mutex, OnceLock};
7use std::time::{Duration, Instant};
8
9use anyhow::{Context, Result, bail, ensure};
10
11use mj_core::config::{TargetTemplate, atomic_write, data_dir};
12use mj_core::state::{SessionState, State, TargetLocator};
13
14use crate::targets::{
15 self, CancellableProcessExecutor, CommandExecutor, CommandOutput, CommandSpec, ProvisionStage,
16 ProvisionStageGuard,
17};
18
19use super::backend::{
20 ContainerOverrides, TargetCheck, backend_bundle, backend_locator, backend_target,
21 configure_github_token_environment, controller_github_token, locator_after_provision,
22 preflight_target, use_github_https_urls,
23};
24use super::git_cache;
25use super::readiness::{connect_started_worker, wait_for_native_session_in_stage};
26use super::worker_binary::{bridge_readiness_stage, start_worker, worker_probe_diagnosis};
27use super::{Controller, execute_checked, now};
28
29const INHERITED_GIT_SETTINGS: &[&str] = &[
30 "diff.algorithm",
31 "fetch.prune",
32 "fetch.prunetags",
33 "init.defaultbranch",
34 "merge.conflictstyle",
35 "pull.ff",
36 "pull.rebase",
37 "push.autosetupremote",
38 "push.default",
39 "rebase.autostash",
40 "rerere.autoupdate",
41 "rerere.enabled",
42 "user.email",
43 "user.name",
44];
45
46const CONTAINER_START_ADMISSION: usize = 2;
57
58pub(super) fn container_start_gate(
64 locator: &targets::TargetLocator,
65) -> Option<Arc<tokio::sync::Semaphore>> {
66 static GATES: OnceLock<Mutex<HashMap<String, Arc<tokio::sync::Semaphore>>>> = OnceLock::new();
67 let container = match locator {
68 targets::TargetLocator::LocalPodman { container_id, .. }
69 | targets::TargetLocator::LocalDocker { container_id, .. }
70 | targets::TargetLocator::AppleContainer { container_id, .. }
71 | targets::TargetLocator::SshPodman { container_id, .. }
72 | targets::TargetLocator::SshDocker { container_id, .. } => container_id.clone(),
73 targets::TargetLocator::LocalBare { .. }
74 | targets::TargetLocator::SshBare { .. }
75 | targets::TargetLocator::AwsEc2 { .. } => return None,
76 };
77 let gates = GATES.get_or_init(|| Mutex::new(HashMap::new()));
78 let mut gates = gates
79 .lock()
80 .unwrap_or_else(std::sync::PoisonError::into_inner);
81 Some(Arc::clone(gates.entry(container).or_insert_with(|| {
82 Arc::new(tokio::sync::Semaphore::new(CONTAINER_START_ADMISSION))
83 })))
84}
85
86fn subagent_start_is_retryable(error: &anyhow::Error) -> bool {
94 if mj_core::refusal::Refusal::of(error).is_some() {
95 return false;
96 }
97 if format!("{error:#}").contains("operation cancelled") {
98 return false;
99 }
100 error
101 .downcast_ref::<super::readiness::WorkerStartupFailure>()
102 .is_some_and(|failure| !failure.reached_socket)
103}
104
105#[derive(Debug, Clone, Copy, PartialEq, Eq)]
106pub(super) enum ProvisioningFailureDisposition {
107 Discard,
109 Preserve,
111}
112
113impl Controller {
114 pub async fn provision_session_controlled_with_commit(
115 &mut self,
116 session_id: &str,
117 executor: &(impl CommandExecutor + Sync),
118 grant_commit: impl FnOnce() -> Result<()>,
119 ) -> Result<()> {
120 let github_token = controller_github_token();
121 let repositories = self
122 .provision_session_target_with_failure_disposition(
123 session_id,
124 executor,
125 github_token.as_deref(),
126 ProvisioningFailureDisposition::Discard,
127 )
128 .await?;
129 let setup = execute_concurrent_lanes(
130 || execute_repository_setup(&repositories, executor),
131 || self.install_worker_payload(session_id, executor),
132 );
133 let result = match setup {
134 Ok(((), (backend, worker_root))) => {
135 self.connect_and_start_worker(session_id, executor, &backend, &worker_root, true)
136 .await
137 }
138 Err(error) => Err(error),
139 };
140 match result {
141 Ok(native_session_id) => {
142 if let Err(error) = grant_commit() {
143 return Err(self.rollback_failed_new_session(session_id, error, executor)?);
144 }
145 self.mark_worker_connected(session_id, native_session_id)
146 }
147 Err(error) => Err(self.rollback_failed_new_session(session_id, error, executor)?),
148 }
149 }
150
151 pub async fn provision_subagent_session_controlled(
154 &mut self,
155 session_id: &str,
156 executor: &(impl CommandExecutor + Sync),
157 ) -> Result<()> {
158 let mut attempts: Vec<String> = Vec::new();
165 let (result, placement) = loop {
166 let attempt = self.attempt_subagent_start(session_id, executor).await;
167 let (result, placement) = attempt;
168 let Err(error) = &result else {
169 break (result, placement);
170 };
171 if attempts.len() == 1 || !subagent_start_is_retryable(error) {
172 if !attempts.is_empty() {
173 let combined = attempts
174 .iter()
175 .enumerate()
176 .map(|(index, attempt)| format!("attempt {}: {attempt}", index + 1))
177 .chain(std::iter::once(format!(
178 "attempt {}: {error:#}",
179 attempts.len() + 1
180 )))
181 .collect::<Vec<_>>()
182 .join("; ");
183 break (Err(anyhow::anyhow!("{combined}")), placement);
184 }
185 break (result, placement);
186 }
187 tracing::warn!(
188 session_id,
189 error = format!("{error:#}"),
190 "sub-agent worker never started; retrying once"
191 );
192 attempts.push(format!("{error:#}"));
193 if let Some((backend, worker_root)) = &placement
196 && let Err(stop_error) =
197 super::worker_binary::stop_worker(executor, backend, worker_root)
198 {
199 tracing::debug!(
200 session_id,
201 error = format!("{stop_error:#}"),
202 "could not stop the worker of a retried sub-agent start"
203 );
204 }
205 };
206 match result {
207 Ok(native_session_id) => self.mark_worker_connected(session_id, native_session_id),
208 Err(error) => {
209 if let Some((backend, worker_root)) = &placement
211 && let Err(stop_error) =
212 super::worker_binary::stop_worker(executor, backend, worker_root)
213 {
214 tracing::warn!(
215 session_id,
216 error = format!("{stop_error:#}"),
217 "failed sub-agent worker could not be stopped cleanly"
218 );
219 }
220 tracing::warn!(
221 session_id,
222 error = format!("{error:#}"),
223 "sub-agent startup failed"
224 );
225 let record = self
226 .state
227 .sessions
228 .get_mut(session_id)
229 .context("failed sub-agent session disappeared")?;
230 record.state = SessionState::Error;
231 record.updated_at = super::now();
232 let error = match &placement {
235 Some((backend, _)) => {
236 super::subagent_park::explain_process_exhaustion(error, backend, session_id)
237 }
238 None => error,
239 };
240 record.last_error = Some(format!("sub-agent startup failed: {error:#}"));
241 crate::database::save_lifecycle_session(record)?;
242 Err(error)
243 }
244 }
245 }
246
247 async fn attempt_subagent_start(
251 &mut self,
252 session_id: &str,
253 executor: &(impl CommandExecutor + Sync),
254 ) -> (
255 Result<Option<String>>,
256 Option<(targets::TargetLocator, String)>,
257 ) {
258 match self.worker_placement(session_id) {
261 Ok((backend, worker_root)) => {
262 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
263 let prepared = self
264 .prepare_worker_files(session_id, &backend, &worker_root, syncing)
265 .and_then(|()| self.prepare_subagent_report_dir(session_id, &backend, syncing));
266 let result = match prepared {
267 Ok(()) => {
268 let gate = container_start_gate(&backend);
271 let _admitted = match &gate {
272 Some(gate) => gate.acquire().await.ok(),
273 None => None,
274 };
275 self.connect_and_start_worker(
276 session_id,
277 executor,
278 &backend,
279 &worker_root,
280 false,
281 )
282 .await
283 }
284 Err(error) => Err(error),
285 };
286 (result, Some((backend, worker_root)))
287 }
288 Err(error) => (Err(error), None),
289 }
290 }
291
292 fn rollback_failed_new_session(
293 &mut self,
294 session_id: &str,
295 error: anyhow::Error,
296 executor: &impl CommandExecutor,
297 ) -> Result<anyhow::Error> {
298 self.rollback_failed_new_session_with(
299 session_id,
300 error,
301 executor,
302 &CancellableProcessExecutor::with_timeout(Duration::from_secs(15)),
305 )
306 }
307
308 fn rollback_failed_new_session_with(
309 &mut self,
310 session_id: &str,
311 error: anyhow::Error,
312 executor: &impl CommandExecutor,
313 target_cleanup_executor: &impl CommandExecutor,
314 ) -> Result<anyhow::Error> {
315 let session = self
316 .state
317 .sessions
318 .get(session_id)
319 .with_context(|| format!("unknown session {session_id}"))?
320 .clone();
321 let original = note_new_session_launch_failure(session_id, &error);
328 apply_failed_new_session_launch(&mut self.state, session_id, &original);
329 self.persist_session_transition_or_restore(
330 session_id,
331 &session,
332 "record the failed launch before removing its target",
333 )
334 .map_err(|persist_error| {
335 persist_error.context(format!(
336 "{original}; its target was left in place because the failure could not be recorded"
337 ))
338 })?;
339 let target_cleanup = match session.target.as_ref() {
340 Some(locator) => (|| -> Result<()> {
341 let backend = backend_locator(locator, &session, &self.config)?;
342 targets::close_plan(&backend, session_id)?
343 .execute(target_cleanup_executor)
344 .map(|_| ())
345 })(),
346 None => Ok(()),
347 };
348 let worktree_cleanup =
349 self.cleanup_new_session_worktree_after_failure(session_id, executor);
350 let cleanup_error = [target_cleanup, worktree_cleanup]
351 .into_iter()
352 .filter_map(Result::err)
353 .map(|error| format!("{error:#}"))
354 .collect::<Vec<_>>()
355 .join("; ");
356 if !cleanup_error.is_empty() {
357 tracing::warn!(
358 session_id,
359 error = %cleanup_error,
360 "new-session rollback cleanup reported failures"
361 );
362 }
363 let failure = apply_failed_new_session_rollback(
364 &mut self.state,
365 session_id,
366 &original,
367 (!cleanup_error.is_empty()).then_some(cleanup_error),
368 );
369 self.persist_session_state(session_id)?;
370 Ok(failure)
371 }
372
373 pub(super) async fn provision_session_with_failure_disposition(
374 &mut self,
375 session_id: &str,
376 executor: &(impl CommandExecutor + Sync),
377 github_token: Option<&str>,
378 failure_disposition: ProvisioningFailureDisposition,
379 ) -> Result<()> {
380 let repositories = self
381 .provision_session_target_with_failure_disposition(
382 session_id,
383 executor,
384 github_token,
385 failure_disposition,
386 )
387 .await?;
388 match execute_repository_setup(&repositories, executor) {
389 Ok(()) => Ok(()),
390 Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
391 Err(self.rollback_failed_new_session(session_id, error, executor)?)
392 }
393 Err(error) => Err(error),
394 }
395 }
396
397 async fn provision_session_target_with_failure_disposition(
398 &mut self,
399 session_id: &str,
400 executor: &(impl CommandExecutor + Sync),
401 github_token: Option<&str>,
402 failure_disposition: ProvisioningFailureDisposition,
403 ) -> Result<targets::CommandPlan> {
404 let session = self
405 .state
406 .sessions
407 .get(session_id)
408 .with_context(|| format!("unknown session {session_id}"))?
409 .clone();
410 if session.state != SessionState::Provisioning {
411 bail!("session {session_id} is not provisioning");
412 }
413 let preparation = (|| {
414 let selected = self
415 .config
416 .targets
417 .get(&session.target_template_id)
418 .context("target template disappeared before provisioning")?;
419 let runtime = mj_core::state::TargetRuntimeSettings::from(selected);
420 if let Some(recorded) = &session.target_runtime {
421 ensure!(
422 recorded == &runtime,
423 "target access settings changed before provisioning; retry with the selected target"
424 );
425 } else {
426 self.state
427 .sessions
428 .get_mut(session_id)
429 .unwrap()
430 .target_runtime = Some(runtime);
431 self.persist_session_state(session_id)?;
432 }
433 let template = self
434 .config
435 .targets
436 .get(&session.target_template_id)
437 .context("target template disappeared during provisioning")?;
438 let profile = self
439 .config
440 .profiles
441 .get(&session.last_profile)
442 .context("harness profile disappeared during provisioning")?;
443 super::worker_binary::preflight_worker_binary(template)?;
444 super::worker_binary::preflight_harness(template, profile, executor)?;
445 self.prepare_managed_raw_worktree(session_id, executor)
446 })();
447 let created_worktree = match preparation {
448 Ok(created) => created,
449 Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
450 return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
451 }
452 Err(error) => return Err(error),
453 };
454 let session = self
455 .state
456 .sessions
457 .get(session_id)
458 .expect("session retained after managed worktree preparation")
459 .clone();
460 let result = (|| {
463 let template = self
464 .config
465 .targets
466 .get(&session.target_template_id)
467 .context("target template disappeared during provisioning")?;
468 if matches!(template, TargetTemplate::AwsEc2 { .. }) {
469 for resource in &session.additional_mounts {
470 ensure!(
471 resource.source.is_dir(),
472 "attached resource source is not a directory: {}",
473 resource.source.display()
474 );
475 }
476 }
477 let mut target = backend_target(
478 template,
479 session.resource_allocation.as_ref(),
480 ContainerOverrides::for_session(&session),
481 )?;
482 let mut runtime_mounts = if matches!(target, targets::TargetTemplate::AwsEc2(_)) {
483 Vec::new()
484 } else {
485 session.additional_mounts.clone()
486 };
487 for notice in enforce_overlay_capable_mounts(&target, &mut runtime_mounts, executor) {
492 executor.notify_notice(¬ice);
493 }
494 let image_user = podman_image_user(&target, executor);
498 let mut bundle = if session.project_directory.is_some() {
499 None
500 } else if failure_disposition == ProvisioningFailureDisposition::Preserve {
501 Some(super::network_git::checkpoint_bundle(&session)?)
502 } else {
503 Some(backend_bundle(
504 self.config
505 .bundles
506 .get(&session.bundle_id)
507 .context("session bundle is missing")?,
508 executor,
509 )?)
510 };
511 let container_github_token =
512 github_token.filter(|_| configure_github_token_environment(&mut target));
513 if container_github_token.is_some()
514 && let Some(bundle) = bundle.as_mut()
515 {
516 use_github_https_urls(bundle);
517 }
518 preflight_target(template, executor, TargetCheck::Launch)?;
519 let prepared_cache = bundle.as_mut().and_then(|bundle| {
520 git_cache::prepare(
521 &target,
522 session_id,
523 bundle,
524 &mut runtime_mounts,
525 container_github_token,
526 executor,
527 )
528 });
529 let build_cache = super::mbx::prepare(
532 &target,
533 &self.config.build_cache,
534 &session,
535 bundle.as_ref(),
536 prepared_cache.as_ref(),
537 &mut runtime_mounts,
538 executor,
539 );
540 let provision = if let Some(project_directory) = &session.project_directory {
541 targets::provision_bare_project_plan(
542 &target,
543 session_id,
544 &project_directory.to_string_lossy(),
545 )
546 } else {
547 bundle
548 .as_ref()
549 .context("project bundle disappeared during provisioning")
550 .and_then(|bundle| {
551 targets::provision_plan(
552 &target,
553 session_id,
554 bundle,
555 &runtime_mounts,
556 image_user,
557 session.container_workspace.as_deref(),
558 )
559 })
560 };
561 let mut provision = match provision {
562 Ok(provision) => provision,
563 Err(error) => {
564 if let Some(cache) = &prepared_cache {
565 let _ = cache.cleanup(executor);
566 }
567 return Err(error);
568 }
569 };
570 if let Some(token) = container_github_token
571 && let Err(error) =
572 provision.provide_target_environment_secret(&target, "GH_TOKEN", token)
573 {
574 if let Some(cache) = &prepared_cache {
575 let _ = cache.cleanup(executor);
576 }
577 return Err(error);
578 }
579
580 let started = Instant::now();
581 let result =
582 provision_target_creation(&provision, &target, session_id, executor, |outputs| {
583 locator_after_provision(
584 template,
585 &target,
586 session_id,
587 outputs.first(),
588 executor,
589 )
590 })
591 .map(|(locator, remainder)| (locator, remainder, bundle, build_cache));
592 if result.is_err()
593 && let Some(cache) = &prepared_cache
594 {
595 if let Some(locator) = provisioned_locator(&target, session_id, None) {
596 let _ = targets::close_plan(&locator, session_id)
597 .and_then(|plan| plan.execute(executor).map(|_| ()));
598 } else {
599 let _ = cache.cleanup(executor);
600 }
601 }
602 tracing::debug!(
603 session_id,
604 elapsed_ms = started.elapsed().as_millis(),
605 "provisioning plan execution completed"
606 );
607 result
608 })();
609 let result = match result {
610 Err(error)
611 if created_worktree
612 && failure_disposition == ProvisioningFailureDisposition::Discard =>
613 {
614 return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
615 }
616 Err(error) if failure_disposition == ProvisioningFailureDisposition::Preserve => {
617 Err(error)
618 }
619 Err(error) => {
620 let detail = note_new_session_launch_failure(session_id, &error);
624 {
625 let record = self.state.sessions.get_mut(session_id).unwrap();
626 record.state = SessionState::Error;
627 record.target = None;
628 record.updated_at = super::now();
629 record.last_error = Some(format!("session provisioning failed: {detail}"));
630 }
631 return match self.persist_session_state(session_id) {
632 Ok(()) => Err(error),
633 Err(persistence_error) => Err(error.context(format!(
634 "persist removal of failed provisioning session {session_id}: {persistence_error:#}"
635 ))),
636 };
637 }
638 Ok((locator, remainder, bundle, build_cache)) => {
639 apply_new_session_provisioning_result(&mut self.state, session_id, Ok(locator))?;
640 self.state
641 .sessions
642 .get_mut(session_id)
643 .expect("session retained after provisioning")
644 .build_cache = build_cache;
645 let session = &self.state.sessions[session_id];
646 let backend = backend_locator(
647 session
648 .target
649 .as_ref()
650 .context("provisioned target disappeared")?,
651 session,
652 &self.config,
653 )?;
654 if matches!(backend, targets::TargetLocator::AwsEc2 { .. }) {
655 targets::provision_on_locator_plan(
656 &backend,
657 session_id,
658 bundle
659 .as_ref()
660 .context("AWS provisioning requires a project bundle")?,
661 )
662 } else {
663 Ok(remainder)
664 }
665 }
666 };
667 let result = match result {
668 Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
669 return Err(self.rollback_failed_new_session(session_id, error, executor)?);
670 }
671 result => result,
672 };
673 if result.is_ok()
674 && let Some(session) = self.state.sessions.get(session_id)
675 && let Some(directory) = session
676 .managed_worktree
677 .as_ref()
678 .map(|worktree| worktree.source_project_directory.clone())
679 .or_else(|| session.project_directory.clone())
680 && let Some(template) = self.config.targets.get(&session.target_template_id)
681 {
682 let host = match template {
683 TargetTemplate::LocalBare => Some("local"),
684 TargetTemplate::SshBare { ssh, .. } => Some(ssh.host.as_str()),
685 _ => None,
686 };
687 if let Some(host) = host {
688 self.state.remember_project_directory(host, &directory);
689 crate::database::remember_project_directory(host, &directory)?;
690 }
691 }
692 self.persist_session_state(session_id)?;
693 result
694 }
695
696 pub fn mark_worker_connected(
697 &mut self,
698 session_id: &str,
699 native_session_id: Option<String>,
700 ) -> Result<()> {
701 let session = self
702 .state
703 .sessions
704 .get(session_id)
705 .with_context(|| format!("unknown session {session_id}"))?;
706 if session.target.is_none() {
707 bail!("session {session_id} has no provisioned target");
708 }
709 let updated_at = now();
710 crate::database::mark_session_worker_connected(
711 session_id,
712 native_session_id.as_deref(),
713 &updated_at,
714 )?;
715 let session = self
716 .state
717 .sessions
718 .get_mut(session_id)
719 .expect("session disappeared after its worker connection was saved");
720 session.state = SessionState::Running;
721 if native_session_id.is_some() {
722 session.native_session_id = native_session_id;
723 }
724 session.updated_at = updated_at;
725 session.last_error = None;
726 Ok(())
727 }
728
729 fn install_worker_payload(
730 &self,
731 session_id: &str,
732 executor: &impl CommandExecutor,
733 ) -> Result<(targets::TargetLocator, String)> {
734 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
736 let (backend, worker_root) = self.worker_placement(session_id)?;
737 self.prepare_worker_files(session_id, &backend, &worker_root, syncing)?;
738 install_attached_resources(&self.state, session_id, &backend, &worker_root, syncing)?;
739 Ok((backend, worker_root))
740 }
741
742 async fn connect_and_start_worker(
743 &self,
744 session_id: &str,
745 executor: &impl CommandExecutor,
746 backend: &targets::TargetLocator,
747 worker_root: &str,
748 initialize_workspace: bool,
749 ) -> Result<Option<String>> {
750 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
751 if initialize_workspace {
752 install_inherited_git_settings(executor, backend, session_id)?;
753 self.initialize_network_workspaces(session_id, backend, syncing)?;
754 }
755 let session = self
756 .state
757 .sessions
758 .get(session_id)
759 .with_context(|| format!("unknown session {session_id}"))?;
760 let profile = self
761 .config
762 .profiles
763 .get(&session.last_profile)
764 .with_context(|| format!("unknown profile {}", session.last_profile))?;
765 let readiness_stage = bridge_readiness_stage(profile);
766 let reconnect = &targets::reconnect_plan(backend, session_id)?.commands[0];
767 let readiness = async {
768 let mut relay = {
769 let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
770 start_worker(executor, backend, worker_root)?;
771 connect_started_worker(reconnect, session_id, executor, backend, worker_root)
772 .await?
773 };
774 let native_session_id =
775 wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
776 Ok(Some(native_session_id))
777 }
778 .await;
779 match readiness {
780 Ok(native_session_id) => Ok(native_session_id),
781 Err(error) => Err(worker_probe_diagnosis(
782 executor,
783 backend,
784 worker_root,
785 error,
786 )),
787 }
788 }
789}
790
791const MAX_LAUNCH_DIAGNOSTIC_BYTES: usize = 64 * 1024;
792
793const RETAINED_LAUNCH_DIAGNOSTICS: usize = 20;
794
795pub(super) fn note_new_session_launch_failure(session_id: &str, error: &anyhow::Error) -> String {
801 note_new_session_launch_failure_in(&data_dir().join("diagnostics"), session_id, error)
802}
803
804fn note_new_session_launch_failure_in(
805 directory: &Path,
806 session_id: &str,
807 error: &anyhow::Error,
808) -> String {
809 let original = format!("{error:#}");
810 tracing::warn!(session_id, error = %original, "session launch failed");
811 match persist_launch_failure_to(directory, session_id, &original) {
812 Ok(path) => format!("{original}; full diagnostic saved to {}", path.display()),
813 Err(save_error) => {
814 format!("{original}; saving the local diagnostic failed: {save_error:#}")
815 }
816 }
817}
818
819fn persist_launch_failure_to(directory: &Path, session_id: &str, detail: &str) -> Result<PathBuf> {
820 mj_core::config::validate_id("session", session_id)?;
821 std::fs::create_dir_all(directory).with_context(|| {
822 format!(
823 "create launch diagnostics directory {}",
824 directory.display()
825 )
826 })?;
827 #[cfg(unix)]
828 {
829 use std::os::unix::fs::PermissionsExt;
830 std::fs::set_permissions(directory, std::fs::Permissions::from_mode(0o700))?;
831 }
832 let path = directory.join(format!("{session_id}-launch-error.txt"));
833 let detail = bounded_launch_diagnostic(detail);
834 let body = format!(
835 "Mjolnir session launch failure\nsession: {session_id}\nat: {}\n\n{detail}\n",
836 now()
837 );
838 atomic_write(&path, body.as_bytes())?;
839 prune_launch_diagnostics(directory)?;
840 Ok(path)
841}
842
843fn bounded_launch_diagnostic(detail: &str) -> String {
844 if detail.len() <= MAX_LAUNCH_DIAGNOSTIC_BYTES {
845 return detail.to_owned();
846 }
847 let mut head_end = MAX_LAUNCH_DIAGNOSTIC_BYTES / 4;
848 while !detail.is_char_boundary(head_end) {
849 head_end -= 1;
850 }
851 let tail_bytes = MAX_LAUNCH_DIAGNOSTIC_BYTES - head_end;
852 let mut tail_start = detail.len() - tail_bytes;
853 while !detail.is_char_boundary(tail_start) {
854 tail_start += 1;
855 }
856 format!(
857 "{}\n\n[... launch diagnostic truncated ...]\n\n{}",
858 &detail[..head_end],
859 &detail[tail_start..]
860 )
861}
862
863fn prune_launch_diagnostics(directory: &Path) -> Result<()> {
864 let mut diagnostics = Vec::new();
865 for entry in std::fs::read_dir(directory)? {
866 let entry = entry?;
867 if !entry
868 .file_name()
869 .to_str()
870 .is_some_and(|name| name.ends_with("-launch-error.txt"))
871 {
872 continue;
873 }
874 diagnostics.push((entry.metadata()?.modified()?, entry.path()));
875 }
876 diagnostics.sort_by_key(|entry| std::cmp::Reverse(entry.0));
877 for (_, path) in diagnostics.into_iter().skip(RETAINED_LAUNCH_DIAGNOSTICS) {
878 std::fs::remove_file(&path)
879 .with_context(|| format!("prune old launch diagnostic {}", path.display()))?;
880 }
881 Ok(())
882}
883
884fn apply_new_session_provisioning_result(
885 state: &mut State,
886 session_id: &str,
887 result: Result<TargetLocator>,
888) -> Result<()> {
889 match result {
890 Ok(locator) => {
891 let record = state.sessions.get_mut(session_id).unwrap();
892 record.target = Some(locator);
893 record.state = SessionState::Disconnected;
896 record.updated_at = now();
897 record.last_error = None;
898 Ok(())
899 }
900 Err(error) => {
901 let record = state.sessions.get_mut(session_id).unwrap();
902 record.state = SessionState::Error;
903 record.target = None;
904 record.updated_at = now();
905 record.last_error = Some(format!("session provisioning failed: {error:#}"));
906 Err(error)
907 }
908 }
909}
910
911fn apply_failed_new_session_launch(state: &mut State, session_id: &str, original_error: &str) {
914 let record = state.sessions.get_mut(session_id).unwrap();
915 record.state = SessionState::Error;
916 record.updated_at = now();
917 record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
918}
919
920pub(super) fn apply_failed_new_session_rollback(
921 state: &mut State,
922 session_id: &str,
923 original_error: &str,
924 cleanup_error: Option<String>,
925) -> anyhow::Error {
926 match cleanup_error {
927 None => {
928 let record = state.sessions.get_mut(session_id).unwrap();
929 record.state = SessionState::Error;
930 record.target = None;
931 record.updated_at = now();
932 record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
933 anyhow::anyhow!("{original_error}; partial target removed and failed session retained")
934 }
935 Some(cleanup_error) => {
936 let failure = format!(
937 "{original_error}; cleanup of the failed session target failed: {cleanup_error}"
938 );
939 let record = state.sessions.get_mut(session_id).unwrap();
940 record.state = SessionState::Error;
941 record.updated_at = now();
942 record.last_error = Some(format!("worker bootstrap failed: {failure}"));
943 anyhow::anyhow!(failure)
944 }
945 }
946}
947
948pub(super) fn install_attached_resources(
949 state: &State,
950 session_id: &str,
951 backend: &targets::TargetLocator,
952 worker_root: &str,
953 executor: &impl CommandExecutor,
954) -> Result<()> {
955 let targets::TargetLocator::AwsEc2 { .. } = backend else {
956 return Ok(());
957 };
958 let session = state
959 .sessions
960 .get(session_id)
961 .with_context(|| format!("unknown session {session_id}"))?;
962 if session.additional_mounts.is_empty() {
963 return Ok(());
964 }
965 for resource in &session.additional_mounts {
966 let install = targets::command_on_locator(
967 backend,
968 session_id,
969 vec![
970 format!("{worker_root}/hel"),
971 "worker".into(),
972 "install-resource".into(),
973 "--destination".into(),
974 resource.destination.to_string_lossy().into_owned(),
975 ],
976 "stream attached resource",
977 )?;
978 mj_checkpoint::resources::stream_resource(&resource.source, |stream| {
979 execute_checked_with_stdin(executor, &install, stream).map(|_| ())
980 })
981 .with_context(|| format!("stream attached resource {}", resource.source.display()))?;
982 }
983 Ok(())
984}
985
986pub(super) fn execute_concurrent_lanes<A: Send, B: Send>(
990 first: impl FnOnce() -> Result<A> + Send,
991 second: impl FnOnce() -> Result<B> + Send,
992) -> Result<(A, B)> {
993 std::thread::scope(|scope| {
994 let second = scope.spawn(second);
995 let first = first();
996 let second = second.join().unwrap_or_else(|panic| {
997 Err(anyhow::anyhow!(
998 "concurrent target lane panicked: {}",
999 targets::command_thread_panic_message(panic.as_ref())
1000 ))
1001 });
1002 match (first, second) {
1003 (Err(error), _) => Err(error),
1004 (Ok(_), Err(error)) => Err(error),
1005 (Ok(first), Ok(second)) => Ok((first, second)),
1006 }
1007 })
1008}
1009
1010fn execute_repository_setup(
1011 plan: &targets::CommandPlan,
1012 executor: &(impl CommandExecutor + Sync),
1013) -> Result<()> {
1014 if plan.commands.is_empty() {
1015 return Ok(());
1016 }
1017 let _cloning = ProvisionStageGuard::new(executor, ProvisionStage::Cloning);
1018 plan.execute_concurrent(executor).map(|_| ())
1019}
1020
1021#[cfg(test)]
1028fn provision_target(
1029 plan: &targets::CommandPlan,
1030 target: &targets::TargetTemplate,
1031 session_id: &str,
1032 executor: &(impl CommandExecutor + Sync),
1033 discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1034) -> Result<TargetLocator> {
1035 let Some((creation, remainder)) = plan.split_at_target_creation() else {
1036 return discover(&plan.execute_concurrent(executor)?);
1038 };
1039 let mut outputs = creation.execute_concurrent(executor)?;
1040 let result = match remainder.execute_concurrent(executor) {
1041 Ok(rest) => {
1042 outputs.extend(rest);
1043 discover(&outputs)
1044 }
1045 Err(error) => Err(error),
1046 };
1047 result.map_err(|error| {
1048 match cleanup_failed_provision(target, session_id, outputs.first(), executor) {
1049 Some(note) => error.context(note),
1050 None => error,
1051 }
1052 })
1053}
1054
1055fn provision_target_creation(
1059 plan: &targets::CommandPlan,
1060 target: &targets::TargetTemplate,
1061 session_id: &str,
1062 executor: &(impl CommandExecutor + Sync),
1063 discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1064) -> Result<(TargetLocator, targets::CommandPlan)> {
1065 let Some((creation, remainder)) = plan.split_at_target_creation() else {
1066 let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1069 plan.execute_concurrent(executor)
1070 })?;
1071 return discover(&outputs).map(|locator| {
1072 (
1073 locator,
1074 targets::CommandPlan {
1075 description: plan.description.clone(),
1076 commands: Vec::new(),
1077 },
1078 )
1079 });
1080 };
1081 let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1085 creation.execute_concurrent(executor)
1086 })?;
1087 discover(&outputs)
1088 .map(|locator| (locator, remainder))
1089 .map_err(|error| {
1090 match cleanup_failed_provision(target, session_id, outputs.first(), executor) {
1091 Some(note) => error.context(note),
1092 None => error,
1093 }
1094 })
1095}
1096
1097fn cleanup_failed_provision(
1104 target: &targets::TargetTemplate,
1105 session_id: &str,
1106 create_output: Option<&CommandOutput>,
1107 executor: &impl CommandExecutor,
1108) -> Option<String> {
1109 let locator = provisioned_locator(target, session_id, create_output)?;
1110 let leak = format!(
1111 "the resource may still exist; find it via its dev.mj.session={session_id} label/tag"
1112 );
1113 let plan = match targets::close_plan(&locator, session_id) {
1114 Ok(plan) => plan,
1115 Err(error) => {
1116 tracing::warn!(
1117 session_id,
1118 error = format!("{error:#}"),
1119 "could not build provisioning cleanup plan"
1120 );
1121 return Some(format!("cleanup FAILED: {error:#}; {leak}"));
1122 }
1123 };
1124 let purpose = plan
1125 .commands
1126 .iter()
1127 .map(|command| command.purpose.clone())
1128 .collect::<Vec<_>>()
1129 .join("; ");
1130 let Err(error) = plan.execute(executor) else {
1131 return Some(format!("cleanup succeeded: {purpose}"));
1132 };
1133 match targets::cleanup_target_is_confirmed_absent(&locator, session_id, executor) {
1134 Ok(true) => Some(format!("cleanup succeeded: {purpose}")),
1135 Ok(false) => {
1136 tracing::warn!(
1137 session_id,
1138 error = format!("{error:#}"),
1139 "provisioning cleanup failed and the target may still exist"
1140 );
1141 Some(format!("cleanup FAILED ({purpose}): {error:#}; {leak}"))
1142 }
1143 Err(confirm_error) => {
1144 tracing::warn!(
1145 session_id,
1146 error = format!("{confirm_error:#}"),
1147 "could not confirm whether the failed provisioning target was removed"
1148 );
1149 Some(format!(
1150 "cleanup FAILED ({purpose}): {error:#}; checking whether it was removed also failed: {confirm_error:#}; {leak}"
1151 ))
1152 }
1153 }
1154}
1155
1156fn provisioned_locator(
1161 target: &targets::TargetTemplate,
1162 session_id: &str,
1163 create_output: Option<&CommandOutput>,
1164) -> Option<targets::TargetLocator> {
1165 let container_id = || targets::resource_name(session_id).ok();
1166 Some(match target {
1167 targets::TargetTemplate::LocalBare => return None,
1170 targets::TargetTemplate::LocalPodman(container) => targets::TargetLocator::LocalPodman {
1171 borrowed_from: None,
1172 container_id: container_id()?,
1173 workspace_storage: targets::podman_workspace_locator(container, session_id).ok()?,
1174 },
1175 targets::TargetTemplate::LocalDocker(_) => targets::TargetLocator::LocalDocker {
1176 borrowed_from: None,
1177 container_id: container_id()?,
1178 },
1179 targets::TargetTemplate::AppleContainer(_) => targets::TargetLocator::AppleContainer {
1180 borrowed_from: None,
1181 container_id: container_id()?,
1182 },
1183 targets::TargetTemplate::SshPodman { ssh, container } => {
1184 targets::TargetLocator::SshPodman {
1185 borrowed_from: None,
1186 ssh: ssh.clone(),
1187 container_id: container_id()?,
1188 workspace_storage: targets::podman_workspace_locator(container, session_id).ok()?,
1189 }
1190 }
1191 targets::TargetTemplate::SshDocker { ssh, .. } => targets::TargetLocator::SshDocker {
1192 borrowed_from: None,
1193 ssh: ssh.clone(),
1194 container_id: container_id()?,
1195 },
1196 targets::TargetTemplate::SshBare { ssh, .. } => targets::TargetLocator::SshBare {
1197 ssh: ssh.clone(),
1198 workspace: targets::workspace_for(target, session_id).ok()?,
1199 worker_id: None,
1200 },
1201 targets::TargetTemplate::AwsEc2(aws) => targets::TargetLocator::AwsEc2 {
1202 profile: aws.profile.clone(),
1203 region: aws.region.clone(),
1204 instance_id: serde_json::from_slice::<serde_json::Value>(&create_output?.stdout)
1205 .ok()?
1206 .pointer("/Instances/0/InstanceId")?
1207 .as_str()?
1208 .to_owned(),
1209 ssh: aws.ssh.clone(),
1210 workspace: targets::workspace_for(target, session_id).ok()?,
1211 },
1212 })
1213}
1214
1215pub(super) fn enforce_overlay_capable_mounts(
1226 target: &targets::TargetTemplate,
1227 mounts: &mut [targets::AdditionalMount],
1228 executor: &impl CommandExecutor,
1229) -> Vec<String> {
1230 let ssh = match target {
1231 targets::TargetTemplate::LocalPodman(_) | targets::TargetTemplate::LocalDocker(_) => None,
1232 targets::TargetTemplate::SshPodman { ssh, .. }
1233 | targets::TargetTemplate::SshDocker { ssh, .. } => Some(ssh),
1234 _ => return Vec::new(),
1235 };
1236 let overlaid = mounts
1237 .iter()
1238 .filter(|mount| mount.access == targets::MountAccess::Cow)
1239 .map(|mount| mount.source.clone())
1240 .collect::<Vec<_>>();
1241 if overlaid.is_empty() {
1242 return Vec::new();
1243 }
1244 let filesystems = match targets::probe_filesystem_types(ssh, &overlaid, executor) {
1245 Ok(filesystems) => filesystems,
1246 Err(error) => {
1247 tracing::warn!(
1248 error = format!("{error:#}"),
1249 "could not probe attached-directory filesystems; preserving overlay mounts"
1250 );
1251 return vec![format!(
1252 "Could not read the filesystem under the attached directories, so they keep the \
1253 copy-on-write overlay: {error:#}"
1254 )];
1255 }
1256 };
1257 let mut notices = Vec::new();
1258 for (mount, filesystem) in mounts
1259 .iter_mut()
1260 .filter(|mount| mount.access == targets::MountAccess::Cow)
1261 .zip(filesystems)
1262 {
1263 let Some(reason) = targets::overlay_unsupported_filesystem(&filesystem) else {
1264 continue;
1265 };
1266 mount.access = mount.access.without_overlay();
1267 notices.push(format!(
1268 "Mounted {} read-only: the overlay is unreliable on {filesystem} ({reason}).",
1269 mount.source.display()
1270 ));
1271 }
1272 notices
1273}
1274
1275static IMAGE_USERS: std::sync::LazyLock<std::sync::Mutex<BTreeMap<String, targets::ImageUser>>> =
1280 std::sync::LazyLock::new(std::sync::Mutex::default);
1281
1282pub(super) fn podman_image_user(
1289 target: &targets::TargetTemplate,
1290 executor: &impl CommandExecutor,
1291) -> Option<targets::ImageUser> {
1292 let (ssh, container) = match target {
1293 targets::TargetTemplate::LocalPodman(container) => (None, container),
1294 targets::TargetTemplate::SshPodman { ssh, container } => (Some(ssh), container),
1295 _ => return None,
1296 };
1297 let image = container.image.as_str();
1298 let key = format!(
1299 "{}|{image}",
1300 ssh.map_or("local", |ssh| ssh.destination.as_str())
1301 );
1302 if let Some(cached) = IMAGE_USERS.lock().expect("image user cache").get(&key) {
1303 return Some(*cached);
1304 }
1305 match crate::image_pull_gate::with_image_ready(target, executor, || {
1309 targets::probe_image_user(ssh, container, executor)
1310 }) {
1311 Ok(user) => {
1312 IMAGE_USERS
1313 .lock()
1314 .expect("image user cache")
1315 .insert(key, user);
1316 Some(user)
1317 }
1318 Err(error) => {
1319 tracing::warn!(
1320 image,
1321 error = format!("{error:#}"),
1322 "could not read the container image user; keeping Podman's default user mapping"
1323 );
1324 executor.notify_notice(&format!(
1325 "Could not read the user of image {image}, so the container runs with Podman's \
1326 default user mapping and may not be able to write to an attached directory: \
1327 {error:#}"
1328 ));
1329 None
1330 }
1331 }
1332}
1333
1334pub(super) struct StagedExecutor<'a, E: CommandExecutor> {
1338 inner: &'a E,
1339 stage: ProvisionStage,
1340 _guard: ProvisionStageGuard<'a, E>,
1341}
1342
1343impl<'a, E: CommandExecutor> StagedExecutor<'a, E> {
1344 pub(crate) fn new(inner: &'a E, stage: ProvisionStage) -> Self {
1345 Self {
1346 inner,
1347 stage,
1348 _guard: ProvisionStageGuard::new(inner, stage),
1349 }
1350 }
1351
1352 fn staged(&self, command: &CommandSpec) -> CommandSpec {
1353 if command.stage.is_some() {
1354 return command.clone();
1355 }
1356 command.clone().stage(self.stage)
1357 }
1358}
1359
1360impl<E: CommandExecutor> CommandExecutor for StagedExecutor<'_, E> {
1361 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1362 self.inner.execute(&self.staged(command))
1363 }
1364
1365 fn cancellation_requested(&self) -> bool {
1366 self.inner.cancellation_requested()
1367 }
1368
1369 fn stage_started(&self, stage: ProvisionStage) {
1370 self.inner.stage_started(stage);
1371 }
1372
1373 fn stage_finished(&self, stage: ProvisionStage) {
1374 self.inner.stage_finished(stage);
1375 }
1376
1377 fn notify_notice(&self, notice: &str) {
1378 self.inner.notify_notice(notice);
1379 }
1380
1381 fn execute_with_stdin(
1382 &self,
1383 command: &CommandSpec,
1384 input: &mut (dyn std::io::Read + Send),
1385 ) -> Result<CommandOutput> {
1386 self.inner.execute_with_stdin(&self.staged(command), input)
1387 }
1388}
1389
1390fn execute_checked_with_stdin(
1391 executor: &impl CommandExecutor,
1392 command: &CommandSpec,
1393 input: &mut (dyn std::io::Read + Send),
1394) -> Result<CommandOutput> {
1395 let output = executor.execute_with_stdin(command, input)?;
1396 if output.status != 0 {
1397 bail!(
1398 "{} failed with status {}: {}",
1399 command.purpose,
1400 output.status,
1401 String::from_utf8_lossy(&output.stderr)
1402 );
1403 }
1404 Ok(output)
1405}
1406
1407pub(super) fn install_inherited_git_settings(
1408 executor: &impl CommandExecutor,
1409 locator: &targets::TargetLocator,
1410 session_id: &str,
1411) -> Result<()> {
1412 let settings = if inherits_controller_git_settings(locator) {
1413 controller_git_settings()?
1414 } else {
1415 BTreeMap::new()
1416 };
1417 for command in inherited_git_setting_commands(locator, session_id, settings)? {
1418 execute_checked(executor, command)?;
1419 }
1420 Ok(())
1421}
1422
1423fn inherits_controller_git_settings(locator: &targets::TargetLocator) -> bool {
1424 !matches!(
1425 locator,
1426 targets::TargetLocator::LocalBare { .. } | targets::TargetLocator::SshBare { .. }
1427 )
1428}
1429
1430fn inherited_git_setting_commands(
1431 locator: &targets::TargetLocator,
1432 session_id: &str,
1433 settings: BTreeMap<String, String>,
1434) -> Result<Vec<CommandSpec>> {
1435 if matches!(locator, targets::TargetLocator::SshBare { .. }) {
1436 return Ok(Vec::new());
1437 }
1438 settings
1439 .into_iter()
1440 .map(|(key, value)| {
1441 targets::command_on_locator(
1442 locator,
1443 session_id,
1444 vec![
1445 "git".into(),
1446 "config".into(),
1447 "--global".into(),
1448 "--replace-all".into(),
1449 "--".into(),
1450 key.clone(),
1451 value,
1452 ],
1453 format!("inherit Git setting {key}"),
1454 )
1455 })
1456 .collect()
1457}
1458
1459fn controller_git_settings() -> Result<BTreeMap<String, String>> {
1460 let output = match Command::new("git")
1461 .args(["config", "--global", "--includes", "--null", "--list"])
1462 .stdin(Stdio::null())
1463 .output()
1464 {
1465 Ok(output) => output,
1466 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(BTreeMap::new()),
1467 Err(error) => return Err(error).context("read controller Git configuration"),
1468 };
1469 if !output.status.success() {
1470 bail!(
1471 "read controller Git configuration failed with status {}: {}",
1472 output.status,
1473 String::from_utf8_lossy(&output.stderr).trim()
1474 );
1475 }
1476 parse_inherited_git_settings(&output.stdout)
1477}
1478
1479fn parse_inherited_git_settings(output: &[u8]) -> Result<BTreeMap<String, String>> {
1480 let mut settings = BTreeMap::new();
1481 for entry in output
1482 .split(|byte| *byte == 0)
1483 .filter(|entry| !entry.is_empty())
1484 {
1485 let entry = std::str::from_utf8(entry).context("decode controller Git configuration")?;
1486 let (key, value) = entry
1487 .split_once('\n')
1488 .with_context(|| format!("controller Git returned malformed entry {entry:?}"))?;
1489 let key = key.to_ascii_lowercase();
1490 if INHERITED_GIT_SETTINGS.contains(&key.as_str()) {
1491 settings.insert(key, value.to_owned());
1492 }
1493 }
1494 Ok(settings)
1495}
1496
1497#[cfg(test)]
1498mod tests;