Skip to main content

lenso_service/
extraction_quiescence.rs

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}