Skip to main content

mj_controller/controller/
recovery_scan.rs

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