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