mj_controller/controller/checkpoint/
capture.rs1use 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 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 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 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 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}