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