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