1use crate::{
2 ExtractionBackfillRun, ExtractionBackfillStatus, ExtractionPlan,
3 ExtractionReconciliationResult, ExtractionReconciliationStatus, extraction_input_digest,
4 extraction_plan_integrity_is_valid,
5};
6use schemars::JsonSchema;
7use serde::{Deserialize, Serialize};
8use std::fmt;
9
10pub const EXTRACTION_QUIESCENCE_PROTOCOL: &str = "lenso.extraction-quiescence.v1";
11
12#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
13#[serde(rename_all = "camelCase")]
14pub struct ExtractionDrainSnapshot {
15 pub in_flight_requests: u64,
16 pub outbox_messages: u64,
17 pub inbox_messages: u64,
18 pub scheduled_functions: u64,
19 pub timers: u64,
20 pub durable_workflows: u64,
21 #[serde(default)]
22 pub unresolved: Vec<String>,
23 pub timed_out: bool,
24}
25
26impl ExtractionDrainSnapshot {
27 #[must_use]
28 pub const fn empty() -> Self {
29 Self {
30 in_flight_requests: 0,
31 outbox_messages: 0,
32 inbox_messages: 0,
33 scheduled_functions: 0,
34 timers: 0,
35 durable_workflows: 0,
36 unresolved: Vec::new(),
37 timed_out: false,
38 }
39 }
40
41 fn pending_count(&self) -> u64 {
42 self.in_flight_requests
43 + self.outbox_messages
44 + self.inbox_messages
45 + self.scheduled_functions
46 + self.timers
47 + self.durable_workflows
48 }
49}
50
51#[derive(
52 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
53)]
54#[serde(rename_all = "snake_case")]
55pub enum ExtractionQuiescenceStatus {
56 MutationsPaused,
57 Draining,
58 Blocked,
59 Drained,
60 Quiesced,
61 Cancelled,
62}
63
64#[derive(
65 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
66)]
67#[serde(rename_all = "snake_case")]
68pub enum ExtractionQuiescenceIssueCode {
69 PlanInvalid,
70 PlanStale,
71 AuthorityChanged,
72 DrainIncomplete,
73 DrainTimedOut,
74 FinalBackfillIncomplete,
75 FinalReconciliationMismatch,
76}
77
78#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
79#[serde(rename_all = "camelCase")]
80pub struct ExtractionQuiescenceIssue {
81 pub code: ExtractionQuiescenceIssueCode,
82 pub subject: String,
83 pub detail: String,
84 pub next_actions: Vec<String>,
85}
86
87#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
88#[serde(rename_all = "camelCase")]
89pub struct ExtractionQuiescenceEvidence {
90 pub kind: String,
91 pub subject: String,
92 pub digest: String,
93 pub detail: String,
94}
95
96#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
97#[serde(rename_all = "camelCase")]
98pub struct ExtractionQuiescenceEffects {
99 pub pauses_linked_mutations: bool,
100 pub drains_in_flight_work: bool,
101 pub copies_final_delta: bool,
102 pub routes_candidate_traffic: bool,
103 pub changes_authority: bool,
104 pub requires_runtime_console: bool,
105 pub requires_system_plane_for_business_execution: bool,
106}
107
108#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
109#[serde(rename_all = "camelCase")]
110pub struct ExtractionQuiescenceRun {
111 pub protocol: String,
112 pub quiescence_id: String,
113 pub quiescence_digest: String,
114 pub revision: u64,
115 pub status: ExtractionQuiescenceStatus,
116 pub plan_id: String,
117 pub plan_digest: String,
118 pub expected_authority_revision: String,
119 pub linked_mutations_paused: bool,
120 pub linked_read_inspection_available: bool,
121 pub linked_authority_remains_authoritative: bool,
122 pub candidate_authoritative: bool,
123 #[serde(default, skip_serializing_if = "Option::is_none")]
124 pub drain: Option<ExtractionDrainSnapshot>,
125 #[serde(default, skip_serializing_if = "Option::is_none")]
126 pub stable_source_high_water_mark: Option<String>,
127 #[serde(default, skip_serializing_if = "Option::is_none")]
128 pub destination_checkpoint: Option<String>,
129 #[serde(default)]
130 pub issues: Vec<ExtractionQuiescenceIssue>,
131 #[serde(default)]
132 pub evidence: Vec<ExtractionQuiescenceEvidence>,
133 #[serde(default)]
134 pub next_actions: Vec<String>,
135 pub effects: ExtractionQuiescenceEffects,
136}
137
138#[derive(Debug, Clone, PartialEq, Eq)]
139pub struct ExtractionQuiescenceStartError {
140 pub code: ExtractionQuiescenceIssueCode,
141 pub message: String,
142}
143
144impl fmt::Display for ExtractionQuiescenceStartError {
145 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
146 formatter.write_str(&self.message)
147 }
148}
149
150impl std::error::Error for ExtractionQuiescenceStartError {}
151
152pub fn start_extraction_quiescence(
153 plan: &ExtractionPlan,
154 current_authority_revision: &str,
155) -> Result<ExtractionQuiescenceRun, ExtractionQuiescenceStartError> {
156 if !extraction_plan_integrity_is_valid(plan) {
157 return Err(ExtractionQuiescenceStartError {
158 code: ExtractionQuiescenceIssueCode::PlanInvalid,
159 message: "Extraction Plan integrity validation failed before write pause.".to_owned(),
160 });
161 }
162 if plan.expected_authority.revision != current_authority_revision {
163 return Err(ExtractionQuiescenceStartError {
164 code: ExtractionQuiescenceIssueCode::AuthorityChanged,
165 message: "Authority changed after the Extraction Plan was generated.".to_owned(),
166 });
167 }
168 let identity = digest(&(
169 plan.plan_id.as_str(),
170 plan.plan_digest.as_str(),
171 current_authority_revision,
172 ));
173 let mut run = ExtractionQuiescenceRun {
174 protocol: EXTRACTION_QUIESCENCE_PROTOCOL.to_owned(),
175 quiescence_id: format!("extraction-quiescence:{identity}"),
176 quiescence_digest: String::new(),
177 revision: 1,
178 status: ExtractionQuiescenceStatus::MutationsPaused,
179 plan_id: plan.plan_id.clone(),
180 plan_digest: plan.plan_digest.clone(),
181 expected_authority_revision: current_authority_revision.to_owned(),
182 linked_mutations_paused: true,
183 linked_read_inspection_available: true,
184 linked_authority_remains_authoritative: true,
185 candidate_authoritative: false,
186 drain: None,
187 stable_source_high_water_mark: None,
188 destination_checkpoint: None,
189 issues: Vec::new(),
190 evidence: vec![evidence(
191 "write_pause",
192 "linked-mutations",
193 current_authority_revision,
194 "New linked mutations are paused while read-only inspection remains available.",
195 )],
196 next_actions: vec![
197 "Drain in-flight requests, Outbox, Inbox, schedules, timers, and Durable Workflows."
198 .to_owned(),
199 ],
200 effects: ExtractionQuiescenceEffects {
201 pauses_linked_mutations: true,
202 ..ExtractionQuiescenceEffects::default()
203 },
204 };
205 refresh(&mut run);
206 Ok(run)
207}
208
209#[must_use]
210pub fn record_extraction_drain(
211 mut run: ExtractionQuiescenceRun,
212 mut snapshot: ExtractionDrainSnapshot,
213) -> ExtractionQuiescenceRun {
214 snapshot.unresolved.sort();
215 snapshot.unresolved.dedup();
216 run.drain = Some(snapshot.clone());
217 run.effects.drains_in_flight_work = true;
218 run.issues.clear();
219 if snapshot.timed_out {
220 push_issue(
221 &mut run,
222 ExtractionQuiescenceIssueCode::DrainTimedOut,
223 "drain",
224 "The bounded drain timed out.",
225 "Cancel safely or remediate unresolved work before retrying.",
226 );
227 } else if snapshot.pending_count() > 0 || !snapshot.unresolved.is_empty() {
228 push_issue(
229 &mut run,
230 ExtractionQuiescenceIssueCode::DrainIncomplete,
231 "drain",
232 "Eligible work has not drained completely.",
233 "Resolve each recorded work identity; do not abandon it.",
234 );
235 } else {
236 run.status = ExtractionQuiescenceStatus::Drained;
237 run.evidence.push(evidence(
238 "drain_complete",
239 "linked-runtime-work",
240 &digest(&snapshot),
241 "Requests, Outbox, Inbox, schedules, timers, and Durable Workflows are drained.",
242 ));
243 run.next_actions = vec![
244 "Copy the final Postgres delta and reconcile at the stable high-water mark.".to_owned(),
245 ];
246 }
247 run.revision += 1;
248 refresh(&mut run);
249 run
250}
251
252#[must_use]
253pub fn complete_extraction_quiescence(
254 mut run: ExtractionQuiescenceRun,
255 backfill: &ExtractionBackfillRun,
256 reconciliation: &ExtractionReconciliationResult,
257 current_plan_digest: &str,
258 current_authority_revision: &str,
259) -> ExtractionQuiescenceRun {
260 run.issues.clear();
261 if run.plan_digest != current_plan_digest {
262 push_issue(
263 &mut run,
264 ExtractionQuiescenceIssueCode::PlanStale,
265 "plan",
266 "Plan inputs changed during drain.",
267 "Regenerate the plan before retrying Cutover.",
268 );
269 }
270 if run.expected_authority_revision != current_authority_revision {
271 push_issue(
272 &mut run,
273 ExtractionQuiescenceIssueCode::AuthorityChanged,
274 "authority",
275 "Authority changed during drain.",
276 "Return to the recorded linked authority before retrying.",
277 );
278 }
279 if backfill.status != ExtractionBackfillStatus::Succeeded
280 || backfill.scope.plan_id != run.plan_id
281 {
282 push_issue(
283 &mut run,
284 ExtractionQuiescenceIssueCode::FinalBackfillIncomplete,
285 "final-delta",
286 "Final delta backfill is incomplete or belongs to another plan.",
287 "Complete the final checkpointed delta under the write pause.",
288 );
289 }
290 if reconciliation.status != ExtractionReconciliationStatus::Matched
291 || reconciliation.plan_id != run.plan_id
292 || reconciliation.source_high_water_mark != backfill.progress.source_high_water_mark
293 || reconciliation.destination_checkpoint
294 != backfill
295 .progress
296 .destination_checkpoint
297 .clone()
298 .unwrap_or_default()
299 {
300 push_issue(
301 &mut run,
302 ExtractionQuiescenceIssueCode::FinalReconciliationMismatch,
303 "reconciliation",
304 "Final reconciliation does not match the stable delta checkpoint.",
305 "Reconcile the exact final high-water mark and checkpoint.",
306 );
307 }
308 if run.issues.is_empty() && run.status == ExtractionQuiescenceStatus::Drained {
309 run.status = ExtractionQuiescenceStatus::Quiesced;
310 run.stable_source_high_water_mark = Some(backfill.progress.source_high_water_mark.clone());
311 run.destination_checkpoint = backfill.progress.destination_checkpoint.clone();
312 run.effects.copies_final_delta = true;
313 run.evidence.push(evidence(
314 "quiesced",
315 "linked-authority",
316 &reconciliation.reconciliation_digest,
317 "Final delta and reconciliation are stable; authority has not transferred.",
318 ));
319 run.next_actions = vec![
320 "Use this evidence for provisional Cutover while external mutations remain paused."
321 .to_owned(),
322 ];
323 }
324 run.revision += 1;
325 refresh(&mut run);
326 run
327}
328
329#[must_use]
330pub fn cancel_extraction_quiescence(
331 mut run: ExtractionQuiescenceRun,
332 reason: impl Into<String>,
333) -> ExtractionQuiescenceRun {
334 let reason = reason.into();
335 run.status = ExtractionQuiescenceStatus::Cancelled;
336 run.linked_mutations_paused = false;
337 run.effects.pauses_linked_mutations = false;
338 run.evidence.push(evidence(
339 "write_pause_released",
340 "linked-mutations",
341 &reason,
342 &reason,
343 ));
344 run.next_actions = vec![
345 "Linked implementation remains authoritative; create fresh evidence before retrying."
346 .to_owned(),
347 ];
348 run.revision += 1;
349 refresh(&mut run);
350 run
351}
352
353fn push_issue(
354 run: &mut ExtractionQuiescenceRun,
355 code: ExtractionQuiescenceIssueCode,
356 subject: &str,
357 detail: &str,
358 next_action: &str,
359) {
360 run.status = ExtractionQuiescenceStatus::Blocked;
361 run.issues.push(ExtractionQuiescenceIssue {
362 code,
363 subject: subject.to_owned(),
364 detail: detail.to_owned(),
365 next_actions: vec![next_action.to_owned()],
366 });
367 run.next_actions = vec![next_action.to_owned()];
368}
369
370fn evidence(kind: &str, subject: &str, value: &str, detail: &str) -> ExtractionQuiescenceEvidence {
371 ExtractionQuiescenceEvidence {
372 kind: kind.to_owned(),
373 subject: subject.to_owned(),
374 digest: extraction_input_digest(value.as_bytes()),
375 detail: detail.to_owned(),
376 }
377}
378
379fn refresh(run: &mut ExtractionQuiescenceRun) {
380 run.quiescence_digest.clear();
381 run.quiescence_digest = digest(run);
382}
383
384fn digest(value: &impl Serialize) -> String {
385 extraction_input_digest(&serde_json::to_vec(value).expect("quiescence values serialize"))
386}
387
388#[must_use]
389pub fn extraction_quiescence_integrity_is_valid(run: &ExtractionQuiescenceRun) -> bool {
390 if run.protocol != EXTRACTION_QUIESCENCE_PROTOCOL {
391 return false;
392 }
393 let mut value = run.clone();
394 value.quiescence_digest.clear();
395 run.quiescence_digest == digest(&value)
396}