Skip to main content

mj_controller/controller/checkpoint/
capture.rs

1use super::*;
2
3impl Controller {
4    pub(super) async fn checkpoint_session_controlled_with_manager(
5        &mut self,
6        session_id: &str,
7        executor: &(impl CommandExecutor + Sync),
8        manager: Option<&SessionManagerControl>,
9    ) -> Result<CheckpointMetadata> {
10        let previous = self.state.sessions.get(session_id).cloned();
11        let operation_id = new_command_id("checkpoint-operation")?;
12        let result = self
13            .checkpoint_session_owned(session_id, executor, manager, &operation_id)
14            .await;
15        if let Err(error) = &result
16            && let Some(record) = self.state.sessions.get_mut(session_id)
17        {
18            if record.state == SessionState::Checkpointing {
19                record.state = previous
20                    .as_ref()
21                    .map_or(SessionState::Running, |record| record.state);
22            }
23            if !checkpoint_was_deferred(error) {
24                record.last_checkpoint_error = Some(format!("{error:#}"));
25            }
26            record.updated_at = now();
27            crate::database::finish_failed_checkpoint(
28                record,
29                &operation_id,
30                checkpoint_was_deferred(error),
31                format!("{error:#}"),
32            )
33            .with_context(|| format!("persist checkpoint failure after: {error:#}"))?;
34        }
35        result
36    }
37
38    async fn checkpoint_session_owned(
39        &mut self,
40        session_id: &str,
41        executor: &(impl CommandExecutor + Sync),
42        manager: Option<&SessionManagerControl>,
43        operation_id: &str,
44    ) -> Result<CheckpointMetadata> {
45        let previous = self
46            .state
47            .sessions
48            .get(session_id)
49            .with_context(|| format!("unknown session {session_id}"))?
50            .clone();
51        previous.validate_configuration(&self.config)?;
52        ensure!(
53            !matches!(
54                previous.state,
55                SessionState::Closing | SessionState::Destroying
56            ),
57            "session {session_id} is already closing; resume that close instead of starting an ordinary checkpoint"
58        );
59        let record = self.state.sessions.get_mut(session_id).unwrap();
60        record.state = SessionState::Checkpointing;
61        record.updated_at = now();
62        record.last_checkpoint_error = None;
63        if let Err(error) = crate::database::begin_checkpoint_operation(record, operation_id) {
64            self.state.sessions.insert(session_id.to_owned(), previous);
65            return Err(error.context("persist checkpoint operation before capture"));
66        }
67
68        match self
69            .checkpoint_session_latched_for_operation(
70                session_id,
71                executor,
72                manager,
73                LatchExclusivity::ReleaseAfterLatch,
74                CheckpointExportPolicy::Always,
75                Some(operation_id),
76                None,
77            )
78            .await
79        {
80            Ok(latched) => {
81                let artifact = latched.artifact.clone();
82                if let Err(error) = mj_core::test_hooks::reach_test_hook(
83                    "checkpoint_archive_before_database_publication",
84                ) {
85                    latched.abandon(session_id).await;
86                    return Err(remove_uninstalled_checkpoint(
87                        &artifact.metadata.archive_path,
88                        error,
89                    ));
90                }
91                {
92                    let record = self.state.sessions.get_mut(session_id).unwrap();
93                    record.state = SessionState::Running;
94                    record.native_session_id = Some(artifact.native_session_id.clone());
95                    record.checkpoint = Some(artifact.metadata.clone());
96                    record.updated_at = now();
97                    record.last_error = None;
98                    record.last_checkpoint_error = None;
99                }
100                let persist_started = Instant::now();
101                if let Err(error) = crate::database::save_requested_checkpoint(
102                    self.state
103                        .sessions
104                        .get(session_id)
105                        .expect("checkpoint session exists"),
106                    operation_id,
107                ) {
108                    self.state
109                        .sessions
110                        .insert(session_id.to_owned(), previous.clone());
111                    latched.abandon(session_id).await;
112                    return Err(error);
113                }
114                tracing::info!(
115                    session_id,
116                    persist_ms = persist_started.elapsed().as_millis() as u64,
117                    "checkpoint metadata persisted"
118                );
119                prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
120                release_projection_behind_checkpoint(session_id, &artifact.metadata);
121                if let Err(error) = latched.complete().await {
122                    // Only journal retention is at stake. A barrier that is
123                    // still open cannot dangle: the actor retries a failed
124                    // submission over a fresh connection, and the worker
125                    // cancels barriers whose submitting connection dropped.
126                    // The next checkpoint moves the recovery floor again.
127                    tracing::warn!(
128                        session_id,
129                        "verified checkpoint was saved, but the relay could not be told to release the history it covers: {error:#}"
130                    );
131                }
132                Ok(artifact.metadata)
133            }
134            Err(error) => {
135                // A deferred checkpoint says the agent was working, not that
136                // anything failed. Recording it would leave a warning on the
137                // session row until the next successful copy, so the caller is
138                // told and the row is left alone.
139                let deferred = checkpoint_was_deferred(&error);
140                if let Some(record) = self.state.sessions.get_mut(session_id) {
141                    record.state = if previous.state == SessionState::Checkpointing {
142                        SessionState::Running
143                    } else {
144                        previous.state
145                    };
146                    record.updated_at = now();
147                    if !deferred {
148                        record.last_checkpoint_error = Some(format!("{error:#}"));
149                    }
150                }
151                Err(error)
152            }
153        }
154    }
155
156    /// Create, checksum, and durably install a recovery archive before
157    /// allowing the relay to garbage-collect through its event frontier.
158    pub async fn create_recovery_checkpoint_managed_controlled(
159        &self,
160        session_id: &str,
161        manager: &SessionManagerControl,
162        executor: &(impl CommandExecutor + Sync),
163    ) -> Result<CheckpointArtifact> {
164        self.create_recovery_checkpoint_with_manager(session_id, Some(manager), executor)
165            .await
166    }
167
168    pub(super) async fn create_recovery_checkpoint_with_manager(
169        &self,
170        session_id: &str,
171        manager: Option<&SessionManagerControl>,
172        executor: &(impl CommandExecutor + Sync),
173    ) -> Result<CheckpointArtifact> {
174        let observed = self
175            .state
176            .sessions
177            .get(session_id)
178            .with_context(|| format!("unknown session {session_id}"))?;
179        let expected_target = observed
180            .target
181            .as_ref()
182            .context("recovery session has no target")?;
183        let previous_checkpoint = observed.checkpoint.clone();
184        let latched = self
185            .checkpoint_session_latched_with_recovery_stage(
186                session_id,
187                executor,
188                manager,
189                LatchExclusivity::ReleaseAfterLatch,
190                CheckpointExportPolicy::Always,
191                true,
192                None,
193                None,
194            )
195            .await?;
196        let artifact = latched.artifact.clone();
197        let verification = {
198            let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
199            verify_checkpoint_artifact(session_id, &artifact)
200        };
201        if let Err(error) = verification {
202            latched.abandon(session_id).await;
203            return Err(remove_uninstalled_checkpoint(
204                &artifact.metadata.archive_path,
205                error.context("final recovery checkpoint verification"),
206            ));
207        }
208        if let Err(error) =
209            mj_core::test_hooks::reach_test_hook("checkpoint_archive_before_database_publication")
210        {
211            latched.abandon(session_id).await;
212            return Err(remove_uninstalled_checkpoint(
213                &artifact.metadata.archive_path,
214                error,
215            ));
216        }
217        let persist_started = Instant::now();
218        let installed = crate::database::record_recovery_success_if_current(
219            session_id,
220            expected_target,
221            previous_checkpoint.as_ref(),
222            &artifact.native_session_id,
223            &artifact.metadata,
224        );
225        match installed {
226            Ok(true) => {}
227            Ok(false) => {
228                latched.abandon(session_id).await;
229                return Err(remove_uninstalled_checkpoint(
230                    &artifact.metadata.archive_path,
231                    anyhow::anyhow!(
232                        "recovery checkpoint belongs to an earlier session or checkpoint generation"
233                    ),
234                ));
235            }
236            Err(error) => {
237                latched.abandon(session_id).await;
238                return Err(error.context(
239                    "persist verified recovery checkpoint before releasing relay history",
240                ));
241            }
242        }
243        tracing::info!(
244            session_id,
245            persist_ms = persist_started.elapsed().as_millis() as u64,
246            "recovery checkpoint metadata persisted"
247        );
248        if let Err(error) = latched.complete().await {
249            // Only journal retention is at stake. A barrier that is still open
250            // cannot dangle: the actor retries a failed submission over a fresh
251            // connection, and the worker cancels barriers whose submitting
252            // connection dropped. The next checkpoint moves the floor again.
253            tracing::warn!(
254                session_id,
255                "recovery checkpoint was saved, but the relay could not be told to release the history it covers: {error:#}"
256            );
257        }
258        prune_replaced_checkpoint(previous_checkpoint.as_ref(), &artifact.metadata);
259        release_projection_behind_checkpoint(session_id, &artifact.metadata);
260        Ok(artifact)
261    }
262}