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