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        checkout: None,
517        expected_runtime_identity: None,
518        publication: None,
519        build_cache: None,
520        subagents: None,
521        // The adopting caller probes the running container for this.
522        container_workspace: None,
523        create_managed_worktree: None,
524        workspace_id,
525        archived: false,
526        container_cpus: None,
527        container_memory: None,
528        id: session_id.to_owned(),
529        title: format!("Recovered {}", &session_id[..session_id.len().min(8)]),
530        harness_kind,
531        last_profile: profile_id,
532        bundle_id,
533        project_directory: None,
534        managed_worktree: None,
535        target_template_id: target_id.to_owned(),
536        resource_allocation: None,
537        additional_mounts: Vec::new(),
538        state: SessionState::Disconnected,
539        target: Some(locator),
540        native_session_id: None,
541        acp_session_title: None,
542        session_title_override: None,
543        created_at: now.clone(),
544        updated_at: now,
545        viewed_through_event_ordinal: 0,
546        draft_input: String::new(),
547        last_error: None,
548        last_checkpoint_error: None,
549        checkpoint: None,
550    }
551}
552
553/// Whether a tracked session is one an adoption committed and never finished:
554/// it names this target, carries the locator the scan found, and no harness
555/// session has ever been observed on it. Such a record is the retry, so
556/// adoption completes it instead of refusing it as already tracked.
557fn adoption_unfinished(record: &SessionRecord, target_id: &str) -> bool {
558    record.state == SessionState::Disconnected
559        && record.native_session_id.is_none()
560        && record.target_template_id == target_id
561        && record.target.is_some()
562}
563
564fn scan_target_workers(
565    target_id: &str,
566    template: &TargetTemplate,
567    executor: &impl CommandExecutor,
568) -> Result<Vec<RecoveryCandidate>> {
569    let mut candidates = match template {
570        // Local bare sessions persist their locator in the controller database.
571        // Do not infer an adoptable project from Hel's transient worker directory.
572        TargetTemplate::LocalBare => Vec::new(),
573        TargetTemplate::LocalPodman { .. } => scan_container_engine(
574            target_id,
575            template,
576            "podman",
577            vec![
578                "ps".into(),
579                "--all".into(),
580                "--filter".into(),
581                format!("label={}=true", targets::MANAGED_LABEL),
582                "--format".into(),
583                "json".into(),
584            ],
585            executor,
586        )?,
587        TargetTemplate::LocalDocker { .. } => scan_container_engine(
588            target_id,
589            template,
590            "docker",
591            vec![
592                "ps".into(),
593                "--all".into(),
594                "--filter".into(),
595                format!("label={}=true", targets::MANAGED_LABEL),
596                "--format".into(),
597                "json".into(),
598            ],
599            executor,
600        )?,
601        TargetTemplate::AppleContainer { .. } => scan_container_engine(
602            target_id,
603            template,
604            "container",
605            vec![
606                "list".into(),
607                "--all".into(),
608                "--format".into(),
609                "json".into(),
610            ],
611            executor,
612        )?,
613        TargetTemplate::SshPodman { ssh, .. } => {
614            let remote = targets::join_remote_command(&[
615                "podman".into(),
616                "ps".into(),
617                "--all".into(),
618                "--filter".into(),
619                format!("label={}=true", targets::MANAGED_LABEL),
620                "--format".into(),
621                "json".into(),
622            ]);
623            let output = execute_scan(
624                executor,
625                ssh_spec(ssh, [remote]),
626                "scan remote Podman workers",
627            )?;
628            candidates_from_container_json(target_id, template, &output.stdout)?
629        }
630        TargetTemplate::SshDocker { ssh, .. } => {
631            let remote = targets::join_remote_command(&[
632                "docker".into(),
633                "ps".into(),
634                "--all".into(),
635                "--filter".into(),
636                format!("label={}=true", targets::MANAGED_LABEL),
637                "--format".into(),
638                "json".into(),
639            ]);
640            let output = execute_scan(
641                executor,
642                ssh_spec(ssh, [remote]),
643                "scan remote Docker workers",
644            )?;
645            candidates_from_container_json(target_id, template, &output.stdout)?
646        }
647        TargetTemplate::AwsEc2 {
648            aws_profile,
649            region,
650            address_source,
651            ..
652        } => {
653            let profile = aws_profile.clone().unwrap_or_else(|| "default".into());
654            let output = execute_scan(
655                executor,
656                CommandSpec::new(
657                    "aws",
658                    [
659                        "--profile".into(),
660                        profile,
661                        "--region".into(),
662                        region.clone(),
663                        "ec2".into(),
664                        "describe-instances".into(),
665                        "--filters".into(),
666                        format!("Name=tag:{},Values=true", targets::MANAGED_TAG),
667                        "Name=instance-state-name,Values=pending,running,stopping,stopped".into(),
668                        "--output".into(),
669                        "json".into(),
670                    ],
671                )
672                .purpose("scan managed EC2 workers"),
673                "scan managed EC2 workers",
674            )?;
675            candidates_from_aws_json(target_id, address_source.clone(), &output.stdout)?
676        }
677        TargetTemplate::SshBare { ssh, .. } => {
678            let output = execute_scan(
679                executor,
680                ssh_spec(
681                    ssh,
682                    [targets::join_remote_command(&[
683                        "find".into(),
684                        ".local/share/hel/workers".into(),
685                        "-mindepth".into(),
686                        "2".into(),
687                        "-maxdepth".into(),
688                        "2".into(),
689                        "-name".into(),
690                        "ownership.json".into(),
691                        "-print".into(),
692                    ])],
693                ),
694                "scan bare SSH worker markers",
695            )?;
696            output
697                .stdout
698                .split(|byte| *byte == b'\n')
699                .filter_map(|line| {
700                    let path = match std::str::from_utf8(line) {
701                        Ok(path) => path.trim(),
702                        Err(error) => {
703                            tracing::debug!(%error, "recovery scan skipped a non-UTF-8 worker marker path");
704                            return None;
705                        }
706                    };
707                    let Some(session_id) = Path::new(path)
708                        .parent()
709                        .and_then(|parent| parent.file_name())
710                        .and_then(|name| name.to_str())
711                    else {
712                        tracing::debug!(path, "recovery scan skipped a malformed worker marker path");
713                        return None;
714                    };
715                    if let Err(error) = targets::resource_name(session_id) {
716                        tracing::debug!(session_id, %error, "recovery scan skipped an invalid session id");
717                        return None;
718                    }
719                    let backend = match backend_target(template, None, ContainerOverrides::default()) {
720                        Ok(backend) => backend,
721                        Err(error) => {
722                            tracing::debug!(session_id, %error, "recovery scan could not construct the target backend");
723                            return None;
724                        }
725                    };
726                    let workspace = match targets::workspace_for(&backend, session_id) {
727                        Ok(workspace) => workspace,
728                        Err(error) => {
729                            tracing::debug!(session_id, %error, "recovery scan could not derive the target workspace");
730                            return None;
731                        }
732                    };
733                    Some(RecoveryCandidate {
734                        session_id: session_id.to_owned(),
735                        target_template_id: target_id.to_owned(),
736                        locator: TargetLocator::SshBare {
737                            host: ssh.host.clone(),
738                            workspace: PathBuf::from(workspace),
739                            // The scan reads the session id from the worker
740                            // directory name and `targets::worker_root` falls
741                            // back to it, so the root still resolves.
742                            worker_id: None,
743                        },
744                        ownership: None,
745                        instance_id: None,
746                        tracked_session: None,
747                    })
748                })
749                .collect()
750        }
751    };
752    for candidate in &mut candidates {
753        candidate.ownership = read_recovery_ownership(template, candidate, executor);
754        // The label is authoritative; the marker only fills in when the
755        // resource carries no stamp (bare SSH workers have no labels at all).
756        if candidate.instance_id.is_none() {
757            candidate.instance_id = candidate
758                .ownership
759                .as_ref()
760                .and_then(|marker| marker.instance_id.clone());
761        }
762    }
763    Ok(candidates)
764}
765
766fn scan_container_engine(
767    target_id: &str,
768    template: &TargetTemplate,
769    engine: &str,
770    args: Vec<String>,
771    executor: &impl CommandExecutor,
772) -> Result<Vec<RecoveryCandidate>> {
773    let output = execute_scan(
774        executor,
775        CommandSpec::new(engine, args).purpose("scan managed container workers"),
776        "scan managed container workers",
777    )?;
778    candidates_from_container_json(target_id, template, &output.stdout)
779}
780
781/// Where a Podman session's workspace really lives, derived the same way
782/// provisioning derives it: the backend container template plus the session id.
783fn recovery_workspace_storage(
784    template: &TargetTemplate,
785    session_id: &str,
786) -> Result<PodmanWorkspaceLocator> {
787    let backend = backend_target(template, None, ContainerOverrides::default())?;
788    let container = match &backend {
789        targets::TargetTemplate::LocalPodman(container) => container,
790        targets::TargetTemplate::SshPodman { container, .. } => container,
791        _ => bail!("target template is not a Podman target"),
792    };
793    Ok(PodmanWorkspaceLocator::from(
794        targets::podman_workspace_locator(container, session_id)?,
795    ))
796}
797
798fn candidates_from_container_json(
799    target_id: &str,
800    template: &TargetTemplate,
801    stdout: &[u8],
802) -> Result<Vec<RecoveryCandidate>> {
803    let sessions = managed_sessions_from_container_json(stdout)?;
804    Ok(sessions
805        .into_iter()
806        .filter_map(|(session_id, instance_id)| {
807            let generated = match targets::resource_name(&session_id) {
808                Ok(generated) => generated,
809                Err(error) => {
810                    tracing::debug!(%session_id, %error, "recovery scan skipped an invalid managed session id");
811                    return None;
812                }
813            };
814            // A Podman workspace may live in a volume or a host directory, and
815            // destroying the container alone would leave it behind. The storage
816            // follows from the target template and the session id, so derive it
817            // here rather than recording the container layer by default.
818            let workspace_storage = match template {
819                TargetTemplate::LocalPodman { .. } | TargetTemplate::SshPodman { .. } => {
820                    match recovery_workspace_storage(template, &session_id) {
821                        Ok(storage) => storage,
822                        Err(error) => {
823                            tracing::debug!(%session_id, %error, "recovery scan could not derive the Podman workspace storage");
824                            return None;
825                        }
826                    }
827                }
828                _ => Default::default(),
829            };
830            let locator = match template {
831                TargetTemplate::LocalPodman { .. } => TargetLocator::LocalPodman {
832                    borrowed_from: None,
833                    container_id: generated,
834                    workspace_storage,
835                },
836                TargetTemplate::LocalDocker { .. } => TargetLocator::LocalDocker {
837                    borrowed_from: None,
838                    container_id: generated,
839                },
840                TargetTemplate::AppleContainer { .. } => TargetLocator::AppleContainer {
841                    borrowed_from: None,
842                    container_id: generated,
843                },
844                TargetTemplate::SshPodman { ssh, .. } => TargetLocator::SshPodman {
845                    borrowed_from: None,
846                    host: ssh.host.clone(),
847                    container_id: generated,
848                    workspace_storage,
849                },
850                TargetTemplate::SshDocker { ssh, .. } => TargetLocator::SshDocker {
851                    borrowed_from: None,
852                    host: ssh.host.clone(),
853                    container_id: generated,
854                },
855                _ => return None,
856            };
857            Some(RecoveryCandidate {
858                session_id,
859                target_template_id: target_id.to_owned(),
860                locator,
861                ownership: None,
862                instance_id,
863                tracked_session: None,
864            })
865        })
866        .collect())
867}
868
869/// Managed session IDs in a container listing, each with the instance label
870/// that created it when the container carries one.
871pub(super) fn managed_sessions_from_container_json(
872    stdout: &[u8],
873) -> Result<Vec<(String, Option<String>)>> {
874    let values = serde_json::Deserializer::from_slice(stdout)
875        .into_iter::<serde_json::Value>()
876        .collect::<std::result::Result<Vec<_>, _>>()
877        .context("parse container list JSON")?;
878    let mut sessions = Vec::new();
879    for value in &values {
880        collect_managed_sessions(value, &mut sessions);
881    }
882    sessions.sort();
883    sessions.dedup();
884    Ok(sessions)
885}
886
887pub(super) fn collect_managed_sessions(
888    value: &serde_json::Value,
889    sessions: &mut Vec<(String, Option<String>)>,
890) {
891    match value {
892        serde_json::Value::Array(values) => {
893            for value in values {
894                collect_managed_sessions(value, sessions);
895            }
896        }
897        serde_json::Value::Object(object) => {
898            for label_key in ["Labels", "labels"] {
899                if let Some(labels) = object.get(label_key) {
900                    let managed = label_value(labels, targets::MANAGED_LABEL)
901                        .is_some_and(|value| value == "true");
902                    if managed && let Some(session) = label_value(labels, targets::SESSION_LABEL) {
903                        sessions.push((session, label_value(labels, targets::INSTANCE_LABEL)));
904                    }
905                }
906            }
907            for value in object.values() {
908                collect_managed_sessions(value, sessions);
909            }
910        }
911        _ => {}
912    }
913}
914
915fn label_value(labels: &serde_json::Value, key: &str) -> Option<String> {
916    match labels {
917        serde_json::Value::Object(object) => object.get(key)?.as_str().map(str::to_owned),
918        serde_json::Value::String(text) => text
919            .split(',')
920            .find_map(|label| {
921                label
922                    .trim()
923                    .split_once('=')
924                    .filter(|(name, _)| *name == key)
925            })
926            .map(|(_, value)| value.to_owned()),
927        _ => None,
928    }
929}
930
931fn candidates_from_aws_json(
932    target_id: &str,
933    address_source: AwsAddressSource,
934    stdout: &[u8],
935) -> Result<Vec<RecoveryCandidate>> {
936    let value: serde_json::Value =
937        serde_json::from_slice(stdout).context("parse AWS instance JSON")?;
938    let mut result = Vec::new();
939    let reservations = value
940        .get("Reservations")
941        .and_then(serde_json::Value::as_array)
942        .cloned()
943        .unwrap_or_default();
944    for instance in reservations.iter().flat_map(|reservation| {
945        reservation
946            .get("Instances")
947            .and_then(serde_json::Value::as_array)
948            .into_iter()
949            .flatten()
950    }) {
951        let tags = instance
952            .get("Tags")
953            .and_then(serde_json::Value::as_array)
954            .cloned()
955            .unwrap_or_default();
956        let tag = |key: &str| {
957            tags.iter()
958                .find(|tag| tag.get("Key").and_then(serde_json::Value::as_str) == Some(key))
959                .and_then(|tag| tag.get("Value"))
960                .and_then(serde_json::Value::as_str)
961        };
962        if tag(targets::MANAGED_TAG) != Some("true") {
963            continue;
964        }
965        let Some(session_id) = tag(targets::SESSION_TAG).map(str::to_owned) else {
966            continue;
967        };
968        targets::resource_name(&session_id)?;
969        let instance_id = instance
970            .get("InstanceId")
971            .and_then(serde_json::Value::as_str)
972            .context("managed EC2 instance omitted InstanceId")?
973            .to_owned();
974        let field = match address_source {
975            AwsAddressSource::PublicDns => "PublicDnsName",
976            AwsAddressSource::PublicIp => "PublicIpAddress",
977            AwsAddressSource::PrivateDns => "PrivateDnsName",
978            AwsAddressSource::PrivateIp => "PrivateIpAddress",
979        };
980        let address = instance
981            .get(field)
982            .and_then(serde_json::Value::as_str)
983            .filter(|value| !value.is_empty())
984            .map(str::to_owned);
985        let created_by = tag(targets::INSTANCE_TAG).map(str::to_owned);
986        result.push(RecoveryCandidate {
987            session_id,
988            target_template_id: target_id.to_owned(),
989            locator: TargetLocator::AwsEc2 {
990                instance_id,
991                address,
992            },
993            ownership: None,
994            instance_id: created_by,
995            tracked_session: None,
996        });
997    }
998    Ok(result)
999}
1000
1001fn execute_scan(
1002    executor: &impl CommandExecutor,
1003    command: CommandSpec,
1004    operation: &str,
1005) -> Result<CommandOutput> {
1006    let output = executor.execute(&command)?;
1007    if output.status != 0 {
1008        bail!(
1009            "{operation} failed with status {}: {}",
1010            output.status,
1011            String::from_utf8_lossy(&output.stderr).trim()
1012        );
1013    }
1014    Ok(output)
1015}
1016
1017fn ssh_spec(ssh: &SshConnection, remote: impl IntoIterator<Item = String>) -> CommandSpec {
1018    let backend = SshTarget::from(ssh);
1019    let mut args = backend.ssh_args.clone();
1020    args.push(backend.destination.clone());
1021    args.extend(remote);
1022    CommandSpec::new("ssh", args).ssh_session(&backend)
1023}
1024
1025fn read_recovery_ownership(
1026    template: &TargetTemplate,
1027    candidate: &RecoveryCandidate,
1028    executor: &impl CommandExecutor,
1029) -> Option<WorkerOwnership> {
1030    let backend =
1031        match recovery_backend_locator(template, &candidate.locator, &candidate.session_id) {
1032            Ok(backend) => backend,
1033            Err(error) => {
1034                tracing::debug!(
1035                    session_id = %candidate.session_id,
1036                    %error,
1037                    "could not construct a recovery ownership probe"
1038                );
1039                return None;
1040            }
1041        };
1042    let root = match targets::worker_root(&backend, &candidate.session_id) {
1043        Ok(root) => root,
1044        Err(error) => {
1045            tracing::debug!(
1046                session_id = %candidate.session_id,
1047                %error,
1048                "could not derive a recovery worker root"
1049            );
1050            return None;
1051        }
1052    };
1053    let command = match targets::command_on_locator(
1054        &backend,
1055        &candidate.session_id,
1056        vec!["cat".into(), format!("{root}/ownership.json")],
1057        "read worker ownership marker",
1058    ) {
1059        Ok(command) => command,
1060        Err(error) => {
1061            tracing::debug!(
1062                session_id = %candidate.session_id,
1063                %error,
1064                "could not construct a recovery ownership command"
1065            );
1066            return None;
1067        }
1068    };
1069    let output = match executor.execute(&command) {
1070        Ok(output) => output,
1071        Err(error) => {
1072            tracing::debug!(
1073                session_id = %candidate.session_id,
1074                %error,
1075                "could not read a recovery worker ownership marker"
1076            );
1077            return None;
1078        }
1079    };
1080    if output.status != 0 {
1081        tracing::debug!(
1082            session_id = %candidate.session_id,
1083            status = output.status,
1084            "recovery worker ownership probe returned a failure"
1085        );
1086        return None;
1087    }
1088    let marker: WorkerOwnership = match serde_json::from_slice(&output.stdout) {
1089        Ok(marker) => marker,
1090        Err(error) => {
1091            tracing::debug!(
1092                session_id = %candidate.session_id,
1093                %error,
1094                "recovery worker ownership marker was not valid JSON"
1095            );
1096            return None;
1097        }
1098    };
1099    if !(1..=WorkerOwnership::VERSION).contains(&marker.version)
1100        || marker.session_id != candidate.session_id
1101        || marker.target_template_id != candidate.target_template_id
1102    {
1103        tracing::debug!(
1104            session_id = %candidate.session_id,
1105            marker_session_id = %marker.session_id,
1106            marker_target_template_id = %marker.target_template_id,
1107            "recovery worker ownership marker did not match the candidate"
1108        );
1109        return None;
1110    }
1111    Some(marker)
1112}
1113
1114/// The backend locator recovery uses to destroy an orphan worker.
1115///
1116/// This stays separate from `backend_locator` (and from the shared
1117/// `TryFrom<StoredTarget>` conversion) because destroying an orphan has to
1118/// succeed in two cases the shared conversion rejects or cannot answer:
1119///
1120/// - The AWS arm substitutes `unavailable.invalid` for a missing address on
1121///   purpose. AWS cleanup terminates the instance through the `aws` CLI
1122///   (`targets/cleanup.rs`), so an instance with no address must still be
1123///   destroyable, whereas the shared conversion errors on the missing address.
1124/// - `worker_id: None` for `SshBare` is correct here. The scan takes the
1125///   session id from the worker directory name and `targets::worker_root`
1126///   falls back to the session id, so the root resolves. A borrowed target's
1127///   shared parent workspace is not discoverable from a scan, and recomputing
1128///   the child's own (nonexistent) path is the safe direction for destroy.
1129fn recovery_backend_locator(
1130    template: &TargetTemplate,
1131    locator: &TargetLocator,
1132    session_id: &str,
1133) -> Result<targets::TargetLocator> {
1134    Ok(match (template, locator) {
1135        (TargetTemplate::LocalBare, TargetLocator::LocalBare { worker_root }) => {
1136            targets::TargetLocator::LocalBare {
1137                worker_root: worker_root.to_string_lossy().into_owned(),
1138            }
1139        }
1140        (
1141            TargetTemplate::LocalPodman { .. },
1142            TargetLocator::LocalPodman {
1143                container_id,
1144                workspace_storage,
1145                ..
1146            },
1147        ) => targets::TargetLocator::LocalPodman {
1148            borrowed_from: None,
1149            container_id: container_id.clone(),
1150            workspace_storage: workspace_storage.into(),
1151        },
1152        (TargetTemplate::LocalDocker { .. }, TargetLocator::LocalDocker { container_id, .. }) => {
1153            targets::TargetLocator::LocalDocker {
1154                borrowed_from: None,
1155                container_id: container_id.clone(),
1156            }
1157        }
1158        (
1159            TargetTemplate::AppleContainer { .. },
1160            TargetLocator::AppleContainer { container_id, .. },
1161        ) => targets::TargetLocator::AppleContainer {
1162            borrowed_from: None,
1163            container_id: container_id.clone(),
1164        },
1165        (
1166            TargetTemplate::SshPodman { ssh, .. },
1167            TargetLocator::SshPodman {
1168                container_id,
1169                workspace_storage,
1170                ..
1171            },
1172        ) => targets::TargetLocator::SshPodman {
1173            borrowed_from: None,
1174            ssh: SshTarget::from(ssh),
1175            container_id: container_id.clone(),
1176            workspace_storage: workspace_storage.into(),
1177        },
1178        (
1179            TargetTemplate::SshDocker { ssh, .. },
1180            TargetLocator::SshDocker {
1181                host, container_id, ..
1182            },
1183        ) => {
1184            if host != &ssh.host {
1185                bail!("recovery SSH Docker host does not match target template")
1186            }
1187            targets::TargetLocator::SshDocker {
1188                borrowed_from: None,
1189                ssh: SshTarget::from(ssh),
1190                container_id: container_id.clone(),
1191            }
1192        }
1193        (TargetTemplate::SshBare { ssh, .. }, TargetLocator::SshBare { workspace, .. }) => {
1194            targets::TargetLocator::SshBare {
1195                ssh: SshTarget::from(ssh),
1196                workspace: workspace.to_string_lossy().into_owned(),
1197                worker_id: None,
1198            }
1199        }
1200        (
1201            TargetTemplate::AwsEc2 {
1202                aws_profile,
1203                region,
1204                ssh_user,
1205                identity_file,
1206                ssh_args,
1207                ..
1208            },
1209            TargetLocator::AwsEc2 {
1210                instance_id,
1211                address,
1212            },
1213        ) => targets::TargetLocator::AwsEc2 {
1214            profile: aws_profile.clone().unwrap_or_else(|| "default".into()),
1215            region: region.clone(),
1216            instance_id: instance_id.clone(),
1217            ssh: SshTarget {
1218                destination: format!(
1219                    "{ssh_user}@{}",
1220                    address.as_deref().unwrap_or("unavailable.invalid")
1221                ),
1222                ssh_args: targets::ssh_args_with_identity(ssh_args, identity_file.as_deref()),
1223            },
1224            workspace: format!(".local/share/hel/workspaces/{session_id}"),
1225        },
1226        _ => bail!("recovery target locator does not match target template"),
1227    })
1228}
1229
1230#[cfg(test)]
1231mod tests {
1232
1233    use crate::controller::test_support::{IsolatedTest, test_name};
1234    use mj_core::config::{
1235        AwsAddressSource, Config, ContainerTemplate as ConfigContainer, HarnessKind,
1236        PodmanWorkspaceStorage, TargetTemplate,
1237    };
1238    use mj_core::state::{State, TargetLocator};
1239
1240    use crate::targets::ProcessExecutor;
1241
1242    use super::*;
1243
1244    const FAILED_ADOPTION_CHILD: &str = "MJ_TEST_FAILED_ADOPTION_CHILD";
1245
1246    #[tokio::test]
1247    async fn a_failed_adoption_records_the_failure_and_stays_retryable() {
1248        // MJ_DATA_DIR is process-global, so the database-backed half runs in
1249        // an exact child test with its own data directory.
1250        if std::env::var_os(FAILED_ADOPTION_CHILD).is_none() {
1251            let directory = tempfile::tempdir().unwrap();
1252            IsolatedTest::new(test_name(
1253                module_path!(),
1254                "a_failed_adoption_records_the_failure_and_stays_retryable",
1255            ))
1256            .env(FAILED_ADOPTION_CHILD, "1")
1257            .env("MJ_DATA_DIR", directory.path())
1258            .run();
1259            return;
1260        }
1261        // Alone in this child process, so it installs the one writer.
1262        let _writer = crate::database::install_isolated_test_writer();
1263
1264        let session_id = "0123456789abcdef0123456789abcdef";
1265        let workers = tempfile::tempdir().unwrap();
1266        // Exactly what an adoption commits before its handshake, on a worker
1267        // root that holds no worker binary, so the handshake cannot succeed.
1268        let record = adopted_session_record(
1269            session_id,
1270            "local-bare",
1271            "codex".into(),
1272            HarnessKind::Codex,
1273            "project".into(),
1274            mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1275            TargetLocator::LocalBare {
1276                worker_root: workers.path().join(session_id),
1277            },
1278        );
1279        assert!(
1280            adoption_unfinished(&record, "local-bare"),
1281            "the record adoption commits must be the record adoption can retry"
1282        );
1283        crate::database::save_session(&record).unwrap();
1284        let mut config = Config::default();
1285        config
1286            .targets
1287            .insert("local-bare".into(), TargetTemplate::LocalBare);
1288        let mut state = State::default();
1289        state.sessions.insert(session_id.to_owned(), record);
1290        let mut controller = Controller { config, state };
1291
1292        let failure = controller
1293            .adopt_orphan_worker(
1294                session_id,
1295                "local-bare",
1296                None,
1297                None,
1298                false,
1299                &ProcessExecutor,
1300            )
1301            .await
1302            .expect_err("a worker root without a worker cannot complete the handshake");
1303        assert!(
1304            format!("{failure:#}").contains("orphan relay"),
1305            "unexpected failure: {failure:#}"
1306        );
1307        let recorded = controller.state.sessions[session_id]
1308            .last_error
1309            .clone()
1310            .expect("the failed handshake was recorded on the session");
1311        assert!(
1312            recorded.contains("orphan adoption failed"),
1313            "unexpected recorded failure: {recorded}"
1314        );
1315        assert_eq!(
1316            controller.state.sessions[session_id].state,
1317            SessionState::Disconnected
1318        );
1319        let stored = crate::database::load_state().unwrap();
1320        assert_eq!(
1321            stored.sessions[session_id].last_error.as_deref(),
1322            Some(recorded.as_str()),
1323            "the adoption failure was not persisted"
1324        );
1325
1326        let retry = controller
1327            .adopt_orphan_worker(
1328                session_id,
1329                "local-bare",
1330                None,
1331                None,
1332                false,
1333                &ProcessExecutor,
1334            )
1335            .await
1336            .expect_err("the worker is still unreachable");
1337        let retry = format!("{retry:#}");
1338        assert!(
1339            retry.contains("orphan relay"),
1340            "adoption did not retry the handshake: {retry}"
1341        );
1342        assert!(
1343            !retry.contains("already tracked"),
1344            "a session adoption never finished blocked its own retry: {retry}"
1345        );
1346    }
1347
1348    const RECOVERY_WORKSPACE_CHILD: &str = "MJ_TEST_RECOVERY_WORKSPACE_CHILD";
1349
1350    #[tokio::test]
1351    async fn orphan_workspace_ids_are_reconciled_before_adoption_persistence() {
1352        if std::env::var_os(RECOVERY_WORKSPACE_CHILD).is_none() {
1353            let directory = tempfile::tempdir().unwrap();
1354            IsolatedTest::new(test_name(
1355                module_path!(),
1356                "orphan_workspace_ids_are_reconciled_before_adoption_persistence",
1357            ))
1358            .env(RECOVERY_WORKSPACE_CHILD, "1")
1359            .env("MJ_DATA_DIR", directory.path())
1360            .run();
1361            return;
1362        }
1363
1364        let _writer = crate::database::install_isolated_test_writer();
1365        let default =
1366            resolve_recovery_workspace_id(mj_core::workspace::DEFAULT_WORKSPACE_ID).unwrap();
1367        assert_eq!(default, mj_core::workspace::DEFAULT_WORKSPACE_ID);
1368
1369        let known = crate::database::create_or_get_workspace("Known").unwrap();
1370        assert_eq!(resolve_recovery_workspace_id(&known.id).unwrap(), known.id);
1371
1372        let recovered = resolve_recovery_workspace_id("workspace-from-old-controller").unwrap();
1373        let repeated = resolve_recovery_workspace_id("another-old-workspace").unwrap();
1374        assert_eq!(repeated, recovered);
1375        assert_eq!(
1376            crate::database::list_workspaces()
1377                .unwrap()
1378                .iter()
1379                .filter(|workspace| workspace.name == "Recovered")
1380                .count(),
1381            1
1382        );
1383
1384        let session_id = "0123456789abcdef0123456789abcdef";
1385        let workers = tempfile::tempdir().unwrap();
1386        let record = adopted_session_record(
1387            session_id,
1388            "local-bare",
1389            "codex".into(),
1390            HarnessKind::Codex,
1391            "project".into(),
1392            recovered.clone(),
1393            TargetLocator::LocalBare {
1394                worker_root: workers.path().join(session_id),
1395            },
1396        );
1397        crate::database::save_session(&record).unwrap();
1398        assert_eq!(
1399            crate::database::load_state().unwrap().sessions[session_id].workspace_id,
1400            recovered
1401        );
1402
1403        // Adoption has already committed the reconciled workspace. Exercise
1404        // the post-save handshake path with the same fake used by the retry
1405        // regression, which leaves and persists a diagnostic on failure.
1406        let mut state = State::default();
1407        state.sessions.insert(session_id.to_owned(), record);
1408        let mut config = Config::default();
1409        config
1410            .targets
1411            .insert("local-bare".into(), TargetTemplate::LocalBare);
1412        let mut controller = Controller { config, state };
1413        let failure = controller
1414            .adopt_orphan_worker(
1415                session_id,
1416                "local-bare",
1417                None,
1418                None,
1419                false,
1420                &ProcessExecutor,
1421            )
1422            .await
1423            .expect_err("a worker root without a worker cannot complete the handshake");
1424        assert!(
1425            format!("{failure:#}").contains("orphan relay"),
1426            "unexpected failure: {failure:#}"
1427        );
1428        let stored = crate::database::load_state().unwrap();
1429        assert_eq!(stored.sessions[session_id].workspace_id, recovered);
1430        assert!(
1431            stored.sessions[session_id]
1432                .last_error
1433                .as_deref()
1434                .is_some_and(|error| error.contains("orphan adoption failed"))
1435        );
1436    }
1437
1438    #[test]
1439    fn a_session_that_completed_its_handshake_is_not_adoptable_again() {
1440        let mut record = adopted_session_record(
1441            "0123456789abcdef0123456789abcdef",
1442            "local-bare",
1443            "codex".into(),
1444            HarnessKind::Codex,
1445            "project".into(),
1446            mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1447            TargetLocator::LocalBare {
1448                worker_root: std::path::PathBuf::from("/workers/0123456789abcdef0123456789abcdef"),
1449            },
1450        );
1451        record.native_session_id = Some("native-session".into());
1452        assert!(!adoption_unfinished(&record, "local-bare"));
1453
1454        record.native_session_id = None;
1455        assert!(
1456            !adoption_unfinished(&record, "other-target"),
1457            "a record adopted onto another target is not this target's retry"
1458        );
1459    }
1460
1461    #[test]
1462    fn recovery_container_scan_requires_both_managed_and_session_labels() {
1463        let template = TargetTemplate::LocalPodman {
1464            container: ConfigContainer {
1465                build_cache: None,
1466                image: "ignored".into(),
1467                pull_policy: Default::default(),
1468                platform: None,
1469                cpus: None,
1470                memory: None,
1471                environment: Default::default(),
1472                workspace_storage: Default::default(),
1473            },
1474        };
1475        let json = serde_json::json!([
1476            {"Labels": {"dev.mj.managed": "true", "dev.mj.session": "0123456789abcdef0123456789abcdef", "dev.mj.instance": "qa0916"}},
1477            {"Labels": {"dev.mj.managed": "false", "dev.mj.session": "not-owned"}},
1478            {"configuration": {"labels": "dev.mj.managed=true,dev.mj.session=abcdef0123456789abcdef0123456789"}}
1479        ]);
1480        let candidates = candidates_from_container_json(
1481            "local",
1482            &template,
1483            serde_json::to_string(&json).unwrap().as_bytes(),
1484        )
1485        .unwrap();
1486        assert_eq!(candidates.len(), 2);
1487        assert_eq!(candidates[0].session_id, "0123456789abcdef0123456789abcdef");
1488        assert_eq!(candidates[0].instance_id.as_deref(), Some("qa0916"));
1489        assert_eq!(
1490            candidates[1].instance_id, None,
1491            "a container without an instance label is of unknown origin"
1492        );
1493    }
1494
1495    /// Answers the container listing the scan runs and fails everything else,
1496    /// which is what a container with no worker in it does to the ownership
1497    /// probe.
1498    struct ListingExecutor {
1499        listing: String,
1500    }
1501
1502    impl CommandExecutor for ListingExecutor {
1503        fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1504            if command.program == "podman" && command.args.first().map(String::as_str) == Some("ps")
1505            {
1506                return Ok(CommandOutput {
1507                    status: 0,
1508                    stdout: self.listing.clone().into_bytes(),
1509                    stderr: Vec::new(),
1510                });
1511            }
1512            Ok(CommandOutput {
1513                status: 1,
1514                stdout: Vec::new(),
1515                stderr: Vec::new(),
1516            })
1517        }
1518    }
1519
1520    fn podman_scan_controller(record: SessionRecord) -> (Controller, ListingExecutor) {
1521        let mut config = Config::default();
1522        config.targets.insert(
1523            "local".into(),
1524            TargetTemplate::LocalPodman {
1525                container: ConfigContainer {
1526                    build_cache: None,
1527                    image: "ignored".into(),
1528                    pull_policy: Default::default(),
1529                    platform: None,
1530                    cpus: None,
1531                    memory: None,
1532                    environment: Default::default(),
1533                    workspace_storage: Default::default(),
1534                },
1535            },
1536        );
1537        let listing = serde_json::to_string(&serde_json::json!([
1538            {"Labels": {"dev.mj.managed": "true", "dev.mj.session": record.id}}
1539        ]))
1540        .unwrap();
1541        let mut state = State::default();
1542        state.sessions.insert(record.id.clone(), record);
1543        (Controller { config, state }, ListingExecutor { listing })
1544    }
1545
1546    fn errored_podman_session(session_id: &str) -> SessionRecord {
1547        let mut record = adopted_session_record(
1548            session_id,
1549            "local",
1550            "codex".into(),
1551            HarnessKind::Codex,
1552            "project".into(),
1553            mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1554            TargetLocator::LocalPodman {
1555                container_id: targets::resource_name(session_id).unwrap(),
1556                workspace_storage: PodmanWorkspaceLocator::ContainerLayer,
1557                borrowed_from: None,
1558            },
1559        );
1560        // What a failed provision or a failed move writes: the session is
1561        // remembered, the locator that could tear its container down is not.
1562        record.state = SessionState::Error;
1563        record.target = None;
1564        record
1565    }
1566
1567    #[test]
1568    fn a_container_an_errored_session_no_longer_names_is_an_orphan() {
1569        let session_id = "0123456789abcdef0123456789abcdef";
1570        let (controller, executor) = podman_scan_controller(errored_podman_session(session_id));
1571
1572        let scan = controller.scan_orphan_workers(&executor, true);
1573
1574        let [candidate] = scan.candidates.as_slice() else {
1575            panic!("the errored session's container was not offered: {scan:?}");
1576        };
1577        assert_eq!(candidate.session_id, session_id);
1578        assert_eq!(
1579            candidate.tracked_session,
1580            Some(SessionState::Error),
1581            "recovery must say the session id is still taken, so only destroy applies"
1582        );
1583    }
1584
1585    #[test]
1586    fn a_container_a_session_still_names_is_not_an_orphan() {
1587        let session_id = "0123456789abcdef0123456789abcdef";
1588        let mut record = errored_podman_session(session_id);
1589        // The same failure, except the record kept its locator: the session's
1590        // own forced destroy still reaches the container, so recovery stays out.
1591        record.target = Some(TargetLocator::LocalPodman {
1592            container_id: targets::resource_name(session_id).unwrap(),
1593            workspace_storage: PodmanWorkspaceLocator::ContainerLayer,
1594            borrowed_from: None,
1595        });
1596        let (controller, executor) = podman_scan_controller(record);
1597
1598        let scan = controller.scan_orphan_workers(&executor, true);
1599
1600        assert!(
1601            scan.candidates.is_empty(),
1602            "a container the controller can still drive was offered for destruction: {scan:?}"
1603        );
1604    }
1605
1606    #[test]
1607    fn a_container_being_provisioned_is_not_an_orphan() {
1608        let session_id = "0123456789abcdef0123456789abcdef";
1609        let mut record = errored_podman_session(session_id);
1610        // Provisioning writes the locator as it runs, so a scan that catches
1611        // the window between `podman run` and locator discovery must not
1612        // mistake the new container for a leftover.
1613        record.state = SessionState::Provisioning;
1614        let (controller, executor) = podman_scan_controller(record);
1615
1616        let scan = controller.scan_orphan_workers(&executor, true);
1617
1618        assert!(
1619            scan.candidates.is_empty(),
1620            "a container still being provisioned was offered for destruction: {scan:?}"
1621        );
1622    }
1623
1624    fn candidate(session_id: &str, instance_id: Option<&str>) -> RecoveryCandidate {
1625        RecoveryCandidate {
1626            session_id: session_id.to_owned(),
1627            target_template_id: "local".to_owned(),
1628            locator: TargetLocator::LocalDocker {
1629                container_id: format!("mj-{session_id}"),
1630                borrowed_from: None,
1631            },
1632            ownership: None,
1633            instance_id: instance_id.map(str::to_owned),
1634            tracked_session: None,
1635        }
1636    }
1637
1638    #[test]
1639    fn default_scan_scope_hides_other_and_unknown_instances() {
1640        let mut scan = RecoveryScan {
1641            candidates: vec![
1642                candidate("mine", Some("qa")),
1643                candidate("theirs", Some("prod")),
1644                candidate("legacy", None),
1645            ],
1646            instance_id: "qa".to_owned(),
1647            ..RecoveryScan::default()
1648        };
1649        restrict_to_instance(&mut scan);
1650        assert_eq!(
1651            scan.candidates
1652                .iter()
1653                .map(|candidate| candidate.session_id.as_str())
1654                .collect::<Vec<_>>(),
1655            ["mine"]
1656        );
1657        assert_eq!(scan.hidden_other_instances, 2);
1658    }
1659
1660    #[test]
1661    fn acting_on_another_or_unknown_instance_requires_the_explicit_flag() {
1662        require_instance_access(&candidate("mine", Some("qa")), "qa", false).unwrap();
1663
1664        let other = require_instance_access(&candidate("theirs", Some("prod")), "qa", false)
1665            .expect_err("another instance's worker is refused by default");
1666        assert!(
1667            other.to_string().contains("belongs to instance \"prod\""),
1668            "{other}"
1669        );
1670        require_instance_access(&candidate("theirs", Some("prod")), "qa", true).unwrap();
1671
1672        let unknown = require_instance_access(&candidate("legacy", None), "qa", false)
1673            .expect_err("a worker without a stamp is refused by default");
1674        assert!(
1675            unknown.to_string().contains("no instance stamp"),
1676            "{unknown}"
1677        );
1678        require_instance_access(&candidate("legacy", None), "qa", true).unwrap();
1679    }
1680
1681    #[test]
1682    fn a_podman_orphan_destroy_plan_removes_its_workspace_volume() {
1683        // `destroy_orphan_worker` rescans, so this drives the two steps it runs
1684        // on the candidate it finds: the scan's locator and the close plan
1685        // built from it. Those are what carried the container-layer default.
1686        let session = "0123456789abcdef0123456789abcdef";
1687        let template = TargetTemplate::LocalPodman {
1688            container: ConfigContainer {
1689                image: "ignored".into(),
1690                pull_policy: Default::default(),
1691                platform: None,
1692                cpus: None,
1693                memory: None,
1694                environment: Default::default(),
1695                workspace_storage: PodmanWorkspaceStorage::PodmanVolume,
1696                build_cache: None,
1697            },
1698        };
1699        let json = serde_json::json!([
1700            {"Labels": {"dev.mj.managed": "true", "dev.mj.session": session}}
1701        ]);
1702
1703        let candidates = candidates_from_container_json(
1704            "local",
1705            &template,
1706            serde_json::to_string(&json).unwrap().as_bytes(),
1707        )
1708        .unwrap();
1709
1710        let [candidate] = candidates.as_slice() else {
1711            panic!("expected one orphan candidate, got {candidates:?}");
1712        };
1713        let volume = format!("{}-workspace", targets::resource_name(session).unwrap());
1714        assert!(
1715            matches!(
1716                &candidate.locator,
1717                TargetLocator::LocalPodman {
1718                    workspace_storage: PodmanWorkspaceLocator::Volume { name },
1719                    ..
1720                } if name == &volume
1721            ),
1722            "candidate locator lost the volume storage: {:?}",
1723            candidate.locator
1724        );
1725        let backend = recovery_backend_locator(&template, &candidate.locator, session).unwrap();
1726        let plan = targets::close_plan(&backend, session).unwrap();
1727        assert!(
1728            plan.commands.iter().any(|command| {
1729                command.args.iter().any(|argument| argument == &volume)
1730                    && command
1731                        .args
1732                        .iter()
1733                        .any(|argument| argument.contains("podman volume rm"))
1734            }),
1735            "destroy plan does not remove the workspace volume: {plan:?}"
1736        );
1737    }
1738
1739    #[test]
1740    fn recovery_docker_scan_accepts_json_lines_and_builds_a_docker_locator() {
1741        let template = TargetTemplate::LocalDocker {
1742            container: ConfigContainer {
1743                build_cache: None,
1744                image: "ignored".into(),
1745                pull_policy: Default::default(),
1746                platform: None,
1747                cpus: None,
1748                memory: None,
1749                environment: Default::default(),
1750                workspace_storage: Default::default(),
1751            },
1752        };
1753        let session = "0123456789abcdef0123456789abcdef";
1754        let output = format!(
1755            "{{\"Labels\":\"dev.mj.managed=true,dev.mj.session={session},dev.mj.instance=abc123\"}}\n{{\"Labels\":\"dev.mj.managed=false,dev.mj.session=ignored\"}}\n"
1756        );
1757
1758        let candidates =
1759            candidates_from_container_json("docker", &template, output.as_bytes()).unwrap();
1760
1761        assert_eq!(candidates.len(), 1);
1762        assert_eq!(candidates[0].session_id, session);
1763        assert_eq!(candidates[0].instance_id.as_deref(), Some("abc123"));
1764        assert!(matches!(
1765            &candidates[0].locator,
1766            TargetLocator::LocalDocker { container_id, .. }
1767                if container_id == &targets::resource_name(session).unwrap()
1768        ));
1769    }
1770
1771    #[test]
1772    fn recovery_aws_scan_uses_exact_tagged_instance_and_address() {
1773        let json = serde_json::json!({"Reservations": [{"Instances": [{
1774            "InstanceId": "i-exact",
1775            "PrivateIpAddress": "10.0.0.7",
1776            "Tags": [
1777                {"Key": "dev.mj.managed", "Value": "true"},
1778                {"Key": "dev.mj.session", "Value": "0123456789abcdef0123456789abcdef"},
1779                {"Key": "dev.mj.instance", "Value": "qa0916"}
1780            ]
1781        }]}]});
1782        let candidates = candidates_from_aws_json(
1783            "aws",
1784            AwsAddressSource::PrivateIp,
1785            serde_json::to_string(&json).unwrap().as_bytes(),
1786        )
1787        .unwrap();
1788        assert_eq!(candidates[0].instance_id.as_deref(), Some("qa0916"));
1789        assert!(matches!(
1790            &candidates[0].locator,
1791            TargetLocator::AwsEc2 { instance_id, address }
1792                if instance_id == "i-exact" && address.as_deref() == Some("10.0.0.7")
1793        ));
1794    }
1795}