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        mjolnir_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    use std::collections::BTreeMap;
1233
1234    use crate::controller::test_support::{IsolatedTest, test_name};
1235    use mj_core::config::{
1236        AwsAddressSource, Config, ContainerTemplate as ConfigContainer, HarnessKind,
1237        PodmanWorkspaceStorage, TargetTemplate,
1238    };
1239    use mj_core::state::{State, TargetLocator};
1240
1241    use crate::targets::ProcessExecutor;
1242
1243    use super::*;
1244
1245    const FAILED_ADOPTION_CHILD: &str = "MJ_TEST_FAILED_ADOPTION_CHILD";
1246
1247    #[tokio::test]
1248    async fn a_failed_adoption_records_the_failure_and_stays_retryable() {
1249        // MJ_DATA_DIR is process-global, so the database-backed half runs in
1250        // an exact child test with its own data directory.
1251        if std::env::var_os(FAILED_ADOPTION_CHILD).is_none() {
1252            let directory = tempfile::tempdir().unwrap();
1253            IsolatedTest::new(test_name(
1254                module_path!(),
1255                "a_failed_adoption_records_the_failure_and_stays_retryable",
1256            ))
1257            .env(FAILED_ADOPTION_CHILD, "1")
1258            .env("MJ_DATA_DIR", directory.path())
1259            .run();
1260            return;
1261        }
1262        // Alone in this child process, so it installs the one writer.
1263        let _writer = crate::database::install_isolated_test_writer();
1264
1265        let session_id = "0123456789abcdef0123456789abcdef";
1266        let workers = tempfile::tempdir().unwrap();
1267        // Exactly what an adoption commits before its handshake, on a worker
1268        // root that holds no worker binary, so the handshake cannot succeed.
1269        let record = adopted_session_record(
1270            session_id,
1271            "local-bare",
1272            "codex".into(),
1273            HarnessKind::Codex,
1274            "project".into(),
1275            mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1276            TargetLocator::LocalBare {
1277                worker_root: workers.path().join(session_id),
1278            },
1279        );
1280        assert!(
1281            adoption_unfinished(&record, "local-bare"),
1282            "the record adoption commits must be the record adoption can retry"
1283        );
1284        crate::database::save_session(&record).unwrap();
1285        let mut config = Config::default();
1286        config
1287            .targets
1288            .insert("local-bare".into(), TargetTemplate::LocalBare);
1289        let mut state = State::default();
1290        state.sessions.insert(session_id.to_owned(), record);
1291        let mut controller = Controller { config, state };
1292
1293        let failure = controller
1294            .adopt_orphan_worker(
1295                session_id,
1296                "local-bare",
1297                None,
1298                None,
1299                false,
1300                &ProcessExecutor,
1301            )
1302            .await
1303            .expect_err("a worker root without a worker cannot complete the handshake");
1304        assert!(
1305            format!("{failure:#}").contains("orphan relay"),
1306            "unexpected failure: {failure:#}"
1307        );
1308        let recorded = controller.state.sessions[session_id]
1309            .last_error
1310            .clone()
1311            .expect("the failed handshake was recorded on the session");
1312        assert!(
1313            recorded.contains("orphan adoption failed"),
1314            "unexpected recorded failure: {recorded}"
1315        );
1316        assert_eq!(
1317            controller.state.sessions[session_id].state,
1318            SessionState::Disconnected
1319        );
1320        let stored = crate::database::load_state().unwrap();
1321        assert_eq!(
1322            stored.sessions[session_id].last_error.as_deref(),
1323            Some(recorded.as_str()),
1324            "the adoption failure was not persisted"
1325        );
1326
1327        let retry = controller
1328            .adopt_orphan_worker(
1329                session_id,
1330                "local-bare",
1331                None,
1332                None,
1333                false,
1334                &ProcessExecutor,
1335            )
1336            .await
1337            .expect_err("the worker is still unreachable");
1338        let retry = format!("{retry:#}");
1339        assert!(
1340            retry.contains("orphan relay"),
1341            "adoption did not retry the handshake: {retry}"
1342        );
1343        assert!(
1344            !retry.contains("already tracked"),
1345            "a session adoption never finished blocked its own retry: {retry}"
1346        );
1347    }
1348
1349    const RECOVERY_WORKSPACE_CHILD: &str = "MJ_TEST_RECOVERY_WORKSPACE_CHILD";
1350
1351    #[tokio::test]
1352    async fn orphan_workspace_ids_are_reconciled_before_adoption_persistence() {
1353        if std::env::var_os(RECOVERY_WORKSPACE_CHILD).is_none() {
1354            let directory = tempfile::tempdir().unwrap();
1355            IsolatedTest::new(test_name(
1356                module_path!(),
1357                "orphan_workspace_ids_are_reconciled_before_adoption_persistence",
1358            ))
1359            .env(RECOVERY_WORKSPACE_CHILD, "1")
1360            .env("MJ_DATA_DIR", directory.path())
1361            .run();
1362            return;
1363        }
1364
1365        let _writer = crate::database::install_isolated_test_writer();
1366        let default =
1367            resolve_recovery_workspace_id(mj_core::workspace::DEFAULT_WORKSPACE_ID).unwrap();
1368        assert_eq!(default, mj_core::workspace::DEFAULT_WORKSPACE_ID);
1369
1370        let known = crate::database::create_or_get_workspace("Known").unwrap();
1371        assert_eq!(resolve_recovery_workspace_id(&known.id).unwrap(), known.id);
1372
1373        let recovered = resolve_recovery_workspace_id("workspace-from-old-controller").unwrap();
1374        let repeated = resolve_recovery_workspace_id("another-old-workspace").unwrap();
1375        assert_eq!(repeated, recovered);
1376        assert_eq!(
1377            crate::database::list_workspaces()
1378                .unwrap()
1379                .iter()
1380                .filter(|workspace| workspace.name == "Recovered")
1381                .count(),
1382            1
1383        );
1384
1385        let session_id = "0123456789abcdef0123456789abcdef";
1386        let workers = tempfile::tempdir().unwrap();
1387        let record = adopted_session_record(
1388            session_id,
1389            "local-bare",
1390            "codex".into(),
1391            HarnessKind::Codex,
1392            "project".into(),
1393            recovered.clone(),
1394            TargetLocator::LocalBare {
1395                worker_root: workers.path().join(session_id),
1396            },
1397        );
1398        crate::database::save_session(&record).unwrap();
1399        assert_eq!(
1400            crate::database::load_state().unwrap().sessions[session_id].workspace_id,
1401            recovered
1402        );
1403
1404        // Adoption has already committed the reconciled workspace. Exercise
1405        // the post-save handshake path with the same fake used by the retry
1406        // regression, which leaves and persists a diagnostic on failure.
1407        let mut state = State::default();
1408        state.sessions.insert(session_id.to_owned(), record);
1409        let mut config = Config::default();
1410        config
1411            .targets
1412            .insert("local-bare".into(), TargetTemplate::LocalBare);
1413        let mut controller = Controller { config, state };
1414        let failure = controller
1415            .adopt_orphan_worker(
1416                session_id,
1417                "local-bare",
1418                None,
1419                None,
1420                false,
1421                &ProcessExecutor,
1422            )
1423            .await
1424            .expect_err("a worker root without a worker cannot complete the handshake");
1425        assert!(
1426            format!("{failure:#}").contains("orphan relay"),
1427            "unexpected failure: {failure:#}"
1428        );
1429        let stored = crate::database::load_state().unwrap();
1430        assert_eq!(stored.sessions[session_id].workspace_id, recovered);
1431        assert!(
1432            stored.sessions[session_id]
1433                .last_error
1434                .as_deref()
1435                .is_some_and(|error| error.contains("orphan adoption failed"))
1436        );
1437    }
1438
1439    #[test]
1440    fn a_session_that_completed_its_handshake_is_not_adoptable_again() {
1441        let mut record = adopted_session_record(
1442            "0123456789abcdef0123456789abcdef",
1443            "local-bare",
1444            "codex".into(),
1445            HarnessKind::Codex,
1446            "project".into(),
1447            mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1448            TargetLocator::LocalBare {
1449                worker_root: std::path::PathBuf::from("/workers/0123456789abcdef0123456789abcdef"),
1450            },
1451        );
1452        record.native_session_id = Some("native-session".into());
1453        assert!(!adoption_unfinished(&record, "local-bare"));
1454
1455        record.native_session_id = None;
1456        assert!(
1457            !adoption_unfinished(&record, "other-target"),
1458            "a record adopted onto another target is not this target's retry"
1459        );
1460    }
1461
1462    #[test]
1463    fn recovery_container_scan_requires_both_managed_and_session_labels() {
1464        let template = TargetTemplate::LocalPodman {
1465            container: ConfigContainer {
1466                build_cache: None,
1467                image: "ignored".into(),
1468                pull_policy: Default::default(),
1469                platform: None,
1470                cpus: None,
1471                memory: None,
1472                environment: BTreeMap::new(),
1473                workspace_storage: Default::default(),
1474            },
1475        };
1476        let json = serde_json::json!([
1477            {"Labels": {"dev.mj.managed": "true", "dev.mj.session": "0123456789abcdef0123456789abcdef", "dev.mj.instance": "qa0916"}},
1478            {"Labels": {"dev.mj.managed": "false", "dev.mj.session": "not-owned"}},
1479            {"configuration": {"labels": "dev.mj.managed=true,dev.mj.session=abcdef0123456789abcdef0123456789"}}
1480        ]);
1481        let candidates = candidates_from_container_json(
1482            "local",
1483            &template,
1484            serde_json::to_string(&json).unwrap().as_bytes(),
1485        )
1486        .unwrap();
1487        assert_eq!(candidates.len(), 2);
1488        assert_eq!(candidates[0].session_id, "0123456789abcdef0123456789abcdef");
1489        assert_eq!(candidates[0].instance_id.as_deref(), Some("qa0916"));
1490        assert_eq!(
1491            candidates[1].instance_id, None,
1492            "a container without an instance label is of unknown origin"
1493        );
1494    }
1495
1496    /// Answers the container listing the scan runs and fails everything else,
1497    /// which is what a container with no worker in it does to the ownership
1498    /// probe.
1499    struct ListingExecutor {
1500        listing: String,
1501    }
1502
1503    impl CommandExecutor for ListingExecutor {
1504        fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
1505            if command.program == "podman" && command.args.first().map(String::as_str) == Some("ps")
1506            {
1507                return Ok(CommandOutput {
1508                    status: 0,
1509                    stdout: self.listing.clone().into_bytes(),
1510                    stderr: Vec::new(),
1511                });
1512            }
1513            Ok(CommandOutput {
1514                status: 1,
1515                stdout: Vec::new(),
1516                stderr: Vec::new(),
1517            })
1518        }
1519    }
1520
1521    fn podman_scan_controller(record: SessionRecord) -> (Controller, ListingExecutor) {
1522        let mut config = Config::default();
1523        config.targets.insert(
1524            "local".into(),
1525            TargetTemplate::LocalPodman {
1526                container: ConfigContainer {
1527                    build_cache: None,
1528                    image: "ignored".into(),
1529                    pull_policy: Default::default(),
1530                    platform: None,
1531                    cpus: None,
1532                    memory: None,
1533                    environment: BTreeMap::new(),
1534                    workspace_storage: Default::default(),
1535                },
1536            },
1537        );
1538        let listing = serde_json::to_string(&serde_json::json!([
1539            {"Labels": {"dev.mj.managed": "true", "dev.mj.session": record.id}}
1540        ]))
1541        .unwrap();
1542        let mut state = State::default();
1543        state.sessions.insert(record.id.clone(), record);
1544        (Controller { config, state }, ListingExecutor { listing })
1545    }
1546
1547    fn errored_podman_session(session_id: &str) -> SessionRecord {
1548        let mut record = adopted_session_record(
1549            session_id,
1550            "local",
1551            "codex".into(),
1552            HarnessKind::Codex,
1553            "project".into(),
1554            mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1555            TargetLocator::LocalPodman {
1556                container_id: targets::resource_name(session_id).unwrap(),
1557                workspace_storage: PodmanWorkspaceLocator::ContainerLayer,
1558                borrowed_from: None,
1559            },
1560        );
1561        // What a failed provision or a failed move writes: the session is
1562        // remembered, the locator that could tear its container down is not.
1563        record.state = SessionState::Error;
1564        record.target = None;
1565        record
1566    }
1567
1568    #[test]
1569    fn a_container_an_errored_session_no_longer_names_is_an_orphan() {
1570        let session_id = "0123456789abcdef0123456789abcdef";
1571        let (controller, executor) = podman_scan_controller(errored_podman_session(session_id));
1572
1573        let scan = controller.scan_orphan_workers(&executor, true);
1574
1575        let [candidate] = scan.candidates.as_slice() else {
1576            panic!("the errored session's container was not offered: {scan:?}");
1577        };
1578        assert_eq!(candidate.session_id, session_id);
1579        assert_eq!(
1580            candidate.tracked_session,
1581            Some(SessionState::Error),
1582            "recovery must say the session id is still taken, so only destroy applies"
1583        );
1584    }
1585
1586    #[test]
1587    fn a_container_a_session_still_names_is_not_an_orphan() {
1588        let session_id = "0123456789abcdef0123456789abcdef";
1589        let mut record = errored_podman_session(session_id);
1590        // The same failure, except the record kept its locator: the session's
1591        // own forced destroy still reaches the container, so recovery stays out.
1592        record.target = Some(TargetLocator::LocalPodman {
1593            container_id: targets::resource_name(session_id).unwrap(),
1594            workspace_storage: PodmanWorkspaceLocator::ContainerLayer,
1595            borrowed_from: None,
1596        });
1597        let (controller, executor) = podman_scan_controller(record);
1598
1599        let scan = controller.scan_orphan_workers(&executor, true);
1600
1601        assert!(
1602            scan.candidates.is_empty(),
1603            "a container the controller can still drive was offered for destruction: {scan:?}"
1604        );
1605    }
1606
1607    #[test]
1608    fn a_container_being_provisioned_is_not_an_orphan() {
1609        let session_id = "0123456789abcdef0123456789abcdef";
1610        let mut record = errored_podman_session(session_id);
1611        // Provisioning writes the locator as it runs, so a scan that catches
1612        // the window between `podman run` and locator discovery must not
1613        // mistake the new container for a leftover.
1614        record.state = SessionState::Provisioning;
1615        let (controller, executor) = podman_scan_controller(record);
1616
1617        let scan = controller.scan_orphan_workers(&executor, true);
1618
1619        assert!(
1620            scan.candidates.is_empty(),
1621            "a container still being provisioned was offered for destruction: {scan:?}"
1622        );
1623    }
1624
1625    fn candidate(session_id: &str, instance_id: Option<&str>) -> RecoveryCandidate {
1626        RecoveryCandidate {
1627            session_id: session_id.to_owned(),
1628            target_template_id: "local".to_owned(),
1629            locator: TargetLocator::LocalDocker {
1630                container_id: format!("mj-{session_id}"),
1631                borrowed_from: None,
1632            },
1633            ownership: None,
1634            instance_id: instance_id.map(str::to_owned),
1635            tracked_session: None,
1636        }
1637    }
1638
1639    #[test]
1640    fn default_scan_scope_hides_other_and_unknown_instances() {
1641        let mut scan = RecoveryScan {
1642            candidates: vec![
1643                candidate("mine", Some("qa")),
1644                candidate("theirs", Some("prod")),
1645                candidate("legacy", None),
1646            ],
1647            instance_id: "qa".to_owned(),
1648            ..RecoveryScan::default()
1649        };
1650        restrict_to_instance(&mut scan);
1651        assert_eq!(
1652            scan.candidates
1653                .iter()
1654                .map(|candidate| candidate.session_id.as_str())
1655                .collect::<Vec<_>>(),
1656            ["mine"]
1657        );
1658        assert_eq!(scan.hidden_other_instances, 2);
1659    }
1660
1661    #[test]
1662    fn acting_on_another_or_unknown_instance_requires_the_explicit_flag() {
1663        require_instance_access(&candidate("mine", Some("qa")), "qa", false).unwrap();
1664
1665        let other = require_instance_access(&candidate("theirs", Some("prod")), "qa", false)
1666            .expect_err("another instance's worker is refused by default");
1667        assert!(
1668            other.to_string().contains("belongs to instance \"prod\""),
1669            "{other}"
1670        );
1671        require_instance_access(&candidate("theirs", Some("prod")), "qa", true).unwrap();
1672
1673        let unknown = require_instance_access(&candidate("legacy", None), "qa", false)
1674            .expect_err("a worker without a stamp is refused by default");
1675        assert!(
1676            unknown.to_string().contains("no instance stamp"),
1677            "{unknown}"
1678        );
1679        require_instance_access(&candidate("legacy", None), "qa", true).unwrap();
1680    }
1681
1682    #[test]
1683    fn a_podman_orphan_destroy_plan_removes_its_workspace_volume() {
1684        // `destroy_orphan_worker` rescans, so this drives the two steps it runs
1685        // on the candidate it finds: the scan's locator and the close plan
1686        // built from it. Those are what carried the container-layer default.
1687        let session = "0123456789abcdef0123456789abcdef";
1688        let template = TargetTemplate::LocalPodman {
1689            container: ConfigContainer {
1690                image: "ignored".into(),
1691                pull_policy: Default::default(),
1692                platform: None,
1693                cpus: None,
1694                memory: None,
1695                environment: BTreeMap::new(),
1696                workspace_storage: PodmanWorkspaceStorage::PodmanVolume,
1697                build_cache: None,
1698            },
1699        };
1700        let json = serde_json::json!([
1701            {"Labels": {"dev.mj.managed": "true", "dev.mj.session": session}}
1702        ]);
1703
1704        let candidates = candidates_from_container_json(
1705            "local",
1706            &template,
1707            serde_json::to_string(&json).unwrap().as_bytes(),
1708        )
1709        .unwrap();
1710
1711        let [candidate] = candidates.as_slice() else {
1712            panic!("expected one orphan candidate, got {candidates:?}");
1713        };
1714        let volume = format!("{}-workspace", targets::resource_name(session).unwrap());
1715        assert!(
1716            matches!(
1717                &candidate.locator,
1718                TargetLocator::LocalPodman {
1719                    workspace_storage: PodmanWorkspaceLocator::Volume { name },
1720                    ..
1721                } if name == &volume
1722            ),
1723            "candidate locator lost the volume storage: {:?}",
1724            candidate.locator
1725        );
1726        let backend = recovery_backend_locator(&template, &candidate.locator, session).unwrap();
1727        let plan = targets::close_plan(&backend, session).unwrap();
1728        assert!(
1729            plan.commands.iter().any(|command| {
1730                command.args.iter().any(|argument| argument == &volume)
1731                    && command
1732                        .args
1733                        .iter()
1734                        .any(|argument| argument.contains("podman volume rm"))
1735            }),
1736            "destroy plan does not remove the workspace volume: {plan:?}"
1737        );
1738    }
1739
1740    #[test]
1741    fn recovery_docker_scan_accepts_json_lines_and_builds_a_docker_locator() {
1742        let template = TargetTemplate::LocalDocker {
1743            container: ConfigContainer {
1744                build_cache: None,
1745                image: "ignored".into(),
1746                pull_policy: Default::default(),
1747                platform: None,
1748                cpus: None,
1749                memory: None,
1750                environment: BTreeMap::new(),
1751                workspace_storage: Default::default(),
1752            },
1753        };
1754        let session = "0123456789abcdef0123456789abcdef";
1755        let output = format!(
1756            "{{\"Labels\":\"dev.mj.managed=true,dev.mj.session={session},dev.mj.instance=abc123\"}}\n{{\"Labels\":\"dev.mj.managed=false,dev.mj.session=ignored\"}}\n"
1757        );
1758
1759        let candidates =
1760            candidates_from_container_json("docker", &template, output.as_bytes()).unwrap();
1761
1762        assert_eq!(candidates.len(), 1);
1763        assert_eq!(candidates[0].session_id, session);
1764        assert_eq!(candidates[0].instance_id.as_deref(), Some("abc123"));
1765        assert!(matches!(
1766            &candidates[0].locator,
1767            TargetLocator::LocalDocker { container_id, .. }
1768                if container_id == &targets::resource_name(session).unwrap()
1769        ));
1770    }
1771
1772    #[test]
1773    fn recovery_aws_scan_uses_exact_tagged_instance_and_address() {
1774        let json = serde_json::json!({"Reservations": [{"Instances": [{
1775            "InstanceId": "i-exact",
1776            "PrivateIpAddress": "10.0.0.7",
1777            "Tags": [
1778                {"Key": "dev.mj.managed", "Value": "true"},
1779                {"Key": "dev.mj.session", "Value": "0123456789abcdef0123456789abcdef"},
1780                {"Key": "dev.mj.instance", "Value": "qa0916"}
1781            ]
1782        }]}]});
1783        let candidates = candidates_from_aws_json(
1784            "aws",
1785            AwsAddressSource::PrivateIp,
1786            serde_json::to_string(&json).unwrap().as_bytes(),
1787        )
1788        .unwrap();
1789        assert_eq!(candidates[0].instance_id.as_deref(), Some("qa0916"));
1790        assert!(matches!(
1791            &candidates[0].locator,
1792            TargetLocator::AwsEc2 { instance_id, address }
1793                if instance_id == "i-exact" && address.as_deref() == Some("10.0.0.7")
1794        ));
1795    }
1796}