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 let runtime = mj_core::state::TargetRuntimeSettings::from(selected);
525 if let Some(recorded) = &session.target_runtime {
526 ensure!(
527 recorded == &runtime,
528 "target access settings changed before provisioning; retry with the selected target"
529 );
530 } else {
531 self.state
532 .sessions
533 .get_mut(session_id)
534 .unwrap()
535 .target_runtime = Some(runtime);
536 self.persist_session_state(session_id)?;
537 }
538 let template = self
539 .config
540 .targets
541 .get(&session.target_template_id)
542 .context("target template disappeared during provisioning")?;
543 let profile = self
544 .config
545 .profiles
546 .get(&session.last_profile)
547 .context("harness profile disappeared during provisioning")?;
548 super::worker_binary::preflight_worker_binary(template, executor)?;
549 super::worker_binary::preflight_harness(template, profile, executor)?;
550 self.prepare_managed_raw_worktree(session_id, executor)
551 })();
552 let created_worktree = match preparation {
553 Ok(created) => created,
554 Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
555 return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
556 }
557 Err(error) => return Err(error),
558 };
559 let session = self
560 .state
561 .sessions
562 .get(session_id)
563 .expect("session retained after managed worktree preparation")
564 .clone();
565 let result = (|| {
568 let template = self
569 .config
570 .targets
571 .get(&session.target_template_id)
572 .context("target template disappeared during provisioning")?;
573 if matches!(template, TargetTemplate::AwsEc2 { .. }) {
574 for resource in &session.additional_mounts {
575 ensure!(
576 resource.source.is_dir(),
577 "attached resource source is not a directory: {}",
578 resource.source.display()
579 );
580 }
581 }
582 let mut target = backend_target(
583 template,
584 session.resource_allocation.as_ref(),
585 ContainerOverrides::for_session(&session),
586 )?;
587 let mut runtime_mounts = if matches!(target, targets::TargetTemplate::AwsEc2(_)) {
588 Vec::new()
589 } else {
590 session.additional_mounts.clone()
591 };
592 for notice in enforce_overlay_capable_mounts(&target, &mut runtime_mounts, executor) {
597 executor.notify_notice(¬ice);
598 }
599 let image_user = podman_image_user(&target, executor);
603 let mut bundle = if session.project_directory.is_some() {
604 None
605 } else if let Some(bundle) = self.move_destination_bundle(session_id)? {
606 Some(bundle)
607 } else if failure_disposition == ProvisioningFailureDisposition::Preserve {
608 Some(super::network_git::checkpoint_bundle(&session)?)
609 } else {
610 Some(backend_session_bundle(&session, &self.config, executor)?)
611 };
612 let container_github_token =
613 github_token.filter(|_| configure_github_token_environment(&mut target));
614 if container_github_token.is_some()
615 && let Some(bundle) = bundle.as_mut()
616 {
617 use_github_https_urls(bundle);
618 }
619 preflight_target(template, executor, TargetCheck::Launch)?;
620 let resource_name = crate::database::load_move_operation(session_id)?
621 .filter(|op| {
622 op.workspace_transfer.is_some()
623 && op.phase == mj_core::state::MovePhase::ResumingDestination
624 })
625 .map(|op| targets::move_resource_name(session_id, &op.operation_id))
626 .unwrap_or_else(|| targets::resource_name(session_id))?;
627 let moving_workspace = resource_name != targets::resource_name(session_id)?;
628 let prepared_cache = bundle
629 .as_mut()
630 .filter(|_| !moving_workspace)
631 .and_then(|bundle| {
632 git_cache::prepare(
633 &target,
634 session_id,
635 bundle,
636 &mut runtime_mounts,
637 container_github_token,
638 executor,
639 )
640 });
641 let build_cache = super::mbx::prepare(
644 &target,
645 &session,
646 bundle.as_ref(),
647 prepared_cache.as_ref(),
648 &mut runtime_mounts,
649 executor,
650 );
651 let provision = if let Some(project_directory) = &session.project_directory {
652 targets::provision_bare_project_plan(
653 &target,
654 session_id,
655 &project_directory.to_string_lossy(),
656 )
657 } else {
658 bundle
659 .as_ref()
660 .context("project bundle disappeared during provisioning")
661 .and_then(|bundle| {
662 targets::provision_plan_named(
663 &target,
664 session_id,
665 bundle,
666 &runtime_mounts,
667 image_user,
668 session.container_workspace.as_deref(),
669 &resource_name,
670 )
671 })
672 };
673 let mut provision = match provision {
674 Ok(provision) => provision,
675 Err(error) => {
676 if let Some(cache) = &prepared_cache {
677 let _ = cache.cleanup(executor);
678 }
679 return Err(error);
680 }
681 };
682 if let Some(token) = container_github_token
683 && let Err(error) =
684 provision.provide_target_environment_secret(&target, "GH_TOKEN", token)
685 {
686 if let Some(cache) = &prepared_cache {
687 let _ = cache.cleanup(executor);
688 }
689 return Err(error);
690 }
691
692 let started = Instant::now();
693 let result = provision_target_creation_named(
694 &provision,
695 &target,
696 session_id,
697 executor,
698 &resource_name,
699 |outputs| {
700 super::backend::locator_after_provision_named(
701 template,
702 &target,
703 session_id,
704 outputs.first(),
705 executor,
706 &resource_name,
707 )
708 },
709 )
710 .map(|(locator, remainder)| (locator, remainder, bundle, build_cache));
711 if result.is_err()
712 && let Some(cache) = &prepared_cache
713 {
714 if let Some(locator) =
715 provisioned_locator_named(&target, session_id, None, &resource_name)
716 {
717 let _ = targets::close_plan(&locator, session_id)
718 .and_then(|plan| plan.execute(executor).map(|_| ()));
719 } else {
720 let _ = cache.cleanup(executor);
721 }
722 }
723 tracing::debug!(
724 session_id,
725 elapsed_ms = started.elapsed().as_millis(),
726 "provisioning plan execution completed"
727 );
728 result
729 })();
730 let result = match result {
731 Err(error)
732 if created_worktree
733 && failure_disposition == ProvisioningFailureDisposition::Discard =>
734 {
735 return Err(self.fail_new_session_with_cleanup(session_id, error, executor)?);
736 }
737 Err(error) if failure_disposition == ProvisioningFailureDisposition::Preserve => {
738 Err(error)
739 }
740 Err(error) => {
741 let detail = note_new_session_launch_failure(session_id, &error);
745 {
746 let record = self.state.sessions.get_mut(session_id).unwrap();
747 record.state = SessionState::Error;
748 record.target = None;
749 record.updated_at = super::now();
750 record.last_error = Some(format!("session provisioning failed: {detail}"));
751 }
752 return match self.persist_session_state(session_id) {
753 Ok(()) => Err(error),
754 Err(persistence_error) => Err(error.context(format!(
755 "persist removal of failed provisioning session {session_id}: {persistence_error:#}"
756 ))),
757 };
758 }
759 Ok((locator, remainder, bundle, build_cache)) => {
760 apply_new_session_provisioning_result(&mut self.state, session_id, Ok(locator))?;
761 self.state
762 .sessions
763 .get_mut(session_id)
764 .expect("session retained after provisioning")
765 .build_cache = build_cache;
766 let session = &self.state.sessions[session_id];
767 let backend = backend_locator(
768 session
769 .target
770 .as_ref()
771 .context("provisioned target disappeared")?,
772 session,
773 &self.config,
774 )?;
775 if matches!(backend, targets::TargetLocator::AwsEc2 { .. }) {
776 targets::provision_on_locator_plan(
777 &backend,
778 session_id,
779 bundle
780 .as_ref()
781 .context("AWS provisioning requires a project bundle")?,
782 )
783 } else {
784 Ok(remainder)
785 }
786 }
787 };
788 let result = match result {
789 Err(error) if failure_disposition == ProvisioningFailureDisposition::Discard => {
790 return Err(self.rollback_failed_new_session(session_id, error)?);
791 }
792 result => result,
793 };
794 if result.is_ok()
795 && let Some(session) = self.state.sessions.get(session_id)
796 && let Some(directory) = session
797 .managed_worktree
798 .as_ref()
799 .map(|worktree| worktree.source_project_directory.clone())
800 .or_else(|| session.project_directory.clone())
801 && let Some(template) = self.config.targets.get(&session.target_template_id)
802 {
803 let host = match template {
804 TargetTemplate::LocalBare => Some("local"),
805 TargetTemplate::SshBare { ssh, .. } => Some(ssh.host.as_str()),
806 _ => None,
807 };
808 if let Some(host) = host {
809 self.state.remember_project_directory(host, &directory);
810 crate::database::remember_project_directory(host, &directory)?;
811 }
812 }
813 self.persist_session_state(session_id)?;
814 result
815
816 }).await
817 }
818
819 pub fn mark_worker_connected(
820 &mut self,
821 session_id: &str,
822 native_session_id: Option<String>,
823 ) -> Result<()> {
824 let session = self
825 .state
826 .sessions
827 .get(session_id)
828 .with_context(|| format!("unknown session {session_id}"))?;
829 if session.target.is_none() {
830 bail!("session {session_id} has no provisioned target");
831 }
832 let updated_at = now();
833 crate::database::mark_session_worker_connected(
834 session_id,
835 native_session_id.as_deref(),
836 &updated_at,
837 )?;
838 let session = self
839 .state
840 .sessions
841 .get_mut(session_id)
842 .expect("session disappeared after its worker connection was saved");
843 session.state = SessionState::Running;
844 if native_session_id.is_some() {
845 session.native_session_id = native_session_id;
846 }
847 session.updated_at = updated_at;
848 session.last_error = None;
849 Ok(())
850 }
851
852 fn install_worker_payload(
853 &self,
854 session_id: &str,
855 executor: &impl CommandExecutor,
856 ) -> Result<(targets::TargetLocator, String)> {
857 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
859 let (backend, worker_root) = self.worker_placement(session_id)?;
860 self.prepare_worker_files(session_id, &backend, &worker_root, syncing)?;
861 install_attached_resources(&self.state, session_id, &backend, &worker_root, syncing)?;
862 Ok((backend, worker_root))
863 }
864
865 async fn connect_and_start_worker(
866 &self,
867 session_id: &str,
868 executor: &impl CommandExecutor,
869 backend: &targets::TargetLocator,
870 worker_root: &str,
871 initialize_workspace: bool,
872 ) -> Result<Option<String>> {
873 crate::worker_lifecycle::run(session_id, "connect and start worker", executor, async {
874 let syncing = &StagedExecutor::new(executor, ProvisionStage::Syncing);
875 if initialize_workspace {
876 install_inherited_git_settings(executor, backend, session_id)?;
877 self.initialize_network_workspaces(session_id, backend, syncing)?;
878 }
879 let session = self
880 .state
881 .sessions
882 .get(session_id)
883 .with_context(|| format!("unknown session {session_id}"))?;
884 let profile = self
885 .config
886 .profiles
887 .get(&session.last_profile)
888 .with_context(|| format!("unknown profile {}", session.last_profile))?;
889 let readiness_stage = bridge_readiness_stage(profile);
890 let reconnect = &targets::reconnect_plan(backend, session_id)?.commands[0];
891 let readiness = async {
892 let mut relay = {
893 let _starting = ProvisionStageGuard::new(executor, ProvisionStage::Starting);
894 start_worker_durably(
895 &crate::worker_lifecycle::require(session_id)?,
896 self.state.sessions[session_id]
897 .target
898 .as_ref()
899 .context("worker start has no durable target")?,
900 executor,
901 backend,
902 worker_root,
903 )?;
904 connect_started_worker(reconnect, session_id, executor, backend, worker_root)
905 .await?
906 };
907 let native_session_id =
908 wait_for_native_session_in_stage(&mut relay, executor, readiness_stage).await?;
909 let owner = crate::worker_lifecycle::require(session_id)?;
910 crate::database::finish_worker_restart(session_id, owner.operation_id())?;
911 Ok(Some(native_session_id))
912 }
913 .await;
914 match readiness {
915 Ok(native_session_id) => Ok(native_session_id),
916 Err(error) => {
917 let error = worker_probe_diagnosis(executor, backend, worker_root, error);
920 Err(error)
921 }
922 }
923 })
924 .await
925 }
926}
927
928const MAX_LAUNCH_DIAGNOSTIC_BYTES: usize = 64 * 1024;
929
930const RETAINED_LAUNCH_DIAGNOSTICS: usize = 20;
931
932pub(super) fn note_new_session_launch_failure(session_id: &str, error: &anyhow::Error) -> String {
938 note_new_session_launch_failure_in(&data_dir().join("diagnostics"), session_id, error)
939}
940
941fn note_new_session_launch_failure_in(
942 directory: &Path,
943 session_id: &str,
944 error: &anyhow::Error,
945) -> String {
946 let original = format!("{error:#}");
947 tracing::warn!(session_id, error = %original, "session launch failed");
948 match persist_launch_failure_to(directory, session_id, &original) {
949 Ok(path) => format!("{original}; full diagnostic saved to {}", path.display()),
950 Err(save_error) => {
951 format!("{original}; saving the local diagnostic failed: {save_error:#}")
952 }
953 }
954}
955
956fn persist_launch_failure_to(directory: &Path, session_id: &str, detail: &str) -> Result<PathBuf> {
957 mj_core::config::validate_id("session", session_id)?;
958 std::fs::create_dir_all(directory).with_context(|| {
959 format!(
960 "create launch diagnostics directory {}",
961 directory.display()
962 )
963 })?;
964 #[cfg(unix)]
965 {
966 use std::os::unix::fs::PermissionsExt;
967 std::fs::set_permissions(directory, std::fs::Permissions::from_mode(0o700))?;
968 }
969 let path = directory.join(format!("{session_id}-launch-error.txt"));
970 let detail = bounded_launch_diagnostic(detail);
971 let body = format!(
972 "Mjolnir session launch failure\nsession: {session_id}\nat: {}\n\n{detail}\n",
973 now()
974 );
975 atomic_write(&path, body.as_bytes())?;
976 prune_launch_diagnostics(directory)?;
977 Ok(path)
978}
979
980fn bounded_launch_diagnostic(detail: &str) -> String {
981 if detail.len() <= MAX_LAUNCH_DIAGNOSTIC_BYTES {
982 return detail.to_owned();
983 }
984 let mut head_end = MAX_LAUNCH_DIAGNOSTIC_BYTES / 4;
985 while !detail.is_char_boundary(head_end) {
986 head_end -= 1;
987 }
988 let tail_bytes = MAX_LAUNCH_DIAGNOSTIC_BYTES - head_end;
989 let mut tail_start = detail.len() - tail_bytes;
990 while !detail.is_char_boundary(tail_start) {
991 tail_start += 1;
992 }
993 format!(
994 "{}\n\n[... launch diagnostic truncated ...]\n\n{}",
995 &detail[..head_end],
996 &detail[tail_start..]
997 )
998}
999
1000fn prune_launch_diagnostics(directory: &Path) -> Result<()> {
1001 let mut diagnostics = Vec::new();
1002 for entry in std::fs::read_dir(directory)? {
1003 let entry = entry?;
1004 if !entry
1005 .file_name()
1006 .to_str()
1007 .is_some_and(|name| name.ends_with("-launch-error.txt"))
1008 {
1009 continue;
1010 }
1011 diagnostics.push((entry.metadata()?.modified()?, entry.path()));
1012 }
1013 diagnostics.sort_by_key(|entry| std::cmp::Reverse(entry.0));
1014 for (_, path) in diagnostics.into_iter().skip(RETAINED_LAUNCH_DIAGNOSTICS) {
1015 std::fs::remove_file(&path)
1016 .with_context(|| format!("prune old launch diagnostic {}", path.display()))?;
1017 }
1018 Ok(())
1019}
1020
1021fn apply_new_session_provisioning_result(
1022 state: &mut State,
1023 session_id: &str,
1024 result: Result<TargetLocator>,
1025) -> Result<()> {
1026 match result {
1027 Ok(locator) => {
1028 let record = state.sessions.get_mut(session_id).unwrap();
1029 record.target = Some(locator);
1030 record.state = SessionState::Disconnected;
1033 record.updated_at = now();
1034 record.last_error = None;
1035 Ok(())
1036 }
1037 Err(error) => {
1038 let record = state.sessions.get_mut(session_id).unwrap();
1039 record.state = SessionState::Error;
1040 record.target = None;
1041 record.updated_at = now();
1042 record.last_error = Some(format!("session provisioning failed: {error:#}"));
1043 Err(error)
1044 }
1045 }
1046}
1047
1048fn apply_failed_new_session_launch(state: &mut State, session_id: &str, original_error: &str) {
1051 let record = state.sessions.get_mut(session_id).unwrap();
1052 record.state = SessionState::StartupCleanup;
1053 record.updated_at = now();
1054 record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
1055}
1056
1057pub(super) fn apply_failed_new_session_rollback(
1058 state: &mut State,
1059 session_id: &str,
1060 original_error: &str,
1061 cleanup_error: Option<String>,
1062) -> anyhow::Error {
1063 match cleanup_error {
1064 None => {
1065 let record = state.sessions.get_mut(session_id).unwrap();
1066 record.state = SessionState::Error;
1067 record.target = None;
1068 record.updated_at = now();
1069 record.last_error = Some(format!("worker bootstrap failed: {original_error}"));
1070 anyhow::anyhow!("{original_error}; partial target removed and failed session retained")
1071 }
1072 Some(cleanup_error) => {
1073 let failure = format!(
1074 "{original_error}; cleanup of the failed session target failed: {cleanup_error}"
1075 );
1076 let record = state.sessions.get_mut(session_id).unwrap();
1077 record.state = SessionState::StartupCleanup;
1078 record.updated_at = now();
1079 record.last_error = Some(format!("worker bootstrap failed: {failure}"));
1080 anyhow::anyhow!(failure)
1081 }
1082 }
1083}
1084
1085pub(super) fn install_attached_resources(
1086 state: &State,
1087 session_id: &str,
1088 backend: &targets::TargetLocator,
1089 worker_root: &str,
1090 executor: &impl CommandExecutor,
1091) -> Result<()> {
1092 let targets::TargetLocator::AwsEc2 { .. } = backend else {
1093 return Ok(());
1094 };
1095 let session = state
1096 .sessions
1097 .get(session_id)
1098 .with_context(|| format!("unknown session {session_id}"))?;
1099 if session.additional_mounts.is_empty() {
1100 return Ok(());
1101 }
1102 for resource in &session.additional_mounts {
1103 let install = targets::command_on_locator(
1104 backend,
1105 session_id,
1106 vec![
1107 format!("{worker_root}/hel"),
1108 "worker".into(),
1109 "install-resource".into(),
1110 "--destination".into(),
1111 resource.destination.to_string_lossy().into_owned(),
1112 ],
1113 "stream attached resource",
1114 )?;
1115 mj_checkpoint::resources::stream_resource(&resource.source, |stream| {
1116 execute_checked_with_stdin(executor, &install, stream).map(|_| ())
1117 })
1118 .with_context(|| format!("stream attached resource {}", resource.source.display()))?;
1119 }
1120 Ok(())
1121}
1122
1123pub(super) fn execute_concurrent_lanes<A: Send, B: Send>(
1127 first: impl FnOnce() -> Result<A> + Send,
1128 second: impl FnOnce() -> Result<B> + Send,
1129) -> Result<(A, B)> {
1130 std::thread::scope(|scope| {
1131 let owner = crate::worker_lifecycle::capture();
1132 let second = scope.spawn(move || match owner {
1133 Some(owner) => owner.scope_blocking(second),
1134 None => second(),
1135 });
1136 let first = first();
1137 let second = second.join().unwrap_or_else(|panic| {
1138 Err(anyhow::anyhow!(
1139 "concurrent target lane panicked: {}",
1140 targets::command_thread_panic_message(panic.as_ref())
1141 ))
1142 });
1143 match (first, second) {
1144 (Err(error), _) => Err(error),
1145 (Ok(_), Err(error)) => Err(error),
1146 (Ok(first), Ok(second)) => Ok((first, second)),
1147 }
1148 })
1149}
1150
1151fn execute_repository_setup(
1152 plan: &targets::CommandPlan,
1153 executor: &(impl CommandExecutor + Sync),
1154) -> Result<()> {
1155 if plan.commands.is_empty() {
1156 return Ok(());
1157 }
1158 let _cloning = ProvisionStageGuard::new(executor, ProvisionStage::Cloning);
1159 plan.execute_concurrent(executor).map(|_| ())
1160}
1161
1162#[cfg(test)]
1169fn provision_target(
1170 plan: &targets::CommandPlan,
1171 target: &targets::TargetTemplate,
1172 session_id: &str,
1173 executor: &(impl CommandExecutor + Sync),
1174 discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1175) -> Result<TargetLocator> {
1176 let Some((creation, remainder)) = plan.split_at_target_creation() else {
1177 return discover(&plan.execute_concurrent(executor)?);
1179 };
1180 let mut outputs = creation.execute_concurrent(executor)?;
1181 let result = match remainder.execute_concurrent(executor) {
1182 Ok(rest) => {
1183 outputs.extend(rest);
1184 discover(&outputs)
1185 }
1186 Err(error) => Err(error),
1187 };
1188 result.map_err(|error| {
1189 match cleanup_failed_provision(target, session_id, outputs.first(), executor) {
1190 Some(note) => error.context(note),
1191 None => error,
1192 }
1193 })
1194}
1195
1196#[cfg(test)]
1200fn provision_target_creation(
1201 plan: &targets::CommandPlan,
1202 target: &targets::TargetTemplate,
1203 session_id: &str,
1204 executor: &(impl CommandExecutor + Sync),
1205 discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1206) -> Result<(TargetLocator, targets::CommandPlan)> {
1207 provision_target_creation_named(
1208 plan,
1209 target,
1210 session_id,
1211 executor,
1212 &targets::resource_name(session_id)?,
1213 discover,
1214 )
1215}
1216
1217fn provision_target_creation_named(
1218 plan: &targets::CommandPlan,
1219 target: &targets::TargetTemplate,
1220 session_id: &str,
1221 executor: &(impl CommandExecutor + Sync),
1222 name: &str,
1223 discover: impl FnOnce(&[CommandOutput]) -> Result<TargetLocator>,
1224) -> Result<(TargetLocator, targets::CommandPlan)> {
1225 let Some((creation, remainder)) = plan.split_at_target_creation() else {
1226 let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1229 plan.execute_concurrent(executor)
1230 })?;
1231 return discover(&outputs).map(|locator| {
1232 (
1233 locator,
1234 targets::CommandPlan {
1235 description: plan.description.clone(),
1236 commands: Vec::new(),
1237 },
1238 )
1239 });
1240 };
1241 let outputs = crate::image_pull_gate::with_image_ready(target, executor, || {
1245 creation.execute_concurrent(executor)
1246 })?;
1247 discover(&outputs)
1248 .map(|locator| (locator, remainder))
1249 .map_err(|error| {
1250 match cleanup_failed_provision_named(
1251 target,
1252 session_id,
1253 outputs.first(),
1254 executor,
1255 name,
1256 ) {
1257 Some(note) => error.context(note),
1258 None => error,
1259 }
1260 })
1261}
1262
1263#[cfg(test)]
1270fn cleanup_failed_provision(
1271 target: &targets::TargetTemplate,
1272 session_id: &str,
1273 create_output: Option<&CommandOutput>,
1274 executor: &impl CommandExecutor,
1275) -> Option<String> {
1276 cleanup_failed_provision_named(
1277 target,
1278 session_id,
1279 create_output,
1280 executor,
1281 &targets::resource_name(session_id).ok()?,
1282 )
1283}
1284
1285fn cleanup_failed_provision_named(
1286 target: &targets::TargetTemplate,
1287 session_id: &str,
1288 create_output: Option<&CommandOutput>,
1289 executor: &impl CommandExecutor,
1290 name: &str,
1291) -> Option<String> {
1292 let locator = provisioned_locator_named(target, session_id, create_output, name)?;
1293
1294 let leak = format!(
1295 "the resource may still exist; find it via its dev.mj.session={session_id} label/tag"
1296 );
1297 let cleanup = if name == targets::resource_name(session_id).ok()? {
1298 targets::close_plan(&locator, session_id)
1299 } else {
1300 targets::retire_move_target_plan(&locator, session_id)
1301 };
1302 let plan = match cleanup {
1303 Ok(plan) => plan,
1304 Err(error) => {
1305 tracing::warn!(
1306 session_id,
1307 error = format!("{error:#}"),
1308 "could not build provisioning cleanup plan"
1309 );
1310 return Some(format!("cleanup FAILED: {error:#}; {leak}"));
1311 }
1312 };
1313 let purpose = plan
1314 .commands
1315 .iter()
1316 .map(|command| command.purpose.clone())
1317 .collect::<Vec<_>>()
1318 .join("; ");
1319 let Err(error) = plan.execute(executor) else {
1320 return Some(format!("cleanup succeeded: {purpose}"));
1321 };
1322 match targets::cleanup_target_is_confirmed_absent(&locator, session_id, executor) {
1323 Ok(true) => Some(format!("cleanup succeeded: {purpose}")),
1324 Ok(false) => {
1325 tracing::warn!(
1326 session_id,
1327 error = format!("{error:#}"),
1328 "provisioning cleanup failed and the target may still exist"
1329 );
1330 Some(format!("cleanup FAILED ({purpose}): {error:#}; {leak}"))
1331 }
1332 Err(confirm_error) => {
1333 tracing::warn!(
1334 session_id,
1335 error = format!("{confirm_error:#}"),
1336 "could not confirm whether the failed provisioning target was removed"
1337 );
1338 Some(format!(
1339 "cleanup FAILED ({purpose}): {error:#}; checking whether it was removed also failed: {confirm_error:#}; {leak}"
1340 ))
1341 }
1342 }
1343}
1344
1345fn provisioned_locator_named(
1350 target: &targets::TargetTemplate,
1351 session_id: &str,
1352 create_output: Option<&CommandOutput>,
1353 name: &str,
1354) -> Option<targets::TargetLocator> {
1355 let container_id = || Some(name.to_owned());
1356
1357 Some(match target {
1358 targets::TargetTemplate::LocalBare => return None,
1361 targets::TargetTemplate::LocalPodman(container) => targets::TargetLocator::LocalPodman {
1362 borrowed_from: None,
1363 container_id: container_id()?,
1364 workspace_storage: targets::podman_workspace_locator_named(container, name).ok()?,
1365 },
1366 targets::TargetTemplate::LocalDocker(_) => targets::TargetLocator::LocalDocker {
1367 borrowed_from: None,
1368 container_id: container_id()?,
1369 },
1370 targets::TargetTemplate::AppleContainer(_) => targets::TargetLocator::AppleContainer {
1371 borrowed_from: None,
1372 container_id: container_id()?,
1373 },
1374 targets::TargetTemplate::SshPodman { ssh, container } => {
1375 targets::TargetLocator::SshPodman {
1376 borrowed_from: None,
1377 ssh: ssh.clone(),
1378 container_id: container_id()?,
1379 workspace_storage: targets::podman_workspace_locator_named(container, name).ok()?,
1380 }
1381 }
1382 targets::TargetTemplate::SshDocker { ssh, .. } => targets::TargetLocator::SshDocker {
1383 borrowed_from: None,
1384 ssh: ssh.clone(),
1385 container_id: container_id()?,
1386 },
1387 targets::TargetTemplate::SshBare { ssh, .. } => targets::TargetLocator::SshBare {
1388 ssh: ssh.clone(),
1389 workspace: targets::workspace_for(target, session_id).ok()?,
1390 worker_id: None,
1391 },
1392 targets::TargetTemplate::AwsEc2(aws) => targets::TargetLocator::AwsEc2 {
1393 profile: aws.profile.clone(),
1394 region: aws.region.clone(),
1395 instance_id: serde_json::from_slice::<serde_json::Value>(&create_output?.stdout)
1396 .ok()?
1397 .pointer("/Instances/0/InstanceId")?
1398 .as_str()?
1399 .to_owned(),
1400 ssh: aws.ssh.clone(),
1401 workspace: targets::workspace_for(target, session_id).ok()?,
1402 },
1403 })
1404}
1405
1406pub(super) fn enforce_overlay_capable_mounts(
1417 target: &targets::TargetTemplate,
1418 mounts: &mut [targets::AdditionalMount],
1419 executor: &impl CommandExecutor,
1420) -> Vec<String> {
1421 let ssh = match target {
1422 targets::TargetTemplate::LocalPodman(_) | targets::TargetTemplate::LocalDocker(_) => None,
1423 targets::TargetTemplate::SshPodman { ssh, .. }
1424 | targets::TargetTemplate::SshDocker { ssh, .. } => Some(ssh),
1425 _ => return Vec::new(),
1426 };
1427 let overlaid = mounts
1428 .iter()
1429 .filter(|mount| mount.access == targets::MountAccess::Cow)
1430 .map(|mount| mount.source.clone())
1431 .collect::<Vec<_>>();
1432 if overlaid.is_empty() {
1433 return Vec::new();
1434 }
1435 if matches!(target, targets::TargetTemplate::LocalDocker(_)) {
1437 match targets::local_docker_vm_share(executor) {
1438 Ok(Some(reason)) => {
1439 return mounts
1440 .iter_mut()
1441 .filter(|mount| mount.access == targets::MountAccess::Cow)
1442 .map(|mount| {
1443 mount.access = mount.access.without_overlay();
1444 format!(
1445 "Mounted {} read-only: {reason}, which cannot back the \
1446 copy-on-write overlay.",
1447 mount.source.display()
1448 )
1449 })
1450 .collect();
1451 }
1452 Ok(None) => {}
1453 Err(error) => tracing::warn!(
1454 error = format!("{error:#}"),
1455 "could not identify the Docker daemon platform; probing this host's filesystems"
1456 ),
1457 }
1458 }
1459 let filesystems = match targets::probe_filesystem_types(ssh, &overlaid, executor) {
1460 Ok(filesystems) => filesystems,
1461 Err(error) => {
1462 tracing::warn!(
1463 error = format!("{error:#}"),
1464 "could not probe attached-directory filesystems; preserving overlay mounts"
1465 );
1466 return vec![format!(
1467 "Could not read the filesystem under the attached directories, so they keep the \
1468 copy-on-write overlay: {error:#}"
1469 )];
1470 }
1471 };
1472 let mut notices = Vec::new();
1473 for (mount, filesystem) in mounts
1474 .iter_mut()
1475 .filter(|mount| mount.access == targets::MountAccess::Cow)
1476 .zip(filesystems)
1477 {
1478 let Some(reason) = targets::overlay_unsupported_filesystem(&filesystem) else {
1479 continue;
1480 };
1481 mount.access = mount.access.without_overlay();
1482 notices.push(format!(
1483 "Mounted {} read-only: the overlay is unreliable on {filesystem} ({reason}).",
1484 mount.source.display()
1485 ));
1486 }
1487 notices
1488}
1489
1490static IMAGE_USERS: std::sync::LazyLock<std::sync::Mutex<BTreeMap<String, targets::ImageUser>>> =
1495 std::sync::LazyLock::new(std::sync::Mutex::default);
1496
1497pub(super) fn podman_image_user(
1504 target: &targets::TargetTemplate,
1505 executor: &impl CommandExecutor,
1506) -> Option<targets::ImageUser> {
1507 let (ssh, container) = match target {
1508 targets::TargetTemplate::LocalPodman(container) => (None, container),
1509 targets::TargetTemplate::SshPodman { ssh, container } => (Some(ssh), container),
1510 _ => return None,
1511 };
1512 let image = container.image.as_str();
1513 let key = format!(
1514 "{}|{image}",
1515 ssh.map_or("local", |ssh| ssh.destination.as_str())
1516 );
1517 if let Some(cached) = IMAGE_USERS.lock().expect("image user cache").get(&key) {
1518 return Some(*cached);
1519 }
1520 match crate::image_pull_gate::with_image_ready(target, executor, || {
1524 targets::probe_image_user(ssh, container, executor)
1525 }) {
1526 Ok(user) => {
1527 IMAGE_USERS
1528 .lock()
1529 .expect("image user cache")
1530 .insert(key, user);
1531 Some(user)
1532 }
1533 Err(error) => {
1534 tracing::warn!(
1535 image,
1536 error = format!("{error:#}"),
1537 "could not read the container image user; keeping Podman's default user mapping"
1538 );
1539 executor.notify_notice(&format!(
1540 "Could not read the user of image {image}, so the container runs with Podman's \
1541 default user mapping and may not be able to write to an attached directory: \
1542 {error:#}"
1543 ));
1544 None
1545 }
1546 }
1547}
1548
1549pub(super) struct StagedExecutor<'a, E: CommandExecutor> {
1553 inner: &'a E,
1554 stage: ProvisionStage,
1555 _guard: ProvisionStageGuard<'a, E>,
1556}
1557
1558impl<'a, E: CommandExecutor> StagedExecutor<'a, E> {
1559 pub(crate) fn new(inner: &'a E, stage: ProvisionStage) -> Self {
1560 Self {
1561 inner,
1562 stage,
1563 _guard: ProvisionStageGuard::new(inner, stage),
1564 }
1565 }
1566
1567 fn staged(&self, command: &CommandSpec) -> CommandSpec {
1568 if command.stage.is_some() {
1569 return command.clone();
1570 }
1571 command.clone().stage(self.stage)
1572 }
1573}
1574
1575impl<E: CommandExecutor> CommandExecutor for StagedExecutor<'_, E> {
1576 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1577 self.inner.execute(&self.staged(command))
1578 }
1579
1580 fn cancellation_requested(&self) -> bool {
1581 self.inner.cancellation_requested()
1582 }
1583
1584 fn stage_started(&self, stage: ProvisionStage) {
1585 self.inner.stage_started(stage);
1586 }
1587
1588 fn stage_finished(&self, stage: ProvisionStage) {
1589 self.inner.stage_finished(stage);
1590 }
1591
1592 fn notify_notice(&self, notice: &str) {
1593 self.inner.notify_notice(notice);
1594 }
1595
1596 fn execute_with_stdin(
1597 &self,
1598 command: &CommandSpec,
1599 input: &mut (dyn std::io::Read + Send),
1600 ) -> Result<CommandOutput> {
1601 self.inner.execute_with_stdin(&self.staged(command), input)
1602 }
1603}
1604
1605fn execute_checked_with_stdin(
1606 executor: &impl CommandExecutor,
1607 command: &CommandSpec,
1608 input: &mut (dyn std::io::Read + Send),
1609) -> Result<CommandOutput> {
1610 let output = executor.execute_with_stdin(command, input)?;
1611 if output.status != 0 {
1612 bail!(
1613 "{} failed with status {}: {}",
1614 command.purpose,
1615 output.status,
1616 String::from_utf8_lossy(&output.stderr)
1617 );
1618 }
1619 Ok(output)
1620}
1621
1622pub(super) fn install_inherited_git_settings(
1623 executor: &impl CommandExecutor,
1624 locator: &targets::TargetLocator,
1625 session_id: &str,
1626) -> Result<()> {
1627 let settings = if inherits_controller_git_settings(locator) {
1628 controller_git_settings()?
1629 } else {
1630 BTreeMap::new()
1631 };
1632 for command in inherited_git_setting_commands(locator, session_id, settings)? {
1633 execute_checked(executor, command)?;
1634 }
1635 Ok(())
1636}
1637
1638fn inherits_controller_git_settings(locator: &targets::TargetLocator) -> bool {
1639 !matches!(
1640 locator,
1641 targets::TargetLocator::LocalBare { .. } | targets::TargetLocator::SshBare { .. }
1642 )
1643}
1644
1645fn inherited_git_setting_commands(
1646 locator: &targets::TargetLocator,
1647 session_id: &str,
1648 settings: BTreeMap<String, String>,
1649) -> Result<Vec<CommandSpec>> {
1650 if matches!(locator, targets::TargetLocator::SshBare { .. }) {
1651 return Ok(Vec::new());
1652 }
1653 settings
1654 .into_iter()
1655 .map(|(key, value)| {
1656 targets::command_on_locator(
1657 locator,
1658 session_id,
1659 vec![
1660 "git".into(),
1661 "config".into(),
1662 "--global".into(),
1663 "--replace-all".into(),
1664 "--".into(),
1665 key.clone(),
1666 value,
1667 ],
1668 format!("inherit Git setting {key}"),
1669 )
1670 })
1671 .collect()
1672}
1673
1674fn controller_git_settings() -> Result<BTreeMap<String, String>> {
1675 let output = match Command::new("git")
1676 .args(["config", "--global", "--includes", "--null", "--list"])
1677 .stdin(Stdio::null())
1678 .output()
1679 {
1680 Ok(output) => output,
1681 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(BTreeMap::new()),
1682 Err(error) => return Err(error).context("read controller Git configuration"),
1683 };
1684 if !output.status.success() {
1685 bail!(
1686 "read controller Git configuration failed with status {}: {}",
1687 output.status,
1688 String::from_utf8_lossy(&output.stderr).trim()
1689 );
1690 }
1691 parse_inherited_git_settings(&output.stdout)
1692}
1693
1694fn parse_inherited_git_settings(output: &[u8]) -> Result<BTreeMap<String, String>> {
1695 let mut settings = BTreeMap::new();
1696 for entry in output
1697 .split(|byte| *byte == 0)
1698 .filter(|entry| !entry.is_empty())
1699 {
1700 let entry = std::str::from_utf8(entry).context("decode controller Git configuration")?;
1701 let (key, value) = entry
1702 .split_once('\n')
1703 .with_context(|| format!("controller Git returned malformed entry {entry:?}"))?;
1704 let key = key.to_ascii_lowercase();
1705 if INHERITED_GIT_SETTINGS.contains(&key.as_str()) {
1706 settings.insert(key, value.to_owned());
1707 }
1708 }
1709 Ok(settings)
1710}
1711
1712#[cfg(test)]
1713mod tests;