1use std::collections::{BTreeMap, HashMap};
4use std::path::{Path, PathBuf};
5use std::process::{Command, Stdio};
6use std::sync::{Arc, Mutex, OnceLock};
7use std::time::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, CommandExecutor, CommandOutput, CommandSpec, ProvisionStage, ProvisionStageGuard,
16};
17
18use super::backend::{
19 ContainerOverrides, TargetCheck, backend_locator, backend_session_bundle, backend_target,
20 configure_github_token_environment, controller_github_token, preflight_target,
21 use_github_https_urls,
22};
23use super::git_cache;
24use super::readiness::{connect_started_worker, wait_for_native_session_in_stage};
25use super::worker_binary::{bridge_readiness_stage, start_worker_durably, worker_probe_diagnosis};
26use super::{Controller, execute_checked, now};
27
28const INHERITED_GIT_SETTINGS: &[&str] = &[
29 "diff.algorithm",
30 "fetch.prune",
31 "fetch.prunetags",
32 "init.defaultbranch",
33 "merge.conflictstyle",
34 "pull.ff",
35 "pull.rebase",
36 "push.autosetupremote",
37 "push.default",
38 "rebase.autostash",
39 "rerere.autoupdate",
40 "rerere.enabled",
41 "user.email",
42 "user.name",
43];
44
45const CONTAINER_START_ADMISSION: usize = 2;
56
57pub(super) fn container_start_gate(
63 locator: &targets::TargetLocator,
64) -> Option<Arc<tokio::sync::Semaphore>> {
65 static GATES: OnceLock<Mutex<HashMap<String, Arc<tokio::sync::Semaphore>>>> = OnceLock::new();
66 let container = match locator {
67 targets::TargetLocator::LocalPodman { container_id, .. }
68 | targets::TargetLocator::LocalDocker { container_id, .. }
69 | targets::TargetLocator::AppleContainer { container_id, .. }
70 | targets::TargetLocator::SshPodman { container_id, .. }
71 | targets::TargetLocator::SshDocker { container_id, .. } => container_id.clone(),
72 targets::TargetLocator::LocalBare { .. }
73 | targets::TargetLocator::SshBare { .. }
74 | targets::TargetLocator::AwsEc2 { .. } => return None,
75 };
76 let gates = GATES.get_or_init(|| Mutex::new(HashMap::new()));
77 let mut gates = gates
78 .lock()
79 .unwrap_or_else(std::sync::PoisonError::into_inner);
80 Some(Arc::clone(gates.entry(container).or_insert_with(|| {
81 Arc::new(tokio::sync::Semaphore::new(CONTAINER_START_ADMISSION))
82 })))
83}
84
85fn subagent_start_is_retryable(error: &anyhow::Error) -> bool {
93 if mj_core::refusal::Refusal::of(error).is_some() {
94 return false;
95 }
96 if format!("{error:#}").contains("operation cancelled") {
97 return false;
98 }
99 error
100 .downcast_ref::<super::readiness::WorkerStartupFailure>()
101 .is_some_and(|failure| !failure.reached_socket)
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq)]
105pub(super) enum ProvisioningFailureDisposition {
106 Discard,
108 Preserve,
110}
111
112impl Controller {
113 pub async fn provision_session_controlled_with_commit(
114 &mut self,
115 session_id: &str,
116 executor: &(impl CommandExecutor + Sync),
117 grant_commit: impl FnOnce() -> Result<()>,
118 ) -> Result<()> {
119 crate::worker_lifecycle::run(
120 session_id,
121 "provision session controlled with commit",
122 executor,
123 async {
124 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
125 let github_token = controller_github_token();
126 let repositories = self
127 .provision_session_target_with_failure_disposition(
128 session_id,
129 executor,
130 github_token.as_deref(),
131 ProvisioningFailureDisposition::Discard,
132 )
133 .await?;
134 let setup = execute_concurrent_lanes(
135 || execute_repository_setup(&repositories, executor),
136 || self.install_worker_payload(session_id, executor),
137 );
138 let result = match setup {
139 Ok(((), (backend, worker_root))) => {
140 self.connect_and_start_worker(
141 session_id,
142 executor,
143 &backend,
144 &worker_root,
145 true,
146 )
147 .await
148 }
149 Err(error) => Err(error),
150 };
151 match result {
152 Ok(native_session_id) => {
153 if let Err(error) = grant_commit() {
154 return Err(self.rollback_failed_new_session(session_id, error)?);
155 }
156 self.mark_worker_connected(session_id, native_session_id)
157 }
158 Err(error) => Err(self.rollback_failed_new_session(session_id, error)?),
159 }
160 },
161 )
162 .await
163 }
164
165 pub async fn provision_subagent_session_controlled(
168 &mut self,
169 session_id: &str,
170 executor: &(impl CommandExecutor + Sync),
171 ) -> Result<()> {
172 crate::worker_lifecycle::run(
173 session_id,
174 "provision subagent session controlled",
175 executor,
176 async {
177 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
178 let mut attempts: Vec<String> = Vec::new();
185 loop {
186 let (result, placement) =
187 self.attempt_subagent_start(session_id, executor).await;
188 let error = match result {
189 Ok(native_session_id) => {
190 return self.mark_worker_connected(session_id, native_session_id);
191 }
192 Err(error) => error,
193 };
194 let retry = attempts.is_empty() && subagent_start_is_retryable(&error);
195 let error = match &placement {
196 Some((backend, _)) => super::subagent_park::explain_process_exhaustion(
197 error, backend, session_id,
198 ),
199 None => error,
200 };
201 attempts.push(format!("{error:#}"));
202 let error = if attempts.len() > 1 {
203 anyhow::anyhow!(
204 "{}",
205 attempts
206 .iter()
207 .enumerate()
208 .map(|(index, error)| format!("attempt {}: {error}", index + 1))
209 .collect::<Vec<_>>()
210 .join("; ")
211 )
212 } else {
213 error
214 };
215 let diagnostic = note_new_session_launch_failure(session_id, &error);
216 let previous = self
217 .state
218 .sessions
219 .get(session_id)
220 .context("failed child disappeared")?
221 .clone();
222 let record = self.state.sessions.get_mut(session_id).unwrap();
223 record.state = SessionState::StartupCleanup;
224 record.updated_at = now();
225 record.last_error = Some(format!("sub-agent startup failed: {diagnostic}"));
226 self.persist_session_transition_or_restore(
227 session_id,
228 &previous,
229 "record failed startup before teardown",
230 )?;
231 let cleanup = super::failed_launch_cleanup_executor();
232 if let Err(cleanup_error) =
233 self.finish_failed_startup_controlled(session_id, &cleanup, retry)
234 {
235 return Err(error.context(format!(
236 "startup cleanup remains pending: {cleanup_error:#}"
237 )));
238 }
239 if !retry {
240 return Err(error);
241 }
242 tracing::warn!(
243 session_id,
244 error = format!("{error:#}"),
245 "sub-agent worker never started and was stopped; retrying once"
246 );
247 }
248 },
249 )
250 .await
251 }
252
253 pub(crate) fn cleanup_failed_startup_controlled(
255 &mut self,
256 session_id: &str,
257 executor: &impl CommandExecutor,
258 ) -> Result<()> {
259 self.finish_failed_startup_controlled(session_id, executor, false)
260 }
261
262 fn finish_failed_startup_controlled(
263 &mut self,
264 session_id: &str,
265 executor: &impl CommandExecutor,
266 retry_after_stop: bool,
267 ) -> Result<()> {
268 crate::worker_lifecycle::run_blocking(
269 session_id,
270 "finish failed startup controlled",
271 executor,
272 || {
273 let previous = self
274 .state
275 .sessions
276 .get(session_id)
277 .context("cleanup session disappeared")?
278 .clone();
279 ensure!(
280 previous.state == SessionState::StartupCleanup,
281 "session is not awaiting startup cleanup"
282 );
283 let child = self.state.subagents.contains_key(session_id);
284 ensure!(
285 !retry_after_stop || child,
286 "only unaccepted child startup may retry"
287 );
288 let outcome = (|| -> Result<()> {
289 if let Some(locator) = &previous.target {
290 let backend = backend_locator(locator, &previous, &self.config)?;
291 if child {
292 let root = targets::worker_root(&backend, session_id)?;
293 super::worker_binary::stop_worker(
294 &crate::worker_lifecycle::require(session_id)?,
295 executor,
296 &backend,
297 &root,
298 )?;
299 } else {
300 targets::close_plan(&backend, session_id)?.execute(executor)?;
301 }
302 }
303 if !child {
304 self.cleanup_new_session_worktree_after_failure(session_id, executor)?;
305 }
306 Ok(())
307 })();
308 let record = self.state.sessions.get_mut(session_id).unwrap();
309 record.updated_at = now();
310 let cause = previous
311 .last_error
312 .as_deref()
313 .unwrap_or("startup failed")
314 .split("; startup cleanup failed:")
315 .next()
316 .unwrap()
317 .to_owned();
318 match &outcome {
319 Ok(()) => {
320 record.state = if retry_after_stop {
321 SessionState::Provisioning
322 } else {
323 SessionState::Error
324 };
325 record.last_error = (!retry_after_stop).then_some(cause);
326 if !child {
327 record.target = None;
328 }
329 }
330 Err(error) => {
331 record.last_error =
332 Some(format!("{cause}; startup cleanup failed: {error:#}"));
333 }
334 }
335 let record = self.state.sessions.get_mut(session_id).unwrap();
336 if outcome.is_ok() && child && !retry_after_stop {
337 record.archived = true;
338 }
339 if let Err(error) = crate::database::save_startup_cleanup_outcome(
340 record,
341 outcome.is_ok() && child && !retry_after_stop,
342 ) {
343 self.state.sessions.insert(session_id.to_owned(), previous);
344 return Err(error).context("record startup cleanup outcome");
345 }
346 outcome
347 },
348 )
349 }
350
351 async fn attempt_subagent_start(
355 &mut self,
356 session_id: &str,
357 executor: &(impl CommandExecutor + Sync),
358 ) -> (
359 Result<Option<String>>,
360 Option<(targets::TargetLocator, String)>,
361 ) {
362 match self.worker_placement(session_id) {
365 Ok((backend, worker_root)) => {
366 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
367 let prepared = self
368 .prepare_worker_files(session_id, &backend, &worker_root, syncing)
369 .and_then(|()| self.prepare_subagent_report_dir(session_id, &backend, syncing));
370 let result = match prepared {
371 Ok(()) => {
372 let gate = container_start_gate(&backend);
375 let _admitted = match &gate {
376 Some(gate) => gate.acquire().await.ok(),
377 None => None,
378 };
379 self.connect_and_start_worker(
380 session_id,
381 executor,
382 &backend,
383 &worker_root,
384 false,
385 )
386 .await
387 }
388 Err(error) => Err(error),
389 };
390 (result, Some((backend, worker_root)))
391 }
392 Err(error) => (Err(error), None),
393 }
394 }
395
396 fn rollback_failed_new_session(
397 &mut self,
398 session_id: &str,
399 error: anyhow::Error,
400 ) -> Result<anyhow::Error> {
401 self.rollback_failed_new_session_with(
402 session_id,
403 error,
404 &super::failed_launch_cleanup_executor(),
407 )
408 }
409
410 fn rollback_failed_new_session_with(
411 &mut self,
412 session_id: &str,
413 error: anyhow::Error,
414 target_cleanup_executor: &impl CommandExecutor,
415 ) -> Result<anyhow::Error> {
416 let session = self
417 .state
418 .sessions
419 .get(session_id)
420 .with_context(|| format!("unknown session {session_id}"))?
421 .clone();
422 let original = note_new_session_launch_failure(session_id, &error);
429 apply_failed_new_session_launch(&mut self.state, session_id, &original);
430 self.persist_session_transition_or_restore(
431 session_id,
432 &session,
433 "record the failed launch before removing its target",
434 )
435 .map_err(|persist_error| {
436 persist_error.context(format!(
437 "{original}; its target was left in place because the failure could not be recorded"
438 ))
439 })?;
440 let target_cleanup = match session.target.as_ref() {
441 Some(locator) => (|| -> Result<()> {
442 let backend = backend_locator(locator, &session, &self.config)?;
443 targets::close_plan(&backend, session_id)?
444 .execute(target_cleanup_executor)
445 .map(|_| ())
446 })(),
447 None => Ok(()),
448 };
449 let cleanup_error = target_cleanup
452 .and_then(|()| {
453 self.cleanup_new_session_worktree_after_failure(session_id, target_cleanup_executor)
454 })
455 .err()
456 .map(|error| format!("{error:#}"));
457 if let Some(cleanup_error) = &cleanup_error {
458 tracing::warn!(
459 session_id,
460 error = %cleanup_error,
461 "new-session rollback cleanup reported failures"
462 );
463 }
464 let failure = apply_failed_new_session_rollback(
465 &mut self.state,
466 session_id,
467 &original,
468 cleanup_error,
469 );
470 self.persist_session_state(session_id)?;
471 Ok(failure)
472 }
473
474 pub(super) async fn provision_session_with_failure_disposition(
475 &mut self,
476 session_id: &str,
477 executor: &(impl CommandExecutor + Sync),
478 github_token: Option<&str>,
479 failure_disposition: ProvisioningFailureDisposition,
480 ) -> Result<()> {
481 let repositories = self
482 .provision_session_target_with_failure_disposition(
483 session_id,
484 executor,
485 github_token,
486 failure_disposition,
487 )
488 .await?;
489 match execute_repository_setup(&repositories, executor) {
490 Ok(()) => Ok(()),
491 Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
492 Err(self.rollback_failed_new_session(session_id, error)?)
493 }
494 Err(error) => Err(error),
495 }
496 }
497
498 async fn provision_session_target_with_failure_disposition(
499 &mut self,
500 session_id: &str,
501 executor: &(impl CommandExecutor + Sync),
502 github_token: Option<&str>,
503 failure_disposition: ProvisioningFailureDisposition,
504 ) -> Result<targets::CommandPlan> {
505 crate::worker_lifecycle::run(session_id, "provision session target with failure disposition", executor, async {
506 let session = self
507 .state
508 .sessions
509 .get(session_id)
510 .with_context(|| format!("unknown session {session_id}"))?
511 .clone();
512 if session.state != SessionState::Provisioning {
513 bail!("session {session_id} is not provisioning");
514 }
515 if let Some(plan) = self.adopt_prepared_ec2_destination(session_id)? {
516 return Ok(plan);
517 }
518 let preparation = (|| {
519 let selected = self
520 .config
521 .targets
522 .get(&session.target_template_id)
523 .context("target template disappeared before provisioning")?;
524 selected.ensure_ready(&session.target_template_id)?;
527 self.config
528 .profiles
529 .get(&session.last_profile)
530 .context("harness profile disappeared before provisioning")?
531 .ensure_ready(&session.last_profile)?;
532 let runtime = mj_core::state::TargetRuntimeSettings::from(selected);
533 if let Some(recorded) = &session.target_runtime {
534 ensure!(
535 recorded == &runtime,
536 "target access settings changed before provisioning; retry with the selected target"
537 );
538 } else {
539 self.state
540 .sessions
541 .get_mut(session_id)
542 .unwrap()
543 .target_runtime = Some(runtime);
544 self.persist_session_state(session_id)?;
545 }
546 let template = self
547 .config
548 .targets
549 .get(&session.target_template_id)
550 .context("target template disappeared during provisioning")?;
551 let profile = self
552 .config
553 .profiles
554 .get(&session.last_profile)
555 .context("harness profile disappeared during provisioning")?;
556 super::worker_binary::preflight_worker_binary(template, executor)?;
557 super::worker_binary::preflight_harness(template, profile, executor)?;
558 self.prepare_managed_raw_worktree(session_id, executor)
559 })();
560 let created_worktree = match preparation {
561 Ok(created) => created,
562 Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
563 return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
564 }
565 Err(error) => return Err(error),
566 };
567 let session = self
568 .state
569 .sessions
570 .get(session_id)
571 .expect("session retained after managed worktree preparation")
572 .clone();
573 let result = (|| {
576 let template = self
577 .config
578 .targets
579 .get(&session.target_template_id)
580 .context("target template disappeared during provisioning")?;
581 if matches!(template, TargetTemplate::AwsEc2 { .. }) {
582 for resource in &session.additional_mounts {
583 ensure!(
584 resource.source.is_dir(),
585 "attached resource source is not a directory: {}",
586 resource.source.display()
587 );
588 }
589 }
590 let mut target = backend_target(
591 template,
592 session.resource_allocation.as_ref(),
593 ContainerOverrides::for_session(&session),
594 )?;
595 let mut runtime_mounts = if matches!(target, targets::TargetTemplate::AwsEc2(_)) {
596 Vec::new()
597 } else {
598 session.additional_mounts.clone()
599 };
600 for notice in enforce_overlay_capable_mounts(&target, &mut runtime_mounts, executor) {
605 executor.notify_notice(¬ice);
606 }
607 let image_user = podman_image_user(&target, executor);
611 let mut bundle = if session.project_directory.is_some() {
612 None
613 } else if let Some(bundle) = self.move_destination_bundle(session_id)? {
614 Some(bundle)
615 } else if failure_disposition == ProvisioningFailureDisposition::Preserve {
616 Some(super::network_git::checkpoint_bundle(&session)?)
617 } else {
618 Some(backend_session_bundle(&session, &self.config, executor)?)
619 };
620 let container_github_token =
621 github_token.filter(|_| configure_github_token_environment(&mut target));
622 if container_github_token.is_some()
623 && let Some(bundle) = bundle.as_mut()
624 {
625 use_github_https_urls(bundle);
626 }
627 preflight_target(template, executor, TargetCheck::Launch)?;
628 let resource_name = crate::database::load_move_operation(session_id)?
629 .filter(|op| {
630 op.workspace_transfer.is_some()
631 && op.phase == mj_core::state::MovePhase::ResumingDestination
632 })
633 .map(|op| targets::move_resource_name(session_id, &op.operation_id))
634 .unwrap_or_else(|| targets::resource_name(session_id))?;
635 let moving_workspace = resource_name != targets::resource_name(session_id)?;
636 let prepared_cache = bundle
637 .as_mut()
638 .filter(|_| !moving_workspace)
639 .and_then(|bundle| {
640 git_cache::prepare(
641 &target,
642 session_id,
643 bundle,
644 &mut runtime_mounts,
645 container_github_token,
646 executor,
647 )
648 });
649 let build_cache = super::mbx::prepare(
652 &target,
653 &session,
654 bundle.as_ref(),
655 prepared_cache.as_ref(),
656 &mut runtime_mounts,
657 executor,
658 );
659 let provision = if let Some(project_directory) = &session.project_directory {
660 targets::provision_bare_project_plan(
661 &target,
662 session_id,
663 &project_directory.to_string_lossy(),
664 )
665 } else {
666 bundle
667 .as_ref()
668 .context("project bundle disappeared during provisioning")
669 .and_then(|bundle| {
670 targets::provision_plan_named(
671 &target,
672 session_id,
673 bundle,
674 &runtime_mounts,
675 image_user,
676 session.container_workspace.as_deref(),
677 &resource_name,
678 )
679 })
680 };
681 let mut provision = match provision {
682 Ok(provision) => provision,
683 Err(error) => {
684 if let Some(cache) = &prepared_cache {
685 let _ = cache.cleanup(executor);
686 }
687 return Err(error);
688 }
689 };
690 if let Some(token) = container_github_token
691 && let Err(error) =
692 provision.provide_target_environment_secret(&target, "GH_TOKEN", token)
693 {
694 if let Some(cache) = &prepared_cache {
695 let _ = cache.cleanup(executor);
696 }
697 return Err(error);
698 }
699
700 let started = Instant::now();
701 let result = provision_target_creation_named(
702 &provision,
703 &target,
704 session_id,
705 executor,
706 &resource_name,
707 |outputs| {
708 super::backend::locator_after_provision_named(
709 template,
710 &target,
711 session_id,
712 outputs.first(),
713 executor,
714 &resource_name,
715 )
716 },
717 )
718 .map(|(locator, remainder)| (locator, remainder, bundle, build_cache));
719 if result.is_err()
720 && let Some(cache) = &prepared_cache
721 {
722 if let Some(locator) =
723 provisioned_locator_named(&target, session_id, None, &resource_name)
724 {
725 let _ = targets::close_plan(&locator, session_id)
726 .and_then(|plan| plan.execute(executor).map(|_| ()));
727 } else {
728 let _ = cache.cleanup(executor);
729 }
730 }
731 tracing::debug!(
732 session_id,
733 elapsed_ms = started.elapsed().as_millis(),
734 "provisioning plan execution completed"
735 );
736 result
737 })();
738 let result = match result {
739 Err(error)
740 if created_worktree
741 && failure_disposition == ProvisioningFailureDisposition::Discard =>
742 {
743 return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
744 }
745 Err(error) if failure_disposition == ProvisioningFailureDisposition::Preserve => {
746 Err(error)
747 }
748 Err(error) => {
749 let detail = note_new_session_launch_failure(session_id, &error);
753 {
754 let record = self.state.sessions.get_mut(session_id).unwrap();
755 record.state = SessionState::Error;
756 record.target = None;
757 record.updated_at = super::now();
758 record.last_error = Some(format!("session provisioning failed: {detail}"));
759 }
760 return match self.persist_session_state(session_id) {
761 Ok(()) => Err(error),
762 Err(persistence_error) => Err(error.context(format!(
763 "persist removal of failed provisioning session {session_id}: {persistence_error:#}"
764 ))),
765 };
766 }
767 Ok((locator, remainder, bundle, build_cache)) => {
768 apply_new_session_provisioning_result(&mut self.state, session_id, Ok(locator))?;
769 self.state
770 .sessions
771 .get_mut(session_id)
772 .expect("session retained after provisioning")
773 .build_cache = build_cache;
774 let session = &self.state.sessions[session_id];
775 let backend = backend_locator(
776 session
777 .target
778 .as_ref()
779 .context("provisioned target disappeared")?,
780 session,
781 &self.config,
782 )?;
783 if matches!(backend, targets::TargetLocator::AwsEc2 { .. }) {
784 targets::provision_on_locator_plan(
785 &backend,
786 session_id,
787 bundle
788 .as_ref()
789 .context("AWS provisioning requires a project bundle")?,
790 )
791 } else {
792 Ok(remainder)
793 }
794 }
795 };
796 let result = match result {
797 Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
798 return Err(self.rollback_failed_new_session(session_id, error)?);
799 }
800 result => result,
801 };
802 if result.is_ok()
803 && let Some(session) = self.state.sessions.get(session_id)
804 && let Some(directory) = session
805 .managed_worktree
806 .as_ref()
807 .map(|worktree| worktree.source_project_directory.clone())
808 .or_else(|| session.project_directory.clone())
809 && let Some(template) = self.config.targets.get(&session.target_template_id)
810 {
811 let host = match template {
812 TargetTemplate::LocalBare => Some("local"),
813 TargetTemplate::SshBare { ssh, .. } => Some(ssh.host.as_str()),
814 _ => None,
815 };
816 if let Some(host) = host {
817 self.state.remember_project_directory(host, &directory);
818 crate::database::remember_project_directory(host, &directory)?;
819 }
820 }
821 self.persist_session_state(session_id)?;
822 result
823
824 }).await
825 }
826
827 pub fn mark_worker_connected(
828 &mut self,
829 session_id: &str,
830 native_session_id: Option<String>,
831 ) -> Result<()> {
832 let session = self
833 .state
834 .sessions
835 .get(session_id)
836 .with_context(|| format!("unknown session {session_id}"))?;
837 if session.target.is_none() {
838 bail!("session {session_id} has no provisioned target");
839 }
840 let updated_at = now();
841 crate::database::mark_session_worker_connected(
842 session_id,
843 native_session_id.as_deref(),
844 &updated_at,
845 )?;
846 let session = self
847 .state
848 .sessions
849 .get_mut(session_id)
850 .expect("session disappeared after its worker connection was saved");
851 session.state = SessionState::Running;
852 if native_session_id.is_some() {
853 session.native_session_id = native_session_id;
854 }
855 session.updated_at = updated_at;
856 session.last_error = None;
857 Ok(())
858 }
859
860 fn install_worker_payload(
861 &self,
862 session_id: &str,
863 executor: &impl CommandExecutor,
864 ) -> Result<(targets::TargetLocator, String)> {
865 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
867 let (backend, worker_root) = self.worker_placement(session_id)?;
868 self.prepare_worker_files(session_id, &backend, &worker_root, syncing)?;
869 install_attached_resources(&self.state, session_id, &backend, &worker_root, syncing)?;
870 Ok((backend, worker_root))
871 }
872
873 async fn connect_and_start_worker(
874 &self,
875 session_id: &str,
876 executor: &impl CommandExecutor,
877 backend: &targets::TargetLocator,
878 worker_root: &str,
879 initialize_workspace: bool,
880 ) -> Result<Option<String>> {
881 crate::worker_lifecycle::run(session_id, "connect and start worker", executor, async {
882 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
883 if initialize_workspace {
884 install_inherited_git_settings(executor, backend, session_id)?;
885 self.initialize_network_workspaces(session_id, backend, syncing)?;
886 }
887 let session = self
888 .state
889 .sessions
890 .get(session_id)
891 .with_context(|| format!("unknown session {session_id}"))?;
892 let profile = self
893 .config
894 .profiles
895 .get(&session.last_profile)
896 .with_context(|| format!("unknown profile {}", session.last_profile))?;
897 let readiness_stage = bridge_readiness_stage(profile);
898 let reconnect = &targets::reconnect_plan(backend, session_id)?.commands[0];
899 let readiness = async {
900 let mut relay = {
901 let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
902 start_worker_durably(
903 &crate::worker_lifecycle::require(session_id)?,
904 self.state.sessions[session_id]
905 .target
906 .as_ref()
907 .context("worker start has no durable target")?,
908 executor,
909 backend,
910 worker_root,
911 )?;
912 connect_started_worker(reconnect, session_id, executor, backend, worker_root)
913 .await?
914 };
915 let native_session_id =
916 wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
917 let owner = crate::worker_lifecycle::require(session_id)?;
918 crate::database::finish_worker_restart(session_id, owner.operation_id())?;
919 Ok(Some(native_session_id))
920 }
921 .await;
922 match readiness {
923 Ok(native_session_id) => Ok(native_session_id),
924 Err(error) => {
925 let error = worker_probe_diagnosis(executor, backend, worker_root, error);
928 Err(error)
929 }
930 }
931 })
932 .await
933 }
934}
935
936const MAX_LAUNCH_DIAGNOSTIC_BYTES: usize = 64 * 1024;
937
938const RETAINED_LAUNCH_DIAGNOSTICS: usize = 20;
939
940pub(super) fn note_new_session_launch_failure(session_id: &str, error: &anyhow::Error) -> String {
946 note_new_session_launch_failure_in(&data_dir().join("diagnostics"), session_id, error)
947}
948
949fn note_new_session_launch_failure_in(
950 directory: &Path,
951 session_id: &str,
952 error: &anyhow::Error,
953) -> String {
954 let original = format!("{error:#}");
955 tracing::warn!(session_id, error = %original, "session launch failed");
956 match persist_launch_failure_to(directory, session_id, &original) {
957 Ok(path) => format!("{original}; full diagnostic saved to {}", path.display()),
958 Err(save_error) => {
959 format!("{original}; saving the local diagnostic failed: {save_error:#}")
960 }
961 }
962}
963
964fn persist_launch_failure_to(directory: &Path, session_id: &str, detail: &str) -> Result<PathBuf> {
965 mj_core::config::validate_id("session", session_id)?;
966 std::fs::create_dir_all(directory).with_context(|| {
967 format!(
968 "create launch diagnostics directory {}",
969 directory.display()
970 )
971 })?;
972 #[cfg(unix)]
973 {
974 use std::os::unix::fs::PermissionsExt;
975 std::fs::set_permissions(directory, std::fs::Permissions::from_mode(0o700))?;
976 }
977 let path = directory.join(format!("{session_id}-launch-error.txt"));
978 let detail = bounded_launch_diagnostic(detail);
979 let body = format!(
980 "Mjolnir session launch failure\nsession: {session_id}\nat: {}\n\n{detail}\n",
981 now()
982 );
983 atomic_write(&path, body.as_bytes())?;
984 prune_launch_diagnostics(directory)?;
985 Ok(path)
986}
987
988fn bounded_launch_diagnostic(detail: &str) -> String {
989 if detail.len() <= MAX_LAUNCH_DIAGNOSTIC_BYTES {
990 return detail.to_owned();
991 }
992 let mut head_end = MAX_LAUNCH_DIAGNOSTIC_BYTES / 4;
993 while !detail.is_char_boundary(head_end) {
994 head_end -= 1;
995 }
996 let tail_bytes = MAX_LAUNCH_DIAGNOSTIC_BYTES - head_end;
997 let mut tail_start = detail.len() - tail_bytes;
998 while !detail.is_char_boundary(tail_start) {
999 tail_start += 1;
1000 }
1001 format!(
1002 "{}\n\n[... launch diagnostic truncated ...]\n\n{}",
1003 &detail[..head_end],
1004 &detail[tail_start..]
1005 )
1006}
1007
1008fn prune_launch_diagnostics(directory: &Path) -> Result<()> {
1009 let mut diagnostics = Vec::new();
1010 for entry in std::fs::read_dir(directory)? {
1011 let entry = entry?;
1012 if !entry
1013 .file_name()
1014 .to_str()
1015 .is_some_and(|name| name.ends_with("-launch-error.txt"))
1016 {
1017 continue;
1018 }
1019 diagnostics.push((entry.metadata()?.modified()?, entry.path()));
1020 }
1021 diagnostics.sort_by_key(|entry| std::cmp::Reverse(entry.0));
1022 for (_, path) in diagnostics.into_iter().skip(RETAINED_LAUNCH_DIAGNOSTICS) {
1023 std::fs::remove_file(&path)
1024 .with_context(|| format!("prune old launch diagnostic {}", path.display()))?;
1025 }
1026 Ok(())
1027}
1028
1029fn apply_new_session_provisioning_result(
1030 state: &mut State,
1031 session_id: &str,
1032 result: Result<TargetLocator>,
1033) -> Result<()> {
1034 match result {
1035 Ok(locator) => {
1036 let record = state.sessions.get_mut(session_id).unwrap();
1037 record.target = Some(locator);
1038 record.state = SessionState::Disconnected;
1041 record.updated_at = now();
1042 record.last_error = None;
1043 Ok(())
1044 }
1045 Err(error) => {
1046 let record = state.sessions.get_mut(session_id).unwrap();
1047 record.state = SessionState::Error;
1048 record.target = None;
1049 record.updated_at = now();
1050 record.last_error = Some(format!("session provisioning failed: {error:#}"));
1051 Err(error)
1052 }
1053 }
1054}
1055
1056fn apply_failed_new_session_launch(state: &mut State, session_id: &str, original_error: &str) {
1059 let record = state.sessions.get_mut(session_id).unwrap();
1060 record.state = SessionState::StartupCleanup;
1061 record.updated_at = now();
1062 record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
1063}
1064
1065pub(super) fn apply_failed_new_session_rollback(
1066 state: &mut State,
1067 session_id: &str,
1068 original_error: &str,
1069 cleanup_error: Option<String>,
1070) -> anyhow::Error {
1071 match cleanup_error {
1072 None => {
1073 let record = state.sessions.get_mut(session_id).unwrap();
1074 record.state = SessionState::Error;
1075 record.target = None;
1076 record.updated_at = now();
1077 record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
1078 anyhow::anyhow!("{original_error}; partial target removed and failed session retained")
1079 }
1080 Some(cleanup_error) => {
1081 let failure = format!(
1082 "{original_error}; cleanup of the failed session target failed: {cleanup_error}"
1083 );
1084 let record = state.sessions.get_mut(session_id).unwrap();
1085 record.state = SessionState::StartupCleanup;
1086 record.updated_at = now();
1087 record.last_error = Some(format!("worker bootstrap failed: {failure}"));
1088 anyhow::anyhow!(failure)
1089 }
1090 }
1091}
1092
1093pub(super) fn install_attached_resources(
1094 state: &State,
1095 session_id: &str,
1096 backend: &targets::TargetLocator,
1097 worker_root: &str,
1098 executor: &impl CommandExecutor,
1099) -> Result<()> {
1100 let targets::TargetLocator::AwsEc2 { .. } = backend else {
1101 return Ok(());
1102 };
1103 let session = state
1104 .sessions
1105 .get(session_id)
1106 .with_context(|| format!("unknown session {session_id}"))?;
1107 if session.additional_mounts.is_empty() {
1108 return Ok(());
1109 }
1110 for resource in &session.additional_mounts {
1111 let install = targets::command_on_locator(
1112 backend,
1113 session_id,
1114 vec![
1115 format!("{worker_root}/hel"),
1116 "worker".into(),
1117 "install-resource".into(),
1118 "--destination".into(),
1119 resource.destination.to_string_lossy().into_owned(),
1120 ],
1121 "stream attached resource",
1122 )?;
1123 mj_checkpoint::resources::stream_resource(&resource.source, |stream| {
1124 execute_checked_with_stdin(executor, &install, stream).map(|_| ())
1125 })
1126 .with_context(|| format!("stream attached resource {}", resource.source.display()))?;
1127 }
1128 Ok(())
1129}
1130
1131pub(super) fn execute_concurrent_lanes<A: Send, B: Send>(
1135 first: impl FnOnce() -> Result<A> + Send,
1136 second: impl FnOnce() -> Result<B> + Send,
1137) -> Result<(A, B)> {
1138 std::thread::scope(|scope| {
1139 let owner = crate::worker_lifecycle::capture();
1140 let second = scope.spawn(move || match owner {
1141 Some(owner) => owner.scope_blocking(second),
1142 None => second(),
1143 });
1144 let first = first();
1145 let second = second.join().unwrap_or_else(|panic| {
1146 Err(anyhow::anyhow!(
1147 "concurrent target lane panicked: {}",
1148 targets::command_thread_panic_message(panic.as_ref())
1149 ))
1150 });
1151 match (first, second) {
1152 (Err(error), _) => Err(error),
1153 (Ok(_), Err(error)) => Err(error),
1154 (Ok(first), Ok(second)) => Ok((first, second)),
1155 }
1156 })
1157}
1158
1159fn execute_repository_setup(
1160 plan: &targets::CommandPlan,
1161 executor: &(impl CommandExecutor + Sync),
1162) -> Result<()> {
1163 if plan.commands.is_empty() {
1164 return Ok(());
1165 }
1166 let _cloning = ProvisionStageGuard::new(executor, ProvisionStage::Cloning);
1167 plan.execute_concurrent(executor).map(|_| ())
1168}
1169
1170#[cfg(test)]
1177fn provision_target(
1178 plan: &targets::CommandPlan,
1179 target: &targets::TargetTemplate,
1180 session_id: &str,
1181 executor: &(impl CommandExecutor + Sync),
1182 discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1183) -> Result<TargetLocator> {
1184 let Some((creation, remainder)) = plan.split_at_target_creation() else {
1185 return discover(&plan.execute_concurrent(executor)?);
1187 };
1188 let mut outputs = creation.execute_concurrent(executor)?;
1189 let result = match remainder.execute_concurrent(executor) {
1190 Ok(rest) => {
1191 outputs.extend(rest);
1192 discover(&outputs)
1193 }
1194 Err(error) => Err(error),
1195 };
1196 result.map_err(|error| {
1197 match cleanup_failed_provision(target, session_id, outputs.first(), executor) {
1198 Some(note) => error.context(note),
1199 None => error,
1200 }
1201 })
1202}
1203
1204#[cfg(test)]
1208fn provision_target_creation(
1209 plan: &targets::CommandPlan,
1210 target: &targets::TargetTemplate,
1211 session_id: &str,
1212 executor: &(impl CommandExecutor + Sync),
1213 discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1214) -> Result<(TargetLocator, targets::CommandPlan)> {
1215 provision_target_creation_named(
1216 plan,
1217 target,
1218 session_id,
1219 executor,
1220 &targets::resource_name(session_id)?,
1221 discover,
1222 )
1223}
1224
1225fn provision_target_creation_named(
1226 plan: &targets::CommandPlan,
1227 target: &targets::TargetTemplate,
1228 session_id: &str,
1229 executor: &(impl CommandExecutor + Sync),
1230 name: &str,
1231 discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1232) -> Result<(TargetLocator, targets::CommandPlan)> {
1233 let Some((creation, remainder)) = plan.split_at_target_creation() else {
1234 let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1237 plan.execute_concurrent(executor)
1238 })?;
1239 return discover(&outputs).map(|locator| {
1240 (
1241 locator,
1242 targets::CommandPlan {
1243 description: plan.description.clone(),
1244 commands: Vec::new(),
1245 },
1246 )
1247 });
1248 };
1249 let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1253 creation.execute_concurrent(executor)
1254 })?;
1255 discover(&outputs)
1256 .map(|locator| (locator, remainder))
1257 .map_err(|error| {
1258 match cleanup_failed_provision_named(
1259 target,
1260 session_id,
1261 outputs.first(),
1262 executor,
1263 name,
1264 ) {
1265 Some(note) => error.context(note),
1266 None => error,
1267 }
1268 })
1269}
1270
1271#[cfg(test)]
1278fn cleanup_failed_provision(
1279 target: &targets::TargetTemplate,
1280 session_id: &str,
1281 create_output: Option<&CommandOutput>,
1282 executor: &impl CommandExecutor,
1283) -> Option<String> {
1284 cleanup_failed_provision_named(
1285 target,
1286 session_id,
1287 create_output,
1288 executor,
1289 &targets::resource_name(session_id).ok()?,
1290 )
1291}
1292
1293fn cleanup_failed_provision_named(
1294 target: &targets::TargetTemplate,
1295 session_id: &str,
1296 create_output: Option<&CommandOutput>,
1297 executor: &impl CommandExecutor,
1298 name: &str,
1299) -> Option<String> {
1300 let locator = provisioned_locator_named(target, session_id, create_output, name)?;
1301
1302 let leak = format!(
1303 "the resource may still exist; find it via its dev.mj.session={session_id} label/tag"
1304 );
1305 let cleanup = if name == targets::resource_name(session_id).ok()? {
1306 targets::close_plan(&locator, session_id)
1307 } else {
1308 targets::retire_move_target_plan(&locator, session_id)
1309 };
1310 let plan = match cleanup {
1311 Ok(plan) => plan,
1312 Err(error) => {
1313 tracing::warn!(
1314 session_id,
1315 error = format!("{error:#}"),
1316 "could not build provisioning cleanup plan"
1317 );
1318 return Some(format!("cleanup FAILED: {error:#}; {leak}"));
1319 }
1320 };
1321 let purpose = plan
1322 .commands
1323 .iter()
1324 .map(|command| command.purpose.clone())
1325 .collect::<Vec<_>>()
1326 .join("; ");
1327 let Err(error) = plan.execute(executor) else {
1328 return Some(format!("cleanup succeeded: {purpose}"));
1329 };
1330 match targets::cleanup_target_is_confirmed_absent(&locator, session_id, executor) {
1331 Ok(true) => Some(format!("cleanup succeeded: {purpose}")),
1332 Ok(false) => {
1333 tracing::warn!(
1334 session_id,
1335 error = format!("{error:#}"),
1336 "provisioning cleanup failed and the target may still exist"
1337 );
1338 Some(format!("cleanup FAILED ({purpose}): {error:#}; {leak}"))
1339 }
1340 Err(confirm_error) => {
1341 tracing::warn!(
1342 session_id,
1343 error = format!("{confirm_error:#}"),
1344 "could not confirm whether the failed provisioning target was removed"
1345 );
1346 Some(format!(
1347 "cleanup FAILED ({purpose}): {error:#}; checking whether it was removed also failed: {confirm_error:#}; {leak}"
1348 ))
1349 }
1350 }
1351}
1352
1353fn provisioned_locator_named(
1358 target: &targets::TargetTemplate,
1359 session_id: &str,
1360 create_output: Option<&CommandOutput>,
1361 name: &str,
1362) -> Option<targets::TargetLocator> {
1363 let container_id = || Some(name.to_owned());
1364
1365 Some(match target {
1366 targets::TargetTemplate::LocalBare => return None,
1369 targets::TargetTemplate::LocalPodman(container) => targets::TargetLocator::LocalPodman {
1370 borrowed_from: None,
1371 container_id: container_id()?,
1372 workspace_storage: targets::podman_workspace_locator_named(container, name).ok()?,
1373 },
1374 targets::TargetTemplate::LocalDocker(_) => targets::TargetLocator::LocalDocker {
1375 borrowed_from: None,
1376 container_id: container_id()?,
1377 },
1378 targets::TargetTemplate::AppleContainer(_) => targets::TargetLocator::AppleContainer {
1379 borrowed_from: None,
1380 container_id: container_id()?,
1381 },
1382 targets::TargetTemplate::SshPodman { ssh, container } => {
1383 targets::TargetLocator::SshPodman {
1384 borrowed_from: None,
1385 ssh: ssh.clone(),
1386 container_id: container_id()?,
1387 workspace_storage: targets::podman_workspace_locator_named(container, name).ok()?,
1388 }
1389 }
1390 targets::TargetTemplate::SshDocker { ssh, .. } => targets::TargetLocator::SshDocker {
1391 borrowed_from: None,
1392 ssh: ssh.clone(),
1393 container_id: container_id()?,
1394 },
1395 targets::TargetTemplate::SshBare { ssh, .. } => targets::TargetLocator::SshBare {
1396 ssh: ssh.clone(),
1397 workspace: targets::workspace_for(target, session_id).ok()?,
1398 worker_id: None,
1399 },
1400 targets::TargetTemplate::AwsEc2(aws) => targets::TargetLocator::AwsEc2 {
1401 profile: aws.profile.clone(),
1402 region: aws.region.clone(),
1403 instance_id: serde_json::from_slice::<serde_json::Value>(&create_output?.stdout)
1404 .ok()?
1405 .pointer("/Instances/0/InstanceId")?
1406 .as_str()?
1407 .to_owned(),
1408 ssh: aws.ssh.clone(),
1409 workspace: targets::workspace_for(target, session_id).ok()?,
1410 },
1411 })
1412}
1413
1414pub(super) fn enforce_overlay_capable_mounts(
1425 target: &targets::TargetTemplate,
1426 mounts: &mut [targets::AdditionalMount],
1427 executor: &impl CommandExecutor,
1428) -> Vec<String> {
1429 let ssh = match target {
1430 targets::TargetTemplate::LocalPodman(_) | targets::TargetTemplate::LocalDocker(_) => None,
1431 targets::TargetTemplate::SshPodman { ssh, .. }
1432 | targets::TargetTemplate::SshDocker { ssh, .. } => Some(ssh),
1433 _ => return Vec::new(),
1434 };
1435 let overlaid = mounts
1436 .iter()
1437 .filter(|mount| mount.access == targets::MountAccess::Cow)
1438 .map(|mount| mount.source.clone())
1439 .collect::<Vec<_>>();
1440 if overlaid.is_empty() {
1441 return Vec::new();
1442 }
1443 if matches!(target, targets::TargetTemplate::LocalDocker(_)) {
1445 match targets::local_docker_vm_share(executor) {
1446 Ok(Some(reason)) => {
1447 return mounts
1448 .iter_mut()
1449 .filter(|mount| mount.access == targets::MountAccess::Cow)
1450 .map(|mount| {
1451 mount.access = mount.access.without_overlay();
1452 format!(
1453 "Mounted {} read-only: {reason}, which cannot back the \
1454 copy-on-write overlay.",
1455 mount.source.display()
1456 )
1457 })
1458 .collect();
1459 }
1460 Ok(None) => {}
1461 Err(error) => tracing::warn!(
1462 error = format!("{error:#}"),
1463 "could not identify the Docker daemon platform; probing this host's filesystems"
1464 ),
1465 }
1466 }
1467 let filesystems = match targets::probe_filesystem_types(ssh, &overlaid, executor) {
1468 Ok(filesystems) => filesystems,
1469 Err(error) => {
1470 tracing::warn!(
1471 error = format!("{error:#}"),
1472 "could not probe attached-directory filesystems; preserving overlay mounts"
1473 );
1474 return vec![format!(
1475 "Could not read the filesystem under the attached directories, so they keep the \
1476 copy-on-write overlay: {error:#}"
1477 )];
1478 }
1479 };
1480 let mut notices = Vec::new();
1481 for (mount, filesystem) in mounts
1482 .iter_mut()
1483 .filter(|mount| mount.access == targets::MountAccess::Cow)
1484 .zip(filesystems)
1485 {
1486 let Some(reason) = targets::overlay_unsupported_filesystem(&filesystem) else {
1487 continue;
1488 };
1489 mount.access = mount.access.without_overlay();
1490 notices.push(format!(
1491 "Mounted {} read-only: the overlay is unreliable on {filesystem} ({reason}).",
1492 mount.source.display()
1493 ));
1494 }
1495 notices
1496}
1497
1498static IMAGE_USERS: std::sync::LazyLock<std::sync::Mutex<BTreeMap<String, targets::ImageUser>>> =
1503 std::sync::LazyLock::new(std::sync::Mutex::default);
1504
1505pub(super) fn podman_image_user(
1512 target: &targets::TargetTemplate,
1513 executor: &impl CommandExecutor,
1514) -> Option<targets::ImageUser> {
1515 let (ssh, container) = match target {
1516 targets::TargetTemplate::LocalPodman(container) => (None, container),
1517 targets::TargetTemplate::SshPodman { ssh, container } => (Some(ssh), container),
1518 _ => return None,
1519 };
1520 let image = container.image.as_str();
1521 let key = format!(
1522 "{}|{image}",
1523 ssh.map_or("local", |ssh| ssh.destination.as_str())
1524 );
1525 if let Some(cached) = IMAGE_USERS.lock().expect("image user cache").get(&key) {
1526 return Some(*cached);
1527 }
1528 match crate::image_pull_gate::with_image_ready(target, executor, || {
1532 targets::probe_image_user(ssh, container, executor)
1533 }) {
1534 Ok(user) => {
1535 IMAGE_USERS
1536 .lock()
1537 .expect("image user cache")
1538 .insert(key, user);
1539 Some(user)
1540 }
1541 Err(error) => {
1542 tracing::warn!(
1543 image,
1544 error = format!("{error:#}"),
1545 "could not read the container image user; keeping Podman's default user mapping"
1546 );
1547 executor.notify_notice(&format!(
1548 "Could not read the user of image {image}, so the container runs with Podman's \
1549 default user mapping and may not be able to write to an attached directory: \
1550 {error:#}"
1551 ));
1552 None
1553 }
1554 }
1555}
1556
1557pub(super) struct StagedExecutor<'a, E: CommandExecutor> {
1561 inner: &'a E,
1562 stage: ProvisionStage,
1563 _guard: ProvisionStageGuard<'a, E>,
1564}
1565
1566impl<'a, E: CommandExecutor> StagedExecutor<'a, E> {
1567 pub(crate) fn new(inner: &'a E, stage: ProvisionStage) -> Self {
1568 Self {
1569 inner,
1570 stage,
1571 _guard: ProvisionStageGuard::new(inner, stage),
1572 }
1573 }
1574
1575 fn staged(&self, command: &CommandSpec) -> CommandSpec {
1576 if command.stage.is_some() {
1577 return command.clone();
1578 }
1579 command.clone().stage(self.stage)
1580 }
1581}
1582
1583impl<E: CommandExecutor> CommandExecutor for StagedExecutor<'_, E> {
1584 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1585 self.inner.execute(&self.staged(command))
1586 }
1587
1588 fn cancellation_requested(&self) -> bool {
1589 self.inner.cancellation_requested()
1590 }
1591
1592 fn stage_started(&self, stage: ProvisionStage) {
1593 self.inner.stage_started(stage);
1594 }
1595
1596 fn stage_finished(&self, stage: ProvisionStage) {
1597 self.inner.stage_finished(stage);
1598 }
1599
1600 fn notify_notice(&self, notice: &str) {
1601 self.inner.notify_notice(notice);
1602 }
1603
1604 fn execute_with_stdin(
1605 &self,
1606 command: &CommandSpec,
1607 input: &mut (dyn std::io::Read + Send),
1608 ) -> Result<CommandOutput> {
1609 self.inner.execute_with_stdin(&self.staged(command), input)
1610 }
1611}
1612
1613fn execute_checked_with_stdin(
1614 executor: &impl CommandExecutor,
1615 command: &CommandSpec,
1616 input: &mut (dyn std::io::Read + Send),
1617) -> Result<CommandOutput> {
1618 let output = executor.execute_with_stdin(command, input)?;
1619 if output.status != 0 {
1620 bail!(
1621 "{} failed with status {}: {}",
1622 command.purpose,
1623 output.status,
1624 String::from_utf8_lossy(&output.stderr)
1625 );
1626 }
1627 Ok(output)
1628}
1629
1630pub(super) fn install_inherited_git_settings(
1631 executor: &impl CommandExecutor,
1632 locator: &targets::TargetLocator,
1633 session_id: &str,
1634) -> Result<()> {
1635 let settings = if inherits_controller_git_settings(locator) {
1636 controller_git_settings()?
1637 } else {
1638 BTreeMap::new()
1639 };
1640 for command in inherited_git_setting_commands(locator, session_id, settings)? {
1641 execute_checked(executor, command)?;
1642 }
1643 Ok(())
1644}
1645
1646fn inherits_controller_git_settings(locator: &targets::TargetLocator) -> bool {
1647 !matches!(
1648 locator,
1649 targets::TargetLocator::LocalBare { .. } | targets::TargetLocator::SshBare { .. }
1650 )
1651}
1652
1653fn inherited_git_setting_commands(
1654 locator: &targets::TargetLocator,
1655 session_id: &str,
1656 settings: BTreeMap<String, String>,
1657) -> Result<Vec<CommandSpec>> {
1658 if matches!(locator, targets::TargetLocator::SshBare { .. }) {
1659 return Ok(Vec::new());
1660 }
1661 settings
1662 .into_iter()
1663 .map(|(key, value)| {
1664 targets::command_on_locator(
1665 locator,
1666 session_id,
1667 vec![
1668 "git".into(),
1669 "config".into(),
1670 "--global".into(),
1671 "--replace-all".into(),
1672 "--".into(),
1673 key.clone(),
1674 value,
1675 ],
1676 format!("inherit Git setting {key}"),
1677 )
1678 })
1679 .collect()
1680}
1681
1682fn controller_git_settings() -> Result<BTreeMap<String, String>> {
1683 let output = match Command::new("git")
1684 .args(["config", "--global", "--includes", "--null", "--list"])
1685 .stdin(Stdio::null())
1686 .output()
1687 {
1688 Ok(output) => output,
1689 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(BTreeMap::new()),
1690 Err(error) => return Err(error).context("read controller Git configuration"),
1691 };
1692 if !output.status.success() {
1693 bail!(
1694 "read controller Git configuration failed with status {}: {}",
1695 output.status,
1696 String::from_utf8_lossy(&output.stderr).trim()
1697 );
1698 }
1699 parse_inherited_git_settings(&output.stdout)
1700}
1701
1702fn parse_inherited_git_settings(output: &[u8]) -> Result<BTreeMap<String, String>> {
1703 let mut settings = BTreeMap::new();
1704 for entry in output
1705 .split(|byte| *byte == 0)
1706 .filter(|entry| !entry.is_empty())
1707 {
1708 let entry = std::str::from_utf8(entry).context("decode controller Git configuration")?;
1709 let (key, value) = entry
1710 .split_once('\n')
1711 .with_context(|| format!("controller Git returned malformed entry {entry:?}"))?;
1712 let key = key.to_ascii_lowercase();
1713 if INHERITED_GIT_SETTINGS.contains(&key.as_str()) {
1714 settings.insert(key, value.to_owned());
1715 }
1716 }
1717 Ok(settings)
1718}
1719
1720#[cfg(test)]
1721mod tests;