1use std::path::{Path, PathBuf};
4
5use anyhow::{Context, Result, bail};
6
7use crate::session_manager::StandaloneSession;
8use mj_core::config::{AwsAddressSource, SshConnection, TargetTemplate};
9use mj_core::state::{
10 PodmanWorkspaceLocator, SessionRecord, SessionState, TargetLocator, normalize_session_title,
11};
12
13use crate::targets::{self, CommandExecutor, CommandOutput, CommandSpec, SshTarget};
14use mj_core::worker_launch::WorkerOwnership;
15
16use super::backend::{ContainerOverrides, backend_locator, backend_target};
17use super::readiness::wait_for_native_session;
18use super::{Controller, now};
19
20pub use mj_core::state::{RecoveryCandidate, RecoveryScan};
21
22impl Controller {
23 pub fn scan_orphan_workers(
31 &self,
32 executor: &impl CommandExecutor,
33 all_instances: bool,
34 ) -> RecoveryScan {
35 let mut scan = RecoveryScan {
36 instance_id: mj_core::config::instance_identity(),
37 ..RecoveryScan::default()
38 };
39 let moves = match crate::database::load_move_operations() {
40 Ok(moves) => moves,
41 Err(error) => {
42 scan.warnings
43 .push(format!("cannot read Move resource ownership: {error:#}"));
44 return scan;
45 }
46 };
47 for (target_id, template) in &self.config.targets {
48 match scan_target_workers(target_id, template, executor) {
49 Ok(candidates) => {
50 for mut candidate in candidates {
51 let move_owned = moves.iter().any(|op| {
52 op.selection.session_id == candidate.session_id
53 && op.prepared_destination.as_ref().is_some_and(|d| {
54 d.owns_resource()
55 && match &candidate.locator {
56 mj_core::state::TargetLocator::AwsEc2 {
57 instance_id,
58 ..
59 } => d
60 .instance_id()
61 .is_none_or(|owned| owned == instance_id),
62 _ => false,
63 }
64 })
65 });
66 if move_owned || self.state_represents(&candidate) {
67 continue;
68 }
69 candidate.tracked_session = self
70 .state
71 .sessions
72 .get(&candidate.session_id)
73 .map(|record| record.state);
74 scan.candidates.push(candidate);
75 }
76 }
77 Err(error) => scan.warnings.push(format!("target {target_id}: {error:#}")),
78 }
79 }
80 scan.candidates.sort_by(|left, right| {
81 (&left.session_id, &left.target_template_id)
82 .cmp(&(&right.session_id, &right.target_template_id))
83 });
84 scan.candidates.dedup_by(|left, right| {
85 left.session_id == right.session_id
86 && left.target_template_id == right.target_template_id
87 });
88 if !all_instances {
89 restrict_to_instance(&mut scan);
90 }
91 scan
92 }
93
94 fn state_represents(&self, candidate: &RecoveryCandidate) -> bool {
108 if self
109 .state
110 .sessions
111 .get(&candidate.session_id)
112 .is_some_and(|record| locator_in_flight(record.state))
113 {
114 return true;
115 }
116 self.state.sessions.values().any(|record| {
117 record
118 .target
119 .as_ref()
120 .is_some_and(|locator| same_resource(locator, &candidate.locator))
121 })
122 }
123
124 pub async fn adopt_orphan_worker(
125 &mut self,
126 session_id: &str,
127 target_id: &str,
128 profile_override: Option<&str>,
129 bundle_override: Option<&str>,
130 all_instances: bool,
131 executor: &impl CommandExecutor,
132 ) -> Result<()> {
133 let (record, newly_adopted) = match self.state.sessions.get(session_id).cloned() {
134 Some(existing) if adoption_unfinished(&existing, target_id) => {
138 for (flag, requested, adopted) in [
139 ("profile", profile_override, existing.last_profile.as_str()),
140 ("bundle", bundle_override, existing.bundle_id.as_str()),
141 ] {
142 if let Some(requested) = requested
143 && requested != adopted
144 {
145 bail!(
146 "session {session_id} was already adopted with {flag} {adopted:?}; retry without --{flag}"
147 );
148 }
149 }
150 (existing, false)
151 }
152 Some(existing) => bail!(
156 "session {session_id} is already tracked in state {}; use `recover destroy` to remove a resource it left behind",
157 existing.state.as_str()
158 ),
159 None => {
160 let scan = self.scan_orphan_workers(executor, true);
161 let candidate = scan
162 .candidates
163 .into_iter()
164 .find(|candidate| {
165 candidate.session_id == session_id
166 && candidate.target_template_id == target_id
167 })
168 .with_context(|| {
169 format!("no managed orphan {session_id} was found on target {target_id:?}")
170 })?;
171 require_instance_access(&candidate, &scan.instance_id, all_instances)?;
172 let profile_id = profile_override
173 .map(str::to_owned)
174 .or_else(|| {
175 candidate
176 .ownership
177 .as_ref()
178 .map(|marker| marker.profile_id.clone())
179 })
180 .context("orphan has no ownership marker; pass --profile")?;
181 let bundle_id = bundle_override
182 .map(str::to_owned)
183 .or_else(|| {
184 candidate
185 .ownership
186 .as_ref()
187 .map(|marker| marker.bundle_id.clone())
188 })
189 .context("orphan has no ownership marker; pass --bundle")?;
190 let profile = self
191 .config
192 .profiles
193 .get(&profile_id)
194 .with_context(|| format!("unknown profile {profile_id:?}"))?;
195 let (bundle_id, project) =
196 if let Some(saved) = crate::database::saved_project(&bundle_id)? {
197 saved
198 } else {
199 let bundle = self
200 .config
201 .bundles
202 .get(&bundle_id)
203 .with_context(|| format!("unknown project {bundle_id:?}"))?;
204 (
205 bundle_id,
206 crate::project_catalog::snapshot(bundle, executor, false)?,
207 )
208 };
209 let workspace_id = resolve_recovery_workspace_id(
210 candidate
211 .ownership
212 .as_ref()
213 .map(|ownership| ownership.workspace_id.as_str())
214 .unwrap_or(mj_core::workspace::DEFAULT_WORKSPACE_ID),
215 )?;
216 let container_workspace = self
217 .config
218 .targets
219 .get(target_id)
220 .and_then(|template| {
221 recovery_backend_locator(template, &candidate.locator, session_id).ok()
222 })
223 .and_then(|backend| {
224 adopted_container_workspace(&backend, session_id, executor)
225 });
226 let mut record = adopted_session_record(
227 session_id,
228 target_id,
229 profile_id,
230 profile.kind,
231 bundle_id,
232 workspace_id,
233 candidate.locator,
234 );
235 record.target_runtime =
236 Some(record.target_runtime_settings(&self.config)?.into_owned());
237 record.container_workspace = container_workspace;
238 record.project = Some(project);
239 (record, true)
240 }
241 };
242 let locator = record
243 .target
244 .as_ref()
245 .context("adopted session has no target locator")?;
246 let backend = backend_locator(locator, &record, &self.config)?;
247 let spec = targets::reconnect_plan(&backend, session_id)?
248 .commands
249 .into_iter()
250 .next()
251 .context("reconnect plan is empty")?;
252 if newly_adopted {
253 crate::database::save_session(&record)?;
258 self.state.sessions.insert(session_id.to_owned(), record);
259 }
260 match self.complete_adoption(session_id, &spec, executor).await {
261 Ok(()) => Ok(()),
262 Err(error) => Err(self.record_adoption_failure(session_id, error)),
266 }
267 }
268
269 async fn complete_adoption(
271 &mut self,
272 session_id: &str,
273 spec: &CommandSpec,
274 executor: &impl CommandExecutor,
275 ) -> Result<()> {
276 let mut relay = StandaloneSession::connect_command(spec, session_id)
277 .await
278 .context("orphan relay did not complete the v1 handshake")?;
279 let harness = self
280 .state
281 .sessions
282 .get(session_id)
283 .with_context(|| format!("unknown adopted session {session_id}"))?
284 .harness_kind;
285 let native_session_id = wait_for_native_session(&mut relay, executor, harness).await?;
286 self.mark_worker_connected(session_id, Some(native_session_id))?;
287 if let Some(title) = relay
288 .snapshot()
289 .materialized
290 .session_title
291 .as_deref()
292 .and_then(normalize_session_title)
293 {
294 crate::database::set_session_acp_title(session_id, Some(&title))?;
295 self.state
296 .sessions
297 .get_mut(session_id)
298 .expect("adopted session disappeared while saving its ACP title")
299 .acp_session_title = Some(title);
300 }
301 Ok(())
302 }
303
304 fn record_adoption_failure(&mut self, session_id: &str, error: anyhow::Error) -> anyhow::Error {
309 let Some(record) = self.state.sessions.get_mut(session_id) else {
310 return error;
311 };
312 record.updated_at = now();
313 record.last_error = Some(format!("orphan adoption failed: {error:#}"));
314 match self.persist_session_state(session_id) {
315 Ok(()) => error,
316 Err(persist_error) => error.context(format!(
317 "recorded the adoption failure in memory, but failed to persist it: {persist_error:#}"
318 )),
319 }
320 }
321
322 pub fn destroy_orphan_worker(
323 &self,
324 session_id: &str,
325 target_id: &str,
326 confirmation: &str,
327 all_instances: bool,
328 executor: &impl CommandExecutor,
329 ) -> Result<()> {
330 if confirmation != session_id {
331 bail!("refusing destructive recovery: --confirm must exactly match the session ID");
332 }
333 let scan = self.scan_orphan_workers(executor, true);
334 let candidate = scan
335 .candidates
336 .into_iter()
337 .find(|candidate| {
338 candidate.session_id == session_id && candidate.target_template_id == target_id
339 })
340 .with_context(|| {
341 format!("no managed orphan {session_id} was found on target {target_id:?}")
342 })?;
343 require_instance_access(&candidate, &scan.instance_id, all_instances)?;
344 let template = self.config.targets.get(target_id).unwrap();
345 let backend = recovery_backend_locator(template, &candidate.locator, session_id)?;
346 targets::close_plan(&backend, session_id)?
347 .execute(executor)
348 .map(|_| ())
349 }
350}
351
352const fn locator_in_flight(state: SessionState) -> bool {
355 matches!(
356 state,
357 SessionState::Provisioning
358 | SessionState::Checkpointing
359 | SessionState::Closing
360 | SessionState::Destroying
361 )
362}
363
364fn same_resource(left: &TargetLocator, right: &TargetLocator) -> bool {
368 match (left, right) {
369 (
370 TargetLocator::LocalBare { worker_root: left },
371 TargetLocator::LocalBare { worker_root: right },
372 ) => left == right,
373 (
374 TargetLocator::LocalPodman {
375 container_id: left, ..
376 },
377 TargetLocator::LocalPodman {
378 container_id: right,
379 ..
380 },
381 )
382 | (
383 TargetLocator::LocalDocker {
384 container_id: left, ..
385 },
386 TargetLocator::LocalDocker {
387 container_id: right,
388 ..
389 },
390 )
391 | (
392 TargetLocator::AppleContainer {
393 container_id: left, ..
394 },
395 TargetLocator::AppleContainer {
396 container_id: right,
397 ..
398 },
399 ) => left == right,
400 (
401 TargetLocator::AwsEc2 {
402 instance_id: left, ..
403 },
404 TargetLocator::AwsEc2 {
405 instance_id: right, ..
406 },
407 ) => left == right,
408 (
409 TargetLocator::SshBare {
410 host: left_host,
411 workspace: left_workspace,
412 ..
413 },
414 TargetLocator::SshBare {
415 host: right_host,
416 workspace: right_workspace,
417 ..
418 },
419 ) => left_host == right_host && left_workspace == right_workspace,
420 (
421 TargetLocator::SshPodman {
422 host: left_host,
423 container_id: left_container,
424 ..
425 },
426 TargetLocator::SshPodman {
427 host: right_host,
428 container_id: right_container,
429 ..
430 },
431 )
432 | (
433 TargetLocator::SshDocker {
434 host: left_host,
435 container_id: left_container,
436 ..
437 },
438 TargetLocator::SshDocker {
439 host: right_host,
440 container_id: right_container,
441 ..
442 },
443 ) => left_host == right_host && left_container == right_container,
444 _ => false,
445 }
446}
447
448fn restrict_to_instance(scan: &mut RecoveryScan) {
451 let before = scan.candidates.len();
452 scan.candidates
453 .retain(|candidate| candidate.instance_id.as_deref() == Some(scan.instance_id.as_str()));
454 scan.hidden_other_instances = before - scan.candidates.len();
455}
456
457fn require_instance_access(
460 candidate: &RecoveryCandidate,
461 scan_instance: &str,
462 all_instances: bool,
463) -> Result<()> {
464 if all_instances {
465 return Ok(());
466 }
467 match candidate.instance_id.as_deref() {
468 Some(instance) if instance == scan_instance => Ok(()),
469 Some(other) => bail!(
470 "worker {} belongs to instance {other:?}, not this instance {scan_instance:?}; pass --all-instances to act on it",
471 candidate.session_id
472 ),
473 None => bail!(
474 "worker {} has no instance stamp (created by an older build); pass --all-instances to act on it",
475 candidate.session_id
476 ),
477 }
478}
479
480fn resolve_recovery_workspace_id(marked_workspace_id: &str) -> Result<String> {
484 if marked_workspace_id == mj_core::workspace::DEFAULT_WORKSPACE_ID {
485 return Ok(marked_workspace_id.to_owned());
486 }
487 if crate::database::list_workspaces()?
488 .into_iter()
489 .any(|workspace| workspace.id == marked_workspace_id)
490 {
491 return Ok(marked_workspace_id.to_owned());
492 }
493 Ok(crate::database::create_or_get_workspace("Recovered")?.id)
494}
495
496fn adopted_container_workspace(
501 backend: &targets::TargetLocator,
502 session_id: &str,
503 executor: &impl CommandExecutor,
504) -> Option<PathBuf> {
505 if !matches!(
506 backend,
507 targets::TargetLocator::LocalPodman { .. }
508 | targets::TargetLocator::LocalDocker { .. }
509 | targets::TargetLocator::AppleContainer { .. }
510 | targets::TargetLocator::SshPodman { .. }
511 | targets::TargetLocator::SshDocker { .. }
512 ) {
513 return None;
514 }
515 let workspace = targets::new_container_workspace(session_id).ok()?;
516 let command = targets::command_on_locator(
517 backend,
518 session_id,
519 vec![
520 "test".to_owned(),
521 "-d".to_owned(),
522 workspace.to_string_lossy().into_owned(),
523 ],
524 "probe the adopted session workspace",
525 )
526 .ok()?;
527 match executor.execute(&command) {
528 Ok(output) if output.status == 0 => Some(workspace),
529 Ok(_) => None,
530 Err(error) => {
531 tracing::debug!(
532 session_id,
533 %error,
534 "could not probe the adopted session workspace; assuming the shared one"
535 );
536 None
537 }
538 }
539}
540
541fn adopted_session_record(
543 session_id: &str,
544 target_id: &str,
545 profile_id: String,
546 harness_kind: mj_core::config::HarnessKind,
547 bundle_id: String,
548 workspace_id: String,
549 locator: TargetLocator,
550) -> SessionRecord {
551 let now = now();
552 SessionRecord {
553 project: None,
554 target_runtime: None,
555 launch_base: None,
556 launch_branch: None,
557 checkout: None,
558 publication: None,
559 build_cache: None,
560 subagents: None,
561 container_workspace: None,
563 create_managed_worktree: None,
564 workspace_id,
565 archived: false,
566 container_cpus: None,
567 container_memory: None,
568 id: session_id.to_owned(),
569 title: format!("Recovered {}", &session_id[..session_id.len().min(8)]),
570 harness_kind,
571 last_profile: profile_id,
572 bundle_id,
573 project_directory: None,
574 managed_worktree: None,
575 review: None,
576 target_template_id: target_id.to_owned(),
577 resource_allocation: None,
578 additional_mounts: Vec::new(),
579 state: SessionState::Disconnected,
580 target: Some(locator),
581 native_session_id: None,
582 acp_session_title: None,
583 session_title_override: None,
584 created_at: now.clone(),
585 updated_at: now,
586 viewed_through_event_ordinal: 0,
587 draft_input: String::new(),
588 last_error: None,
589 last_checkpoint_error: None,
590 checkpoint: None,
591 }
592}
593
594fn adoption_unfinished(record: &SessionRecord, target_id: &str) -> bool {
599 record.state == SessionState::Disconnected
600 && record.native_session_id.is_none()
601 && record.target_template_id == target_id
602 && record.target.is_some()
603}
604
605fn scan_target_workers(
606 target_id: &str,
607 template: &TargetTemplate,
608 executor: &impl CommandExecutor,
609) -> Result<Vec<RecoveryCandidate>> {
610 let mut candidates = match template {
611 TargetTemplate::LocalBare => Vec::new(),
614 TargetTemplate::LocalPodman { .. } => scan_container_engine(
615 target_id,
616 template,
617 "podman",
618 vec![
619 "ps".into(),
620 "--all".into(),
621 "--filter".into(),
622 format!("label={}=true", targets::MANAGED_LABEL),
623 "--format".into(),
624 "json".into(),
625 ],
626 executor,
627 )?,
628 TargetTemplate::LocalDocker { .. } => scan_container_engine(
629 target_id,
630 template,
631 "docker",
632 vec![
633 "ps".into(),
634 "--all".into(),
635 "--filter".into(),
636 format!("label={}=true", targets::MANAGED_LABEL),
637 "--format".into(),
638 "json".into(),
639 ],
640 executor,
641 )?,
642 TargetTemplate::AppleContainer { .. } => scan_container_engine(
643 target_id,
644 template,
645 "container",
646 vec![
647 "list".into(),
648 "--all".into(),
649 "--format".into(),
650 "json".into(),
651 ],
652 executor,
653 )?,
654 TargetTemplate::SshPodman { ssh, .. } => {
655 let remote = targets::join_remote_command(&[
656 "podman".into(),
657 "ps".into(),
658 "--all".into(),
659 "--filter".into(),
660 format!("label={}=true", targets::MANAGED_LABEL),
661 "--format".into(),
662 "json".into(),
663 ]);
664 let output = execute_scan(
665 executor,
666 ssh_spec(ssh, [remote]),
667 "scan remote Podman workers",
668 )?;
669 candidates_from_container_json(target_id, template, &output.stdout)?
670 }
671 TargetTemplate::SshDocker { ssh, .. } => {
672 let remote = targets::join_remote_command(&[
673 "docker".into(),
674 "ps".into(),
675 "--all".into(),
676 "--filter".into(),
677 format!("label={}=true", targets::MANAGED_LABEL),
678 "--format".into(),
679 "json".into(),
680 ]);
681 let output = execute_scan(
682 executor,
683 ssh_spec(ssh, [remote]),
684 "scan remote Docker workers",
685 )?;
686 candidates_from_container_json(target_id, template, &output.stdout)?
687 }
688 TargetTemplate::AwsEc2 {
689 aws_profile,
690 region,
691 address_source,
692 ..
693 } => {
694 let profile = aws_profile.clone().unwrap_or_else(|| "default".into());
695 let output = execute_scan(
696 executor,
697 CommandSpec::new(
698 "aws",
699 [
700 "--profile".into(),
701 profile,
702 "--region".into(),
703 region.clone(),
704 "ec2".into(),
705 "describe-instances".into(),
706 "--filters".into(),
707 format!("Name=tag:{},Values=true", targets::MANAGED_TAG),
708 "Name=instance-state-name,Values=pending,running,stopping,stopped".into(),
709 "--output".into(),
710 "json".into(),
711 ],
712 )
713 .purpose("scan managed EC2 workers"),
714 "scan managed EC2 workers",
715 )?;
716 candidates_from_aws_json(target_id, address_source.clone(), &output.stdout)?
717 }
718 TargetTemplate::SshBare { ssh, .. } => {
719 let output = execute_scan(
720 executor,
721 ssh_spec(
722 ssh,
723 [targets::join_remote_command(&[
724 "find".into(),
725 ".local/share/hel/workers".into(),
726 "-mindepth".into(),
727 "2".into(),
728 "-maxdepth".into(),
729 "2".into(),
730 "-name".into(),
731 "ownership.json".into(),
732 "-print".into(),
733 ])],
734 ),
735 "scan bare SSH worker markers",
736 )?;
737 output
738 .stdout
739 .split(|byte| *byte == b'\n')
740 .filter_map(|line| {
741 let path = match std::str::from_utf8(line) {
742 Ok(path) => path.trim(),
743 Err(error) => {
744 tracing::debug!(%error, "recovery scan skipped a non-UTF-8 worker marker path");
745 return None;
746 }
747 };
748 let Some(session_id) = Path::new(path)
749 .parent()
750 .and_then(|parent| parent.file_name())
751 .and_then(|name| name.to_str())
752 else {
753 tracing::debug!(path, "recovery scan skipped a malformed worker marker path");
754 return None;
755 };
756 if let Err(error) = targets::resource_name(session_id) {
757 tracing::debug!(session_id, %error, "recovery scan skipped an invalid session id");
758 return None;
759 }
760 let backend = match backend_target(template, None, ContainerOverrides::default()) {
761 Ok(backend) => backend,
762 Err(error) => {
763 tracing::debug!(session_id, %error, "recovery scan could not construct the target backend");
764 return None;
765 }
766 };
767 let workspace = match targets::workspace_for(&backend, session_id) {
768 Ok(workspace) => workspace,
769 Err(error) => {
770 tracing::debug!(session_id, %error, "recovery scan could not derive the target workspace");
771 return None;
772 }
773 };
774 Some(RecoveryCandidate {
775 session_id: session_id.to_owned(),
776 target_template_id: target_id.to_owned(),
777 locator: TargetLocator::SshBare {
778 host: ssh.host.clone(),
779 workspace: PathBuf::from(workspace),
780 worker_id: None,
784 },
785 ownership: None,
786 instance_id: None,
787 tracked_session: None,
788 })
789 })
790 .collect()
791 }
792 };
793 for candidate in &mut candidates {
794 candidate.ownership = read_recovery_ownership(template, candidate, executor);
795 if candidate.instance_id.is_none() {
798 candidate.instance_id = candidate
799 .ownership
800 .as_ref()
801 .and_then(|marker| marker.instance_id.clone());
802 }
803 }
804 Ok(candidates)
805}
806
807fn scan_container_engine(
808 target_id: &str,
809 template: &TargetTemplate,
810 engine: &str,
811 args: Vec<String>,
812 executor: &impl CommandExecutor,
813) -> Result<Vec<RecoveryCandidate>> {
814 let output = execute_scan(
815 executor,
816 CommandSpec::new(engine, args).purpose("scan managed container workers"),
817 "scan managed container workers",
818 )?;
819 candidates_from_container_json(target_id, template, &output.stdout)
820}
821
822fn recovery_workspace_storage(
825 template: &TargetTemplate,
826 session_id: &str,
827) -> Result<PodmanWorkspaceLocator> {
828 let backend = backend_target(template, None, ContainerOverrides::default())?;
829 let container = match &backend {
830 targets::TargetTemplate::LocalPodman(container) => container,
831 targets::TargetTemplate::SshPodman { container, .. } => container,
832 _ => bail!("target template is not a Podman target"),
833 };
834 Ok(PodmanWorkspaceLocator::from(
835 targets::podman_workspace_locator(container, session_id)?,
836 ))
837}
838
839fn candidates_from_container_json(
840 target_id: &str,
841 template: &TargetTemplate,
842 stdout: &[u8],
843) -> Result<Vec<RecoveryCandidate>> {
844 let sessions = managed_sessions_from_container_json(stdout)?;
845 Ok(sessions
846 .into_iter()
847 .filter_map(|(session_id, instance_id)| {
848 let generated = match targets::resource_name(&session_id) {
849 Ok(generated) => generated,
850 Err(error) => {
851 tracing::debug!(%session_id, %error, "recovery scan skipped an invalid managed session id");
852 return None;
853 }
854 };
855 let workspace_storage = match template {
860 TargetTemplate::LocalPodman { .. } | TargetTemplate::SshPodman { .. } => {
861 match recovery_workspace_storage(template, &session_id) {
862 Ok(storage) => storage,
863 Err(error) => {
864 tracing::debug!(%session_id, %error, "recovery scan could not derive the Podman workspace storage");
865 return None;
866 }
867 }
868 }
869 _ => Default::default(),
870 };
871 let locator = match template {
872 TargetTemplate::LocalPodman { .. } => TargetLocator::LocalPodman {
873 borrowed_from: None,
874 container_id: generated,
875 workspace_storage,
876 },
877 TargetTemplate::LocalDocker { .. } => TargetLocator::LocalDocker {
878 borrowed_from: None,
879 container_id: generated,
880 },
881 TargetTemplate::AppleContainer { .. } => TargetLocator::AppleContainer {
882 borrowed_from: None,
883 container_id: generated,
884 },
885 TargetTemplate::SshPodman { ssh, .. } => TargetLocator::SshPodman {
886 borrowed_from: None,
887 host: ssh.host.clone(),
888 container_id: generated,
889 workspace_storage,
890 },
891 TargetTemplate::SshDocker { ssh, .. } => TargetLocator::SshDocker {
892 borrowed_from: None,
893 host: ssh.host.clone(),
894 container_id: generated,
895 },
896 _ => return None,
897 };
898 Some(RecoveryCandidate {
899 session_id,
900 target_template_id: target_id.to_owned(),
901 locator,
902 ownership: None,
903 instance_id,
904 tracked_session: None,
905 })
906 })
907 .collect())
908}
909
910pub(super) fn managed_sessions_from_container_json(
913 stdout: &[u8],
914) -> Result<Vec<(String, Option<String>)>> {
915 let values = serde_json::Deserializer::from_slice(stdout)
916 .into_iter::<serde_json::Value>()
917 .collect::<std::result::Result<Vec<_>, _>>()
918 .context("parse container list JSON")?;
919 let mut sessions = Vec::new();
920 for value in &values {
921 collect_managed_sessions(value, &mut sessions);
922 }
923 sessions.sort();
924 sessions.dedup();
925 Ok(sessions)
926}
927
928pub(super) fn collect_managed_sessions(
929 value: &serde_json::Value,
930 sessions: &mut Vec<(String, Option<String>)>,
931) {
932 match value {
933 serde_json::Value::Array(values) => {
934 for value in values {
935 collect_managed_sessions(value, sessions);
936 }
937 }
938 serde_json::Value::Object(object) => {
939 for label_key in ["Labels", "labels"] {
940 if let Some(labels) = object.get(label_key) {
941 let managed = label_value(labels, targets::MANAGED_LABEL)
942 .is_some_and(|value| value == "true");
943 if managed && let Some(session) = label_value(labels, targets::SESSION_LABEL) {
944 sessions.push((session, label_value(labels, targets::INSTANCE_LABEL)));
945 }
946 }
947 }
948 for value in object.values() {
949 collect_managed_sessions(value, sessions);
950 }
951 }
952 _ => {}
953 }
954}
955
956fn label_value(labels: &serde_json::Value, key: &str) -> Option<String> {
957 match labels {
958 serde_json::Value::Object(object) => object.get(key)?.as_str().map(str::to_owned),
959 serde_json::Value::String(text) => text
960 .split(',')
961 .find_map(|label| {
962 label
963 .trim()
964 .split_once('=')
965 .filter(|(name, _)| *name == key)
966 })
967 .map(|(_, value)| value.to_owned()),
968 _ => None,
969 }
970}
971
972fn candidates_from_aws_json(
973 target_id: &str,
974 address_source: AwsAddressSource,
975 stdout: &[u8],
976) -> Result<Vec<RecoveryCandidate>> {
977 let value: serde_json::Value =
978 serde_json::from_slice(stdout).context("parse AWS instance JSON")?;
979 let mut result = Vec::new();
980 let reservations = value
981 .get("Reservations")
982 .and_then(serde_json::Value::as_array)
983 .cloned()
984 .unwrap_or_default();
985 for instance in reservations.iter().flat_map(|reservation| {
986 reservation
987 .get("Instances")
988 .and_then(serde_json::Value::as_array)
989 .into_iter()
990 .flatten()
991 }) {
992 let tags = instance
993 .get("Tags")
994 .and_then(serde_json::Value::as_array)
995 .cloned()
996 .unwrap_or_default();
997 let tag = |key: &str| {
998 tags.iter()
999 .find(|tag| tag.get("Key").and_then(serde_json::Value::as_str) == Some(key))
1000 .and_then(|tag| tag.get("Value"))
1001 .and_then(serde_json::Value::as_str)
1002 };
1003 if tag(targets::MANAGED_TAG) != Some("true") {
1004 continue;
1005 }
1006 let Some(session_id) = tag(targets::SESSION_TAG).map(str::to_owned) else {
1007 continue;
1008 };
1009 targets::resource_name(&session_id)?;
1010 let instance_id = instance
1011 .get("InstanceId")
1012 .and_then(serde_json::Value::as_str)
1013 .context("managed EC2 instance omitted InstanceId")?
1014 .to_owned();
1015 let field = match address_source {
1016 AwsAddressSource::PublicDns => "PublicDnsName",
1017 AwsAddressSource::PublicIp => "PublicIpAddress",
1018 AwsAddressSource::PrivateDns => "PrivateDnsName",
1019 AwsAddressSource::PrivateIp => "PrivateIpAddress",
1020 };
1021 let address = instance
1022 .get(field)
1023 .and_then(serde_json::Value::as_str)
1024 .filter(|value| !value.is_empty())
1025 .map(str::to_owned);
1026 let created_by = tag(targets::INSTANCE_TAG).map(str::to_owned);
1027 result.push(RecoveryCandidate {
1028 session_id,
1029 target_template_id: target_id.to_owned(),
1030 locator: TargetLocator::AwsEc2 {
1031 instance_id,
1032 address,
1033 },
1034 ownership: None,
1035 instance_id: created_by,
1036 tracked_session: None,
1037 });
1038 }
1039 Ok(result)
1040}
1041
1042fn execute_scan(
1043 executor: &impl CommandExecutor,
1044 command: CommandSpec,
1045 operation: &str,
1046) -> Result<CommandOutput> {
1047 let output = executor.execute(&command)?;
1048 if output.status != 0 {
1049 bail!(
1050 "{operation} failed with status {}: {}",
1051 output.status,
1052 String::from_utf8_lossy(&output.stderr).trim()
1053 );
1054 }
1055 Ok(output)
1056}
1057
1058fn ssh_spec(ssh: &SshConnection, remote: impl IntoIterator<Item = String>) -> CommandSpec {
1059 let backend = SshTarget::from(ssh);
1060 let mut args = backend.ssh_args.clone();
1061 args.push(backend.destination.clone());
1062 args.extend(remote);
1063 CommandSpec::new("ssh", args).ssh_session(&backend)
1064}
1065
1066fn read_recovery_ownership(
1067 template: &TargetTemplate,
1068 candidate: &RecoveryCandidate,
1069 executor: &impl CommandExecutor,
1070) -> Option<WorkerOwnership> {
1071 let backend =
1072 match recovery_backend_locator(template, &candidate.locator, &candidate.session_id) {
1073 Ok(backend) => backend,
1074 Err(error) => {
1075 tracing::debug!(
1076 session_id = %candidate.session_id,
1077 %error,
1078 "could not construct a recovery ownership probe"
1079 );
1080 return None;
1081 }
1082 };
1083 let root = match targets::worker_root(&backend, &candidate.session_id) {
1084 Ok(root) => root,
1085 Err(error) => {
1086 tracing::debug!(
1087 session_id = %candidate.session_id,
1088 %error,
1089 "could not derive a recovery worker root"
1090 );
1091 return None;
1092 }
1093 };
1094 let command = match targets::command_on_locator(
1095 &backend,
1096 &candidate.session_id,
1097 vec!["cat".into(), format!("{root}/ownership.json")],
1098 "read worker ownership marker",
1099 ) {
1100 Ok(command) => command,
1101 Err(error) => {
1102 tracing::debug!(
1103 session_id = %candidate.session_id,
1104 %error,
1105 "could not construct a recovery ownership command"
1106 );
1107 return None;
1108 }
1109 };
1110 let output = match executor.execute(&command) {
1111 Ok(output) => output,
1112 Err(error) => {
1113 tracing::debug!(
1114 session_id = %candidate.session_id,
1115 %error,
1116 "could not read a recovery worker ownership marker"
1117 );
1118 return None;
1119 }
1120 };
1121 if output.status != 0 {
1122 tracing::debug!(
1123 session_id = %candidate.session_id,
1124 status = output.status,
1125 "recovery worker ownership probe returned a failure"
1126 );
1127 return None;
1128 }
1129 let marker: WorkerOwnership = match serde_json::from_slice(&output.stdout) {
1130 Ok(marker) => marker,
1131 Err(error) => {
1132 tracing::debug!(
1133 session_id = %candidate.session_id,
1134 %error,
1135 "recovery worker ownership marker was not valid JSON"
1136 );
1137 return None;
1138 }
1139 };
1140 if !(1..=WorkerOwnership::VERSION).contains(&marker.version)
1141 || marker.session_id != candidate.session_id
1142 || marker.target_template_id != candidate.target_template_id
1143 {
1144 tracing::debug!(
1145 session_id = %candidate.session_id,
1146 marker_session_id = %marker.session_id,
1147 marker_target_template_id = %marker.target_template_id,
1148 "recovery worker ownership marker did not match the candidate"
1149 );
1150 return None;
1151 }
1152 Some(marker)
1153}
1154
1155fn recovery_backend_locator(
1171 template: &TargetTemplate,
1172 locator: &TargetLocator,
1173 session_id: &str,
1174) -> Result<targets::TargetLocator> {
1175 Ok(match (template, locator) {
1176 (TargetTemplate::LocalBare, TargetLocator::LocalBare { worker_root }) => {
1177 targets::TargetLocator::LocalBare {
1178 worker_root: worker_root.to_string_lossy().into_owned(),
1179 }
1180 }
1181 (
1182 TargetTemplate::LocalPodman { .. },
1183 TargetLocator::LocalPodman {
1184 container_id,
1185 workspace_storage,
1186 ..
1187 },
1188 ) => targets::TargetLocator::LocalPodman {
1189 borrowed_from: None,
1190 container_id: container_id.clone(),
1191 workspace_storage: workspace_storage.into(),
1192 },
1193 (TargetTemplate::LocalDocker { .. }, TargetLocator::LocalDocker { container_id, .. }) => {
1194 targets::TargetLocator::LocalDocker {
1195 borrowed_from: None,
1196 container_id: container_id.clone(),
1197 }
1198 }
1199 (
1200 TargetTemplate::AppleContainer { .. },
1201 TargetLocator::AppleContainer { container_id, .. },
1202 ) => targets::TargetLocator::AppleContainer {
1203 borrowed_from: None,
1204 container_id: container_id.clone(),
1205 },
1206 (
1207 TargetTemplate::SshPodman { ssh, .. },
1208 TargetLocator::SshPodman {
1209 container_id,
1210 workspace_storage,
1211 ..
1212 },
1213 ) => targets::TargetLocator::SshPodman {
1214 borrowed_from: None,
1215 ssh: SshTarget::from(ssh),
1216 container_id: container_id.clone(),
1217 workspace_storage: workspace_storage.into(),
1218 },
1219 (
1220 TargetTemplate::SshDocker { ssh, .. },
1221 TargetLocator::SshDocker {
1222 host, container_id, ..
1223 },
1224 ) => {
1225 if host != &ssh.host {
1226 bail!("recovery SSH Docker host does not match target template")
1227 }
1228 targets::TargetLocator::SshDocker {
1229 borrowed_from: None,
1230 ssh: SshTarget::from(ssh),
1231 container_id: container_id.clone(),
1232 }
1233 }
1234 (TargetTemplate::SshBare { ssh, .. }, TargetLocator::SshBare { workspace, .. }) => {
1235 targets::TargetLocator::SshBare {
1236 ssh: SshTarget::from(ssh),
1237 workspace: workspace.to_string_lossy().into_owned(),
1238 worker_id: None,
1239 }
1240 }
1241 (
1242 TargetTemplate::AwsEc2 {
1243 aws_profile,
1244 region,
1245 ssh_user,
1246 identity_file,
1247 ssh_args,
1248 ..
1249 },
1250 TargetLocator::AwsEc2 {
1251 instance_id,
1252 address,
1253 },
1254 ) => targets::TargetLocator::AwsEc2 {
1255 profile: aws_profile.clone().unwrap_or_else(|| "default".into()),
1256 region: region.clone(),
1257 instance_id: instance_id.clone(),
1258 ssh: SshTarget {
1259 destination: format!(
1260 "{ssh_user}@{}",
1261 address.as_deref().unwrap_or("unavailable.invalid")
1262 ),
1263 ssh_args: targets::ssh_args_with_identity(ssh_args, identity_file.as_deref()),
1264 },
1265 workspace: format!(".local/share/hel/workspaces/{session_id}"),
1266 },
1267 _ => bail!("recovery target locator does not match target template"),
1268 })
1269}
1270
1271#[cfg(test)]
1272mod tests {
1273
1274 use crate::controller::test_support::{IsolatedTest, test_name};
1275 use mj_core::config::{
1276 AwsAddressSource, Config, ContainerTemplate as ConfigContainer, HarnessKind,
1277 PodmanWorkspaceStorage, TargetTemplate,
1278 };
1279 use mj_core::state::{State, TargetLocator};
1280
1281 use crate::targets::ProcessExecutor;
1282
1283 use super::*;
1284
1285 const FAILED_ADOPTION_CHILD: &str = "MJ_TEST_FAILED_ADOPTION_CHILD";
1286
1287 #[tokio::test]
1289 async fn a_failed_adoption_records_the_failure_and_stays_retryable() {
1290 if std::env::var_os(FAILED_ADOPTION_CHILD).is_none() {
1293 let directory = tempfile::tempdir().unwrap();
1294 IsolatedTest::new(test_name(
1295 module_path!(),
1296 "a_failed_adoption_records_the_failure_and_stays_retryable",
1297 ))
1298 .env(FAILED_ADOPTION_CHILD, "1")
1299 .env("MJ_DATA_DIR", directory.path())
1300 .run();
1301 return;
1302 }
1303 let _writer = crate::database::install_isolated_test_writer();
1305
1306 let session_id = "0123456789abcdef0123456789abcdef";
1307 let workers = tempfile::tempdir().unwrap();
1308 let record = adopted_session_record(
1311 session_id,
1312 "local-bare",
1313 "codex".into(),
1314 HarnessKind::Codex,
1315 "project".into(),
1316 mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1317 TargetLocator::LocalBare {
1318 worker_root: workers.path().join(session_id),
1319 },
1320 );
1321 assert!(
1322 adoption_unfinished(&record, "local-bare"),
1323 "the record adoption commits must be the record adoption can retry"
1324 );
1325 crate::database::save_session(&record).unwrap();
1326 let mut config = Config::default();
1327 config
1328 .targets
1329 .insert("local-bare".into(), TargetTemplate::LocalBare);
1330 let mut state = State::default();
1331 state.sessions.insert(session_id.to_owned(), record);
1332 let mut controller = Controller { config, state };
1333
1334 let failure = controller
1335 .adopt_orphan_worker(
1336 session_id,
1337 "local-bare",
1338 None,
1339 None,
1340 false,
1341 &ProcessExecutor,
1342 )
1343 .await
1344 .expect_err("a worker root without a worker cannot complete the handshake");
1345 assert!(
1346 format!("{failure:#}").contains("orphan relay"),
1347 "unexpected failure: {failure:#}"
1348 );
1349 let recorded = controller.state.sessions[session_id]
1350 .last_error
1351 .clone()
1352 .expect("the failed handshake was recorded on the session");
1353 assert!(
1354 recorded.contains("orphan adoption failed"),
1355 "unexpected recorded failure: {recorded}"
1356 );
1357 assert_eq!(
1358 controller.state.sessions[session_id].state,
1359 SessionState::Disconnected
1360 );
1361 let stored = crate::database::load_state().unwrap();
1362 assert_eq!(
1363 stored.sessions[session_id].last_error.as_deref(),
1364 Some(recorded.as_str()),
1365 "the adoption failure was not persisted"
1366 );
1367
1368 let retry = controller
1369 .adopt_orphan_worker(
1370 session_id,
1371 "local-bare",
1372 None,
1373 None,
1374 false,
1375 &ProcessExecutor,
1376 )
1377 .await
1378 .expect_err("the worker is still unreachable");
1379 let retry = format!("{retry:#}");
1380 assert!(
1381 retry.contains("orphan relay"),
1382 "adoption did not retry the handshake: {retry}"
1383 );
1384 assert!(
1385 !retry.contains("already tracked"),
1386 "a session adoption never finished blocked its own retry: {retry}"
1387 );
1388 }
1389
1390 const RECOVERY_WORKSPACE_CHILD: &str = "MJ_TEST_RECOVERY_WORKSPACE_CHILD";
1391
1392 #[tokio::test]
1393 async fn orphan_workspace_ids_are_reconciled_before_adoption_persistence() {
1394 if std::env::var_os(RECOVERY_WORKSPACE_CHILD).is_none() {
1395 let directory = tempfile::tempdir().unwrap();
1396 IsolatedTest::new(test_name(
1397 module_path!(),
1398 "orphan_workspace_ids_are_reconciled_before_adoption_persistence",
1399 ))
1400 .env(RECOVERY_WORKSPACE_CHILD, "1")
1401 .env("MJ_DATA_DIR", directory.path())
1402 .run();
1403 return;
1404 }
1405
1406 let _writer = crate::database::install_isolated_test_writer();
1407 let default =
1408 resolve_recovery_workspace_id(mj_core::workspace::DEFAULT_WORKSPACE_ID).unwrap();
1409 assert_eq!(default, mj_core::workspace::DEFAULT_WORKSPACE_ID);
1410
1411 let known = crate::database::create_or_get_workspace("Known").unwrap();
1412 assert_eq!(resolve_recovery_workspace_id(&known.id).unwrap(), known.id);
1413
1414 let recovered = resolve_recovery_workspace_id("workspace-from-old-controller").unwrap();
1415 let repeated = resolve_recovery_workspace_id("another-old-workspace").unwrap();
1416 assert_eq!(repeated, recovered);
1417 assert_eq!(
1418 crate::database::list_workspaces()
1419 .unwrap()
1420 .iter()
1421 .filter(|workspace| workspace.name == "Recovered")
1422 .count(),
1423 1
1424 );
1425
1426 let session_id = "0123456789abcdef0123456789abcdef";
1427 let workers = tempfile::tempdir().unwrap();
1428 let record = adopted_session_record(
1429 session_id,
1430 "local-bare",
1431 "codex".into(),
1432 HarnessKind::Codex,
1433 "project".into(),
1434 recovered.clone(),
1435 TargetLocator::LocalBare {
1436 worker_root: workers.path().join(session_id),
1437 },
1438 );
1439 crate::database::save_session(&record).unwrap();
1440 assert_eq!(
1441 crate::database::load_state().unwrap().sessions[session_id].workspace_id,
1442 recovered
1443 );
1444
1445 let mut state = State::default();
1449 state.sessions.insert(session_id.to_owned(), record);
1450 let mut config = Config::default();
1451 config
1452 .targets
1453 .insert("local-bare".into(), TargetTemplate::LocalBare);
1454 let mut controller = Controller { config, state };
1455 let failure = controller
1456 .adopt_orphan_worker(
1457 session_id,
1458 "local-bare",
1459 None,
1460 None,
1461 false,
1462 &ProcessExecutor,
1463 )
1464 .await
1465 .expect_err("a worker root without a worker cannot complete the handshake");
1466 assert!(
1467 format!("{failure:#}").contains("orphan relay"),
1468 "unexpected failure: {failure:#}"
1469 );
1470 let stored = crate::database::load_state().unwrap();
1471 assert_eq!(stored.sessions[session_id].workspace_id, recovered);
1472 assert!(
1473 stored.sessions[session_id]
1474 .last_error
1475 .as_deref()
1476 .is_some_and(|error| error.contains("orphan adoption failed"))
1477 );
1478 }
1479
1480 #[test]
1481 fn a_session_that_completed_its_handshake_is_not_adoptable_again() {
1482 let mut record = adopted_session_record(
1483 "0123456789abcdef0123456789abcdef",
1484 "local-bare",
1485 "codex".into(),
1486 HarnessKind::Codex,
1487 "project".into(),
1488 mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1489 TargetLocator::LocalBare {
1490 worker_root: std::path::PathBuf::from("/workers/0123456789abcdef0123456789abcdef"),
1491 },
1492 );
1493 record.native_session_id = Some("native-session".into());
1494 assert!(!adoption_unfinished(&record, "local-bare"));
1495
1496 record.native_session_id = None;
1497 assert!(
1498 !adoption_unfinished(&record, "other-target"),
1499 "a record adopted onto another target is not this target's retry"
1500 );
1501 }
1502
1503 #[test]
1504 fn recovery_container_scan_requires_both_managed_and_session_labels() {
1505 let template = TargetTemplate::LocalPodman {
1506 container: ConfigContainer {
1507 build_cache: None,
1508 image: "ignored".into(),
1509 pull_policy: Default::default(),
1510 platform: None,
1511 cpus: None,
1512 memory: None,
1513 environment: Default::default(),
1514 workspace_storage: Default::default(),
1515 },
1516 };
1517 let json = serde_json::json!([
1518 {"Labels": {"dev.mj.managed": "true", "dev.mj.session": "0123456789abcdef0123456789abcdef", "dev.mj.instance": "qa0916"}},
1519 {"Labels": {"dev.mj.managed": "false", "dev.mj.session": "not-owned"}},
1520 {"configuration": {"labels": "dev.mj.managed=true,dev.mj.session=abcdef0123456789abcdef0123456789"}}
1521 ]);
1522 let candidates = candidates_from_container_json(
1523 "local",
1524 &template,
1525 serde_json::to_string(&json).unwrap().as_bytes(),
1526 )
1527 .unwrap();
1528 assert_eq!(candidates.len(), 2);
1529 assert_eq!(candidates[0].session_id, "0123456789abcdef0123456789abcdef");
1530 assert_eq!(candidates[0].instance_id.as_deref(), Some("qa0916"));
1531 assert_eq!(
1532 candidates[1].instance_id, None,
1533 "a container without an instance label is of unknown origin"
1534 );
1535 }
1536
1537 struct ListingExecutor {
1541 listing: String,
1542 }
1543
1544 impl CommandExecutor for ListingExecutor {
1545 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1546 if command.program == "podman" && command.args.first().map(String::as_str) == Some("ps")
1547 {
1548 return Ok(CommandOutput {
1549 status: 0,
1550 stdout: self.listing.clone().into_bytes(),
1551 stderr: Vec::new(),
1552 });
1553 }
1554 Ok(CommandOutput {
1555 status: 1,
1556 stdout: Vec::new(),
1557 stderr: Vec::new(),
1558 })
1559 }
1560 }
1561
1562 fn podman_scan_controller(record: SessionRecord) -> (Controller, ListingExecutor) {
1563 let mut config = Config::default();
1564 config.targets.insert(
1565 "local".into(),
1566 TargetTemplate::LocalPodman {
1567 container: ConfigContainer {
1568 build_cache: None,
1569 image: "ignored".into(),
1570 pull_policy: Default::default(),
1571 platform: None,
1572 cpus: None,
1573 memory: None,
1574 environment: Default::default(),
1575 workspace_storage: Default::default(),
1576 },
1577 },
1578 );
1579 let listing = serde_json::to_string(&serde_json::json!([
1580 {"Labels": {"dev.mj.managed": "true", "dev.mj.session": record.id}}
1581 ]))
1582 .unwrap();
1583 let mut state = State::default();
1584 state.sessions.insert(record.id.clone(), record);
1585 (Controller { config, state }, ListingExecutor { listing })
1586 }
1587
1588 fn errored_podman_session(session_id: &str) -> SessionRecord {
1589 let mut record = adopted_session_record(
1590 session_id,
1591 "local",
1592 "codex".into(),
1593 HarnessKind::Codex,
1594 "project".into(),
1595 mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1596 TargetLocator::LocalPodman {
1597 container_id: targets::resource_name(session_id).unwrap(),
1598 workspace_storage: PodmanWorkspaceLocator::ContainerLayer,
1599 borrowed_from: None,
1600 },
1601 );
1602 record.state = SessionState::Error;
1605 record.target = None;
1606 record
1607 }
1608
1609 fn isolated_scan(name: &str) -> bool {
1610 if std::env::var_os("MJ_TEST_MOVE_OWNERSHIP_SCAN").is_some() {
1611 return true;
1612 }
1613 let directory = tempfile::tempdir().unwrap();
1614 IsolatedTest::new(test_name(module_path!(), name))
1615 .env("MJ_TEST_MOVE_OWNERSHIP_SCAN", "1")
1616 .isolated_store(directory.path())
1617 .run();
1618 false
1619 }
1620
1621 #[test]
1623 fn a_container_an_errored_session_no_longer_names_is_an_orphan() {
1624 if !isolated_scan("a_container_an_errored_session_no_longer_names_is_an_orphan") {
1625 return;
1626 }
1627 let _writer = crate::database::install_isolated_test_writer();
1628 crate::database::save_state(&State::default()).unwrap();
1629 let session_id = "0123456789abcdef0123456789abcdef";
1630 let (controller, executor) = podman_scan_controller(errored_podman_session(session_id));
1631
1632 let scan = controller.scan_orphan_workers(&executor, true);
1633
1634 let [candidate] = scan.candidates.as_slice() else {
1635 panic!("the errored session's container was not offered: {scan:?}");
1636 };
1637 assert_eq!(candidate.session_id, session_id);
1638 assert_eq!(
1639 candidate.tracked_session,
1640 Some(SessionState::Error),
1641 "recovery must say the session id is still taken, so only destroy applies"
1642 );
1643 }
1644
1645 #[test]
1646 fn a_container_a_session_still_names_is_not_an_orphan() {
1647 if !isolated_scan("a_container_a_session_still_names_is_not_an_orphan") {
1648 return;
1649 }
1650 let _writer = crate::database::install_isolated_test_writer();
1651 crate::database::save_state(&State::default()).unwrap();
1652 let session_id = "0123456789abcdef0123456789abcdef";
1653 let mut record = errored_podman_session(session_id);
1654 record.target = Some(TargetLocator::LocalPodman {
1657 container_id: targets::resource_name(session_id).unwrap(),
1658 workspace_storage: PodmanWorkspaceLocator::ContainerLayer,
1659 borrowed_from: None,
1660 });
1661 let (controller, executor) = podman_scan_controller(record);
1662
1663 let scan = controller.scan_orphan_workers(&executor, true);
1664
1665 assert!(
1666 scan.candidates.is_empty(),
1667 "a container the controller can still drive was offered for destruction: {scan:?}"
1668 );
1669 }
1670
1671 #[test]
1672 fn a_container_being_provisioned_is_not_an_orphan() {
1673 if !isolated_scan("a_container_being_provisioned_is_not_an_orphan") {
1674 return;
1675 }
1676 let _writer = crate::database::install_isolated_test_writer();
1677 crate::database::save_state(&State::default()).unwrap();
1678 let session_id = "0123456789abcdef0123456789abcdef";
1679 let mut record = errored_podman_session(session_id);
1680 record.state = SessionState::Provisioning;
1684 let (controller, executor) = podman_scan_controller(record);
1685
1686 let scan = controller.scan_orphan_workers(&executor, true);
1687
1688 assert!(
1689 scan.candidates.is_empty(),
1690 "a container still being provisioned was offered for destruction: {scan:?}"
1691 );
1692 }
1693
1694 fn candidate(session_id: &str, instance_id: Option<&str>) -> RecoveryCandidate {
1695 RecoveryCandidate {
1696 session_id: session_id.to_owned(),
1697 target_template_id: "local".to_owned(),
1698 locator: TargetLocator::LocalDocker {
1699 container_id: format!("mj-{session_id}"),
1700 borrowed_from: None,
1701 },
1702 ownership: None,
1703 instance_id: instance_id.map(str::to_owned),
1704 tracked_session: None,
1705 }
1706 }
1707
1708 #[test]
1710 fn default_scan_scope_hides_other_and_unknown_instances() {
1711 let mut scan = RecoveryScan {
1712 candidates: vec![
1713 candidate("mine", Some("qa")),
1714 candidate("theirs", Some("prod")),
1715 candidate("legacy", None),
1716 ],
1717 instance_id: "qa".to_owned(),
1718 ..RecoveryScan::default()
1719 };
1720 restrict_to_instance(&mut scan);
1721 assert_eq!(
1722 scan.candidates
1723 .iter()
1724 .map(|candidate| candidate.session_id.as_str())
1725 .collect::<Vec<_>>(),
1726 ["mine"]
1727 );
1728 assert_eq!(scan.hidden_other_instances, 2);
1729 }
1730
1731 #[test]
1733 fn acting_on_another_or_unknown_instance_requires_the_explicit_flag() {
1734 require_instance_access(&candidate("mine", Some("qa")), "qa", false).unwrap();
1735
1736 let other = require_instance_access(&candidate("theirs", Some("prod")), "qa", false)
1737 .expect_err("another instance's worker is refused by default");
1738 assert!(
1739 other.to_string().contains("belongs to instance \"prod\""),
1740 "{other}"
1741 );
1742 require_instance_access(&candidate("theirs", Some("prod")), "qa", true).unwrap();
1743
1744 let unknown = require_instance_access(&candidate("legacy", None), "qa", false)
1745 .expect_err("a worker without a stamp is refused by default");
1746 assert!(
1747 unknown.to_string().contains("no instance stamp"),
1748 "{unknown}"
1749 );
1750 require_instance_access(&candidate("legacy", None), "qa", true).unwrap();
1751 }
1752
1753 #[test]
1755 fn a_podman_orphan_destroy_plan_removes_its_workspace_volume() {
1756 let session = "0123456789abcdef0123456789abcdef";
1760 let template = TargetTemplate::LocalPodman {
1761 container: ConfigContainer {
1762 image: "ignored".into(),
1763 pull_policy: Default::default(),
1764 platform: None,
1765 cpus: None,
1766 memory: None,
1767 environment: Default::default(),
1768 workspace_storage: PodmanWorkspaceStorage::PodmanVolume,
1769 build_cache: None,
1770 },
1771 };
1772 let json = serde_json::json!([
1773 {"Labels": {"dev.mj.managed": "true", "dev.mj.session": session}}
1774 ]);
1775
1776 let candidates = candidates_from_container_json(
1777 "local",
1778 &template,
1779 serde_json::to_string(&json).unwrap().as_bytes(),
1780 )
1781 .unwrap();
1782
1783 let [candidate] = candidates.as_slice() else {
1784 panic!("expected one orphan candidate, got {candidates:?}");
1785 };
1786 let volume = format!("{}-workspace", targets::resource_name(session).unwrap());
1787 assert!(
1788 matches!(
1789 &candidate.locator,
1790 TargetLocator::LocalPodman {
1791 workspace_storage: PodmanWorkspaceLocator::Volume { name },
1792 ..
1793 } if name == &volume
1794 ),
1795 "candidate locator lost the volume storage: {:?}",
1796 candidate.locator
1797 );
1798 let backend = recovery_backend_locator(&template, &candidate.locator, session).unwrap();
1799 let plan = targets::close_plan(&backend, session).unwrap();
1800 assert!(
1801 plan.commands.iter().any(|command| {
1802 command.args.iter().any(|argument| argument == &volume)
1803 && command
1804 .args
1805 .iter()
1806 .any(|argument| argument.contains("podman volume rm"))
1807 }),
1808 "destroy plan does not remove the workspace volume: {plan:?}"
1809 );
1810 }
1811
1812 #[test]
1813 fn recovery_docker_scan_accepts_json_lines_and_builds_a_docker_locator() {
1814 let template = TargetTemplate::LocalDocker {
1815 container: ConfigContainer {
1816 build_cache: None,
1817 image: "ignored".into(),
1818 pull_policy: Default::default(),
1819 platform: None,
1820 cpus: None,
1821 memory: None,
1822 environment: Default::default(),
1823 workspace_storage: Default::default(),
1824 },
1825 };
1826 let session = "0123456789abcdef0123456789abcdef";
1827 let output = format!(
1828 "{{\"Labels\":\"dev.mj.managed=true,dev.mj.session={session},dev.mj.instance=abc123\"}}\n{{\"Labels\":\"dev.mj.managed=false,dev.mj.session=ignored\"}}\n"
1829 );
1830
1831 let candidates =
1832 candidates_from_container_json("docker", &template, output.as_bytes()).unwrap();
1833
1834 assert_eq!(candidates.len(), 1);
1835 assert_eq!(candidates[0].session_id, session);
1836 assert_eq!(candidates[0].instance_id.as_deref(), Some("abc123"));
1837 assert!(matches!(
1838 &candidates[0].locator,
1839 TargetLocator::LocalDocker { container_id, .. }
1840 if container_id == &targets::resource_name(session).unwrap()
1841 ));
1842 }
1843
1844 #[test]
1845 fn recovery_aws_scan_uses_exact_tagged_instance_and_address() {
1846 let json = serde_json::json!({"Reservations": [{"Instances": [{
1847 "InstanceId": "i-exact",
1848 "PrivateIpAddress": "10.0.0.7",
1849 "Tags": [
1850 {"Key": "dev.mj.managed", "Value": "true"},
1851 {"Key": "dev.mj.session", "Value": "0123456789abcdef0123456789abcdef"},
1852 {"Key": "dev.mj.instance", "Value": "qa0916"}
1853 ]
1854 }]}]});
1855 let candidates = candidates_from_aws_json(
1856 "aws",
1857 AwsAddressSource::PrivateIp,
1858 serde_json::to_string(&json).unwrap().as_bytes(),
1859 )
1860 .unwrap();
1861 assert_eq!(candidates[0].instance_id.as_deref(), Some("qa0916"));
1862 assert!(matches!(
1863 &candidates[0].locator,
1864 TargetLocator::AwsEc2 { instance_id, address }
1865 if instance_id == "i-exact" && address.as_deref() == Some("10.0.0.7")
1866 ));
1867 }
1868}