Skip to main content

mj_controller/controller/move_session/
destination.rs

1//! Durable preparation of an EC2 destination while the source worker remains live.
2use super::*;
3use crate::targets::{self, CommandSpec};
4use mj_core::config::TargetTemplate;
5use mj_core::state::{
6    PreparedDestinationState, PreparedMoveDestination, TargetLocator, TargetRuntimeSettings,
7};
8
9impl Controller {
10    pub(super) fn prepare_ec2_move_destination(
11        &self,
12        operation: &mut MoveOperation,
13        executor: &(impl CommandExecutor + Sync),
14    ) -> Result<()> {
15        executor.begin_resumable_move_work()?;
16        let result = self.prepare_ec2_move_destination_inner(operation, executor);
17        let admission = executor.end_resumable_move_work();
18        admission?;
19        result
20    }
21
22    fn prepare_ec2_move_destination_inner(
23        &self,
24        operation: &mut MoveOperation,
25        executor: &(impl CommandExecutor + Sync),
26    ) -> Result<()> {
27        let id = &operation.selection.session_id;
28        let template = self
29            .config
30            .targets
31            .get(
32                operation
33                    .selection
34                    .target_template_id
35                    .as_deref()
36                    .context("Move target missing")?,
37            )
38            .context("EC2 Move target disappeared")?;
39        ensure!(
40            matches!(template, TargetTemplate::AwsEc2 { .. }),
41            "Move preparation requires an EC2 target"
42        );
43        ensure!(
44            !executor.cancellation_requested(),
45            "Move cancelled before EC2 preparation"
46        );
47        let fresh_attempt = operation
48            .prepared_destination
49            .as_ref()
50            .is_none_or(|d| matches!(d.state, PreparedDestinationState::Released));
51        if fresh_attempt {
52            let backend = super::super::backend::backend_target(
53                template,
54                operation.selection.resource_allocation.as_ref(),
55                super::super::backend::ContainerOverrides::for_session(&self.state.sessions[id]),
56            )?;
57            let targets::TargetTemplate::AwsEc2(mut aws) = backend else {
58                bail!("EC2 Move backend mismatch");
59            };
60            let version = super::super::execute_checked(
61                executor,
62                CommandSpec::new(
63                    "aws",
64                    [
65                        "--profile".into(),
66                        aws.profile.clone(),
67                        "--region".into(),
68                        aws.region.clone(),
69                        "ec2".into(),
70                        "describe-launch-template-versions".into(),
71                        if aws.launch_template.starts_with("lt-") {
72                            "--launch-template-id"
73                        } else {
74                            "--launch-template-name"
75                        }
76                        .into(),
77                        aws.launch_template.clone(),
78                        "--versions".into(),
79                        aws.launch_template_version
80                            .clone()
81                            .unwrap_or_else(|| "$Default".into()),
82                        "--query".into(),
83                        "LaunchTemplateVersions[0].{LaunchTemplateId:LaunchTemplateId,VersionNumber:VersionNumber}".into(),
84                        "--output".into(),
85                        "json".into(),
86                    ],
87                )
88                .purpose("resolve immutable EC2 launch template version"),
89            )?;
90            let resolved: serde_json::Value = serde_json::from_slice(&version.stdout)
91                .context("parse immutable EC2 launch template identity")?;
92            let template_id = resolved
93                .get("LaunchTemplateId")
94                .and_then(serde_json::Value::as_str)
95                .filter(|id| id.starts_with("lt-"))
96                .context("EC2 launch template omitted immutable template ID")?;
97            let version = resolved
98                .get("VersionNumber")
99                .and_then(serde_json::Value::as_u64)
100                .context("EC2 launch template omitted numeric version")?;
101            aws.launch_template = template_id.to_owned();
102            aws.launch_template_version = Some(version.to_string());
103            let mut command = targets::ec2_launch_command(&aws, id)?;
104            let token = digest(&(&operation.operation_id, new_command_id("ec2-attempt")?))?;
105            command.args.extend([
106                "--client-token".into(),
107                token,
108                "--min-count".into(),
109                "1".into(),
110                "--max-count".into(),
111                "1".into(),
112            ]);
113            operation.prepared_destination = Some(PreparedMoveDestination {
114                launch_args: command.args,
115                runtime: TargetRuntimeSettings::from(template),
116                state: PreparedDestinationState::LaunchPending,
117            });
118            crate::database::save_move_operation(operation)?;
119        }
120        let destination = operation.prepared_destination.as_ref().unwrap().clone();
121        if matches!(destination.state, PreparedDestinationState::LaunchPending) {
122            executor.notify_notice("Creating EC2 destination");
123            let output = executor.execute(
124                &CommandSpec::new("aws", destination.launch_args.clone())
125                    .purpose("launch EC2 Move destination")
126                    .stage(ProvisionStage::Provisioning),
127            )?;
128            if output.status != 0 {
129                // Only a first-call refusal proves this attempt was never accepted.
130                // A refused retry says nothing about a previously lost response.
131                if fresh_attempt && launch_was_refused(&output.stderr) {
132                    operation.prepared_destination.as_mut().unwrap().state =
133                        PreparedDestinationState::Released;
134                    crate::database::save_move_operation(operation)?;
135                }
136                bail!(
137                    "EC2 Move launch failed: {}",
138                    String::from_utf8_lossy(&output.stderr)
139                );
140            }
141            let json: serde_json::Value =
142                serde_json::from_slice(&output.stdout).context("parse EC2 Move launch response")?;
143            let instance_id = json
144                .pointer("/Instances/0/InstanceId")
145                .and_then(serde_json::Value::as_str)
146                .context("EC2 Move launch omitted instance ID")?
147                .to_owned();
148            operation.prepared_destination.as_mut().unwrap().state =
149                PreparedDestinationState::Created { instance_id };
150            crate::database::save_move_operation(operation)?;
151        }
152        if let PreparedDestinationState::Created { instance_id } = operation
153            .prepared_destination
154            .as_ref()
155            .unwrap()
156            .state
157            .clone()
158        {
159            executor.notify_notice("Booting EC2 destination");
160            let target =
161                super::super::backend::ec2_locator_after_launch(template, instance_id, executor)?;
162            let backend = prepared_backend(&target, &destination.runtime, id)?;
163            let targets::TargetLocator::AwsEc2 { ssh, workspace, .. } = &backend else {
164                unreachable!()
165            };
166            super::super::execute_checked(
167                executor,
168                targets::ssh_command(ssh, ["mkdir", "-p", workspace])
169                    .purpose("create prepared EC2 workspace"),
170            )?;
171            targets::install_git_plan(targets::ExecutionBoundary::Ssh(ssh)).execute(executor)?;
172            targets::install_rsync_plan(targets::ExecutionBoundary::Ssh(ssh)).execute(executor)?;
173            operation.prepared_destination.as_mut().unwrap().state =
174                PreparedDestinationState::Checked { target };
175            crate::database::save_move_operation(operation)?;
176        }
177        ensure!(
178            matches!(
179                operation.prepared_destination.as_ref().unwrap().state,
180                PreparedDestinationState::Checked { .. } | PreparedDestinationState::Adopted { .. }
181            ),
182            "EC2 destination cleanup is pending; source retained"
183        );
184        Ok(())
185    }
186
187    /// One bounded cleanup attempt; the daemon retries durable failures with backoff.
188    pub fn cleanup_prepared_move_destination(
189        &self,
190        operation: &mut MoveOperation,
191        executor: &impl CommandExecutor,
192    ) -> Result<()> {
193        let cleanup = MoveCleanupExecutor(executor);
194        let executor = &cleanup;
195        let Some(destination) = operation
196            .prepared_destination
197            .as_ref()
198            .filter(|d| d.owns_resource())
199            .cloned()
200        else {
201            return Ok(());
202        };
203        let mut instance_id = destination.instance_id().map(str::to_owned);
204        operation.prepared_destination.as_mut().unwrap().state =
205            PreparedDestinationState::CleanupPending {
206                instance_id: instance_id.clone(),
207            };
208        crate::database::save_move_operation(operation)?;
209        if instance_id.is_none() {
210            let token = argument(&destination.launch_args, "--client-token")?;
211            let (profile, region) = aws_access(&destination.runtime)?;
212            let output = super::super::execute_checked(
213                executor,
214                CommandSpec::new(
215                    "aws",
216                    [
217                        "--profile",
218                        profile,
219                        "--region",
220                        region,
221                        "ec2",
222                        "describe-instances",
223                        "--filters",
224                        &format!("Name=client-token,Values={token}"),
225                        &format!(
226                            "Name=tag:dev.mj.session,Values={}",
227                            operation.selection.session_id
228                        ),
229                        "Name=tag:dev.mj.managed,Values=true",
230                        "--output",
231                        "json",
232                    ],
233                )
234                .purpose("reconcile EC2 Move launch acknowledgement"),
235            )?;
236            let json: serde_json::Value = serde_json::from_slice(&output.stdout)?;
237            let ids: Vec<&str> = json
238                .get("Reservations")
239                .and_then(serde_json::Value::as_array)
240                .into_iter()
241                .flatten()
242                .flat_map(|r| {
243                    r.get("Instances")
244                        .and_then(serde_json::Value::as_array)
245                        .into_iter()
246                        .flatten()
247                })
248                .filter_map(|i| i.get("InstanceId").and_then(serde_json::Value::as_str))
249                .collect();
250            ensure!(
251                ids.len() == 1,
252                "EC2 launch cleanup cannot yet resolve a unique instance for token {token}; automatic cleanup will retry"
253            );
254            instance_id = Some(ids[0].to_owned());
255            operation.prepared_destination.as_mut().unwrap().state =
256                PreparedDestinationState::CleanupPending {
257                    instance_id: instance_id.clone(),
258                };
259            crate::database::save_move_operation(operation)?;
260        }
261        let instance_id = instance_id.unwrap();
262        let (profile, region) = aws_access(&destination.runtime)?;
263        super::super::execute_checked(
264            executor,
265            targets::terminate_ec2_instance_command(profile, region, &instance_id)?,
266        )?;
267        // Confirm termination before releasing durable resource ownership.
268        super::super::execute_checked(
269            executor,
270            CommandSpec::new(
271                "aws",
272                [
273                    "--profile",
274                    profile,
275                    "--region",
276                    region,
277                    "ec2",
278                    "wait",
279                    "instance-terminated",
280                    "--instance-ids",
281                    &instance_id,
282                ],
283            )
284            .purpose("confirm EC2 Move destination termination"),
285        )?;
286        operation.prepared_destination.as_mut().unwrap().state = PreparedDestinationState::Released;
287        crate::database::save_move_operation(operation)?;
288        Ok(())
289    }
290
291    pub(in crate::controller) fn adopt_prepared_ec2_destination(
292        &mut self,
293        id: &str,
294    ) -> Result<Option<targets::CommandPlan>> {
295        let Some(mut operation) = crate::database::load_move_operation(id)? else {
296            return Ok(None);
297        };
298        if operation.phase != MovePhase::ResumingDestination {
299            return Ok(None);
300        }
301        let Some(destination) = operation.prepared_destination.as_ref() else {
302            return Ok(None);
303        };
304        let target = destination
305            .target()
306            .context("prepared EC2 destination is not available for adoption")?
307            .clone();
308        let runtime = destination.runtime.clone();
309        let backend = prepared_backend(&target, &runtime, id)?;
310        let bundle = self
311            .move_destination_bundle(id)?
312            .context("prepared EC2 Move bundle missing")?;
313        let plan = targets::provision_on_locator_plan(&backend, id, &bundle)?;
314        let session = self
315            .state
316            .sessions
317            .get_mut(id)
318            .context("EC2 Move session missing")?;
319        session.target = Some(target.clone());
320        session.target_runtime = Some(runtime);
321        operation.prepared_destination.as_mut().unwrap().state =
322            PreparedDestinationState::Adopted { target };
323        crate::database::adopt_move_destination(&operation, session)?;
324        Ok(Some(plan))
325    }
326}
327
328struct MoveCleanupExecutor<'a, E>(&'a E);
329
330impl<E: CommandExecutor> CommandExecutor for MoveCleanupExecutor<'_, E> {
331    fn execute(&self, command: &CommandSpec) -> Result<targets::CommandOutput> {
332        self.0.execute_cleanup(command)
333    }
334}
335
336fn argument<'a>(args: &'a [String], key: &str) -> Result<&'a str> {
337    args.windows(2)
338        .find(|a| a[0] == key)
339        .map(|a| a[1].as_str())
340        .with_context(|| format!("persisted EC2 launch lacks {key}"))
341}
342
343fn aws_access(runtime: &TargetRuntimeSettings) -> Result<(&str, &str)> {
344    let mj_core::state::TargetConnection::Aws {
345        profile, region, ..
346    } = &runtime.connection
347    else {
348        bail!("prepared destination has no AWS access settings");
349    };
350    Ok((profile, region))
351}
352
353pub(super) fn prepared_backend(
354    target: &TargetLocator,
355    runtime: &TargetRuntimeSettings,
356    id: &str,
357) -> Result<targets::TargetLocator> {
358    Ok(targets::TargetLocator::try_from(targets::RecordedTarget {
359        locator: target,
360        runtime: Some(runtime),
361        session_id: id,
362    })?)
363}
364
365fn launch_was_refused(stderr: &[u8]) -> bool {
366    let detail = String::from_utf8_lossy(stderr);
367    [
368        "UnauthorizedOperation",
369        "AuthFailure",
370        "InvalidParameterValue",
371        "InvalidParameterCombination",
372        "InsufficientInstanceCapacity",
373        "InstanceLimitExceeded",
374        "InvalidLaunchTemplateName.NotFoundException",
375        "InvalidLaunchTemplateId.NotFound",
376    ]
377    .iter()
378    .any(|code| detail.contains(&format!("An error occurred ({code})")))
379}
380
381#[cfg(all(test, unix))]
382pub(super) mod tests;