Skip to main content

lenso_service/
extraction_provisional_cutover.rs

1use crate::{
2    ExtractionQuiescenceRun, ExtractionQuiescenceStatus, ExtractionVerificationResult,
3    ExtractionVerificationStatus, extraction_input_digest,
4};
5use schemars::JsonSchema;
6use serde::{Deserialize, Serialize};
7use std::fmt;
8
9pub const EXTRACTION_PROVISIONAL_CUTOVER_PROTOCOL: &str = "lenso.extraction-provisional-cutover.v1";
10
11#[derive(
12    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
13)]
14#[serde(rename_all = "snake_case")]
15pub enum ExtractionProvisionalCutoverStatus {
16    Provisional,
17    Verified,
18    RolledBack,
19}
20
21#[derive(
22    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
23)]
24#[serde(rename_all = "snake_case")]
25pub enum ExtractionTrafficRoute {
26    Linked,
27    CandidateVerificationOnly,
28}
29
30#[derive(
31    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
32)]
33#[serde(rename_all = "snake_case")]
34pub enum ExtractionProvisionalCutoverIssueCode {
35    PlanStale,
36    AuthorityChanged,
37    SourceNotQuiesced,
38    FinalReconciliationMissing,
39    CandidateUnhealthy,
40    CompatibilityVerificationFailed,
41    PolicyEvidenceFailed,
42    StoryComparisonFailed,
43}
44
45#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
46#[serde(rename_all = "camelCase")]
47pub struct ExtractionProvisionalCutoverInputs {
48    pub plan_id: String,
49    pub plan_digest: String,
50    pub authority_revision: String,
51    pub routing_revision: String,
52    pub candidate_service_id: String,
53    pub candidate_healthy: bool,
54    pub verification: ExtractionVerificationResult,
55    pub quiescence: ExtractionQuiescenceRun,
56}
57
58#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
59#[serde(rename_all = "camelCase")]
60pub struct ExtractionCutoverReceipt {
61    pub step_id: String,
62    pub step_digest: String,
63    pub from_revision: String,
64    pub to_revision: String,
65    pub outcome: String,
66}
67
68#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
69#[serde(rename_all = "camelCase")]
70pub struct ExtractionCutoverEvidence {
71    pub kind: String,
72    pub subject: String,
73    pub digest: String,
74    pub detail: String,
75}
76
77#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
78#[serde(rename_all = "camelCase")]
79pub struct ExtractionLinkedRollbackValidation {
80    pub validation_digest: String,
81    pub authority_revision: String,
82    pub routing_revision: String,
83    pub business_probe_digest: String,
84    pub passed: bool,
85}
86
87impl ExtractionLinkedRollbackValidation {
88    #[must_use]
89    pub fn bind(
90        run: &ExtractionProvisionalCutoverRun,
91        business_probe_digest: impl Into<String>,
92        passed: bool,
93    ) -> Self {
94        let mut validation = Self {
95            validation_digest: String::new(),
96            authority_revision: run.authority_revision.clone(),
97            routing_revision: run.routing_revision_before.clone(),
98            business_probe_digest: business_probe_digest.into(),
99            passed,
100        };
101        validation.validation_digest = rollback_validation_digest(&validation);
102        validation
103    }
104
105    fn is_valid_for(&self, run: &ExtractionProvisionalCutoverRun) -> bool {
106        self.validation_digest == rollback_validation_digest(self)
107            && self.authority_revision == run.authority_revision
108            && self.routing_revision == run.routing_revision_before
109            && !self.business_probe_digest.trim().is_empty()
110    }
111}
112
113#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
114#[serde(rename_all = "camelCase")]
115pub struct ExtractionProvisionalCutoverRun {
116    pub protocol: String,
117    pub cutover_id: String,
118    pub cutover_digest: String,
119    pub revision: u64,
120    pub status: ExtractionProvisionalCutoverStatus,
121    pub plan_id: String,
122    pub plan_digest: String,
123    pub authority_revision: String,
124    pub routing_revision_before: String,
125    pub routing_revision_current: String,
126    pub candidate_service_id: String,
127    pub verification_digest: String,
128    pub quiescence_digest: String,
129    pub source_high_water_mark: String,
130    pub destination_checkpoint: String,
131    pub route: ExtractionTrafficRoute,
132    pub external_mutations_paused: bool,
133    pub linked_mutations_open: bool,
134    pub linked_authoritative: bool,
135    pub candidate_authoritative: bool,
136    pub candidate_healthy: bool,
137    pub declared_verification_traffic_only: bool,
138    pub verification_effects_isolated: bool,
139    pub linked_business_probe_passed: bool,
140    #[serde(default)]
141    pub apply_receipts: Vec<ExtractionCutoverReceipt>,
142    #[serde(default)]
143    pub rollback_receipts: Vec<ExtractionCutoverReceipt>,
144    #[serde(default)]
145    pub evidence: Vec<ExtractionCutoverEvidence>,
146}
147
148#[derive(Debug, Clone, PartialEq, Eq)]
149pub struct ExtractionProvisionalCutoverError {
150    pub code: ExtractionProvisionalCutoverIssueCode,
151    pub message: String,
152    pub next_actions: Vec<String>,
153}
154
155impl fmt::Display for ExtractionProvisionalCutoverError {
156    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
157        formatter.write_str(&self.message)
158    }
159}
160
161impl std::error::Error for ExtractionProvisionalCutoverError {}
162
163pub fn start_provisional_cutover(
164    inputs: ExtractionProvisionalCutoverInputs,
165) -> Result<ExtractionProvisionalCutoverRun, ExtractionProvisionalCutoverError> {
166    if inputs.plan_id != inputs.quiescence.plan_id
167        || inputs.plan_digest != inputs.quiescence.plan_digest
168        || inputs.plan_id != inputs.verification.plan_id
169    {
170        return Err(error(
171            ExtractionProvisionalCutoverIssueCode::PlanStale,
172            "Provisional Cutover evidence does not belong to one exact plan.",
173            "Regenerate verification and quiescence evidence for the current plan.",
174        ));
175    }
176    if inputs.authority_revision != inputs.quiescence.expected_authority_revision {
177        return Err(error(
178            ExtractionProvisionalCutoverIssueCode::AuthorityChanged,
179            "Authority changed after quiescence began.",
180            "Restore or regenerate evidence for the current linked authority.",
181        ));
182    }
183    if inputs.quiescence.status != ExtractionQuiescenceStatus::Quiesced
184        || !inputs.quiescence.linked_mutations_paused
185        || !inputs.quiescence.linked_authority_remains_authoritative
186    {
187        return Err(error(
188            ExtractionProvisionalCutoverIssueCode::SourceNotQuiesced,
189            "Linked source is not stably quiesced.",
190            "Complete drain, final delta, and reconciliation under the write pause.",
191        ));
192    }
193    if inputs.quiescence.stable_source_high_water_mark.is_none()
194        || inputs.quiescence.destination_checkpoint.is_none()
195    {
196        return Err(error(
197            ExtractionProvisionalCutoverIssueCode::FinalReconciliationMissing,
198            "Stable final reconciliation pins are missing.",
199            "Record the final high-water mark and destination checkpoint.",
200        ));
201    }
202    if !inputs.candidate_healthy {
203        return Err(error(
204            ExtractionProvisionalCutoverIssueCode::CandidateUnhealthy,
205            "Candidate health verification failed.",
206            "Restore candidate health before provisional routing.",
207        ));
208    }
209    if inputs.verification.status != ExtractionVerificationStatus::Verified
210        || !inputs.verification.provisional_cutover_eligible
211    {
212        return Err(error(
213            ExtractionProvisionalCutoverIssueCode::CompatibilityVerificationFailed,
214            "Compatibility and behavior verification did not pass.",
215            "Remediate all verification blockers before provisional routing.",
216        ));
217    }
218    let identity = digest(&(
219        inputs.plan_id.as_str(),
220        inputs.plan_digest.as_str(),
221        inputs.authority_revision.as_str(),
222        inputs.routing_revision.as_str(),
223        inputs.verification.verification_digest.as_str(),
224        inputs.quiescence.quiescence_digest.as_str(),
225    ));
226    let provisional_routing_revision = format!("provisional:{identity}");
227    let apply_receipt = receipt(
228        "route-verification-traffic",
229        &inputs.routing_revision,
230        &provisional_routing_revision,
231        "applied",
232    );
233    let mut run = ExtractionProvisionalCutoverRun {
234        protocol: EXTRACTION_PROVISIONAL_CUTOVER_PROTOCOL.to_owned(),
235        cutover_id: format!("extraction-cutover:{identity}"),
236        cutover_digest: String::new(),
237        revision: 1,
238        status: ExtractionProvisionalCutoverStatus::Provisional,
239        plan_id: inputs.plan_id,
240        plan_digest: inputs.plan_digest,
241        authority_revision: inputs.authority_revision,
242        routing_revision_before: inputs.routing_revision,
243        routing_revision_current: provisional_routing_revision,
244        candidate_service_id: inputs.candidate_service_id,
245        verification_digest: inputs.verification.verification_digest,
246        quiescence_digest: inputs.quiescence.quiescence_digest,
247        source_high_water_mark: inputs
248            .quiescence
249            .stable_source_high_water_mark
250            .expect("validated"),
251        destination_checkpoint: inputs.quiescence.destination_checkpoint.expect("validated"),
252        route: ExtractionTrafficRoute::CandidateVerificationOnly,
253        external_mutations_paused: true,
254        linked_mutations_open: false,
255        linked_authoritative: true,
256        candidate_authoritative: false,
257        candidate_healthy: true,
258        declared_verification_traffic_only: true,
259        verification_effects_isolated: true,
260        linked_business_probe_passed: false,
261        apply_receipts: vec![apply_receipt],
262        rollback_receipts: Vec::new(),
263        evidence: vec![evidence(
264            "provisional_routing",
265            "candidate-verification-only",
266            &identity,
267            "Only declared read-only, recorded, or isolated synthetic verification traffic is routed to the candidate.",
268        )],
269    };
270    refresh(&mut run);
271    Ok(run)
272}
273
274#[must_use]
275pub fn verify_provisional_cutover(
276    mut run: ExtractionProvisionalCutoverRun,
277    audit_identity: &str,
278) -> ExtractionProvisionalCutoverRun {
279    if run.status == ExtractionProvisionalCutoverStatus::Provisional {
280        run.status = ExtractionProvisionalCutoverStatus::Verified;
281        run.evidence.push(evidence(
282            "provisional_verification",
283            audit_identity,
284            &run.verification_digest,
285            "Candidate provisional verification passed while external mutations remained paused.",
286        ));
287        run.revision += 1;
288        refresh(&mut run);
289    }
290    run
291}
292
293#[must_use]
294pub fn fail_provisional_cutover(
295    mut run: ExtractionProvisionalCutoverRun,
296    failure: &str,
297    audit_identity: &str,
298    validation: ExtractionLinkedRollbackValidation,
299) -> ExtractionProvisionalCutoverRun {
300    if run.status == ExtractionProvisionalCutoverStatus::RolledBack {
301        return run;
302    }
303    let restore_routing = receipt(
304        "restore-linked-routing",
305        &run.routing_revision_current,
306        &run.routing_revision_before,
307        "restored",
308    );
309    let validation_passed = validation.is_valid_for(&run) && validation.passed;
310    run.rollback_receipts = vec![restore_routing];
311    if validation_passed {
312        run.rollback_receipts.push(receipt(
313            "reopen-linked-mutations",
314            "paused",
315            "open",
316            "validated",
317        ));
318    }
319    run.status = ExtractionProvisionalCutoverStatus::RolledBack;
320    run.routing_revision_current = run.routing_revision_before.clone();
321    run.route = ExtractionTrafficRoute::Linked;
322    run.external_mutations_paused = !validation_passed;
323    run.linked_mutations_open = validation_passed;
324    run.linked_authoritative = true;
325    run.candidate_authoritative = false;
326    run.linked_business_probe_passed = validation_passed;
327    run.evidence.push(evidence(
328        "rollback",
329        audit_identity,
330        failure,
331        &format!(
332            "Candidate verification failed: {failure}. Linked routing and authority were restored without reverse data movement."
333        ),
334    ));
335    run.revision += 1;
336    refresh(&mut run);
337    run
338}
339
340#[must_use]
341pub fn complete_provisional_rollback_validation(
342    mut run: ExtractionProvisionalCutoverRun,
343    validation: ExtractionLinkedRollbackValidation,
344) -> ExtractionProvisionalCutoverRun {
345    if run.status == ExtractionProvisionalCutoverStatus::RolledBack
346        && run.external_mutations_paused
347        && validation.is_valid_for(&run)
348        && validation.passed
349    {
350        run.rollback_receipts.push(receipt(
351            "reopen-linked-mutations",
352            "paused",
353            "open",
354            "validated",
355        ));
356        run.external_mutations_paused = false;
357        run.linked_mutations_open = true;
358        run.linked_business_probe_passed = true;
359        run.evidence.push(evidence(
360            "linked_business_validation",
361            "linked-authority",
362            &validation.business_probe_digest,
363            "Linked business behavior passed before mutations reopened.",
364        ));
365        run.revision += 1;
366        refresh(&mut run);
367    }
368    run
369}
370
371fn rollback_validation_digest(validation: &ExtractionLinkedRollbackValidation) -> String {
372    let mut value = validation.clone();
373    value.validation_digest.clear();
374    digest(&value)
375}
376
377fn error(
378    code: ExtractionProvisionalCutoverIssueCode,
379    message: &str,
380    next_action: &str,
381) -> ExtractionProvisionalCutoverError {
382    ExtractionProvisionalCutoverError {
383        code,
384        message: message.to_owned(),
385        next_actions: vec![next_action.to_owned()],
386    }
387}
388
389fn receipt(
390    step_id: &str,
391    from_revision: &str,
392    to_revision: &str,
393    outcome: &str,
394) -> ExtractionCutoverReceipt {
395    let step_digest = digest(&(step_id, from_revision, to_revision, outcome));
396    ExtractionCutoverReceipt {
397        step_id: step_id.to_owned(),
398        step_digest,
399        from_revision: from_revision.to_owned(),
400        to_revision: to_revision.to_owned(),
401        outcome: outcome.to_owned(),
402    }
403}
404
405fn evidence(kind: &str, subject: &str, value: &str, detail: &str) -> ExtractionCutoverEvidence {
406    ExtractionCutoverEvidence {
407        kind: kind.to_owned(),
408        subject: subject.to_owned(),
409        digest: extraction_input_digest(value.as_bytes()),
410        detail: detail.to_owned(),
411    }
412}
413
414fn refresh(run: &mut ExtractionProvisionalCutoverRun) {
415    run.cutover_digest.clear();
416    run.cutover_digest = digest(run);
417}
418
419#[must_use]
420pub fn extraction_provisional_cutover_integrity_is_valid(
421    run: &ExtractionProvisionalCutoverRun,
422) -> bool {
423    if run.protocol != EXTRACTION_PROVISIONAL_CUTOVER_PROTOCOL {
424        return false;
425    }
426    let mut value = run.clone();
427    value.cutover_digest.clear();
428    run.cutover_digest == digest(&value)
429}
430
431fn digest(value: &impl Serialize) -> String {
432    extraction_input_digest(&serde_json::to_vec(value).expect("Cutover values serialize"))
433}