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}