Skip to main content

mj_controller/controller/
recovery_scan.rs

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