1use crate::{
2 ExtractionProvisionalCutoverRun, ExtractionProvisionalCutoverStatus, ExtractionQuiescenceRun,
3 ExtractionQuiescenceStatus, ExtractionReconciliationResult, ExtractionReconciliationStatus,
4 ExtractionVerificationResult, ExtractionVerificationStatus, extraction_input_digest,
5 extraction_provisional_cutover_integrity_is_valid, extraction_quiescence_integrity_is_valid,
6 extraction_reconciliation_integrity_is_valid, extraction_verification_integrity_is_valid,
7};
8use schemars::JsonSchema;
9use serde::{Deserialize, Serialize};
10use std::fmt;
11
12pub const EXTRACTION_AUTHORITY_COMMIT_PROTOCOL: &str = "lenso.extraction-authority-commit.v1";
13pub const EXTRACTION_CANDIDATE_HEALTH_PROTOCOL: &str = "lenso.extraction-candidate-health.v1";
14
15#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
16#[serde(rename_all = "camelCase")]
17pub struct ExtractionCandidateHealthEvidence {
18 pub protocol: String,
19 pub evidence_id: String,
20 pub evidence_digest: String,
21 pub plan_id: String,
22 pub candidate_service_id: String,
23 pub endpoint: String,
24 pub endpoint_reachable: bool,
25 pub store_ready: bool,
26 pub healthy: bool,
27}
28
29impl ExtractionCandidateHealthEvidence {
30 #[must_use]
31 pub fn bind(
32 plan_id: impl Into<String>,
33 candidate_service_id: impl Into<String>,
34 endpoint: impl Into<String>,
35 endpoint_reachable: bool,
36 store_ready: bool,
37 ) -> Self {
38 let mut evidence = Self {
39 protocol: EXTRACTION_CANDIDATE_HEALTH_PROTOCOL.to_owned(),
40 evidence_id: String::new(),
41 evidence_digest: String::new(),
42 plan_id: plan_id.into(),
43 candidate_service_id: candidate_service_id.into(),
44 endpoint: endpoint.into(),
45 endpoint_reachable,
46 store_ready,
47 healthy: endpoint_reachable && store_ready,
48 };
49 let identity = digest(&(
50 evidence.plan_id.as_str(),
51 evidence.candidate_service_id.as_str(),
52 evidence.endpoint.as_str(),
53 ));
54 evidence.evidence_id = format!("extraction-candidate-health:{identity}");
55 evidence.evidence_digest = candidate_health_digest(&evidence);
56 evidence
57 }
58}
59
60#[must_use]
61pub fn extraction_candidate_health_integrity_is_valid(
62 evidence: &ExtractionCandidateHealthEvidence,
63) -> bool {
64 evidence.protocol == EXTRACTION_CANDIDATE_HEALTH_PROTOCOL
65 && !evidence.evidence_id.trim().is_empty()
66 && !evidence.plan_id.trim().is_empty()
67 && !evidence.candidate_service_id.trim().is_empty()
68 && !evidence.endpoint.trim().is_empty()
69 && evidence.healthy == (evidence.endpoint_reachable && evidence.store_ready)
70 && evidence.evidence_digest == candidate_health_digest(evidence)
71}
72
73#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
74#[serde(rename_all = "camelCase")]
75pub struct ExtractionApproval {
76 pub approval_id: String,
77 pub approval_digest: String,
78 pub approver: String,
79 pub authorized: bool,
80 pub cutover_id: String,
81 pub cutover_digest: String,
82 pub plan_digest: String,
83 pub authority_revision: String,
84 pub destination_checkpoint: String,
85 pub verification_digest: String,
86 pub quiescence_digest: String,
87 pub candidate_service_id: String,
88 pub candidate_health_digest: String,
89}
90
91impl ExtractionApproval {
92 #[must_use]
93 pub fn bind(
94 cutover: &ExtractionProvisionalCutoverRun,
95 candidate_health: &ExtractionCandidateHealthEvidence,
96 approval_id: impl Into<String>,
97 approver: impl Into<String>,
98 authorized: bool,
99 ) -> Self {
100 let mut approval = Self {
101 approval_id: approval_id.into(),
102 approval_digest: String::new(),
103 approver: approver.into(),
104 authorized,
105 cutover_id: cutover.cutover_id.clone(),
106 cutover_digest: cutover.cutover_digest.clone(),
107 plan_digest: cutover.plan_digest.clone(),
108 authority_revision: cutover.authority_revision.clone(),
109 destination_checkpoint: cutover.destination_checkpoint.clone(),
110 verification_digest: cutover.verification_digest.clone(),
111 quiescence_digest: cutover.quiescence_digest.clone(),
112 candidate_service_id: cutover.candidate_service_id.clone(),
113 candidate_health_digest: candidate_health.evidence_digest.clone(),
114 };
115 approval.approval_digest = approval_digest(&approval);
116 approval
117 }
118}
119
120#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
121#[serde(rename_all = "camelCase")]
122pub struct ExtractionAuthorityCommitInputs {
123 pub cutover: ExtractionProvisionalCutoverRun,
124 pub approval: ExtractionApproval,
125 pub current_authority_revision: String,
126 pub current_routing_revision: String,
127 pub current_system_graph_revision: String,
128 pub revalidation: ExtractionAuthorityCommitRevalidation,
129}
130
131pub trait ExtractionApprovalVerifier: Send + Sync {
132 fn verify(&self, approval: &ExtractionApproval) -> bool;
133}
134
135#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
136#[serde(rename_all = "camelCase")]
137pub struct ExtractionTopologyState {
138 pub authority_revision: String,
139 pub routing_revision: String,
140 pub system_graph_revision: String,
141 pub authority_kind: String,
142 pub owner_id: String,
143}
144
145#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
146#[serde(rename_all = "camelCase")]
147pub struct ExtractionAuthorityCommitRevalidation {
148 pub reconciliation: ExtractionReconciliationResult,
149 pub verification: ExtractionVerificationResult,
150 pub quiescence: ExtractionQuiescenceRun,
151 pub candidate_health: ExtractionCandidateHealthEvidence,
152}
153
154#[derive(
155 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
156)]
157#[serde(rename_all = "snake_case")]
158pub enum ExtractionAuthorityCommitStatus {
159 Committed,
160}
161
162#[derive(
163 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
164)]
165#[serde(rename_all = "snake_case")]
166pub enum ExtractionAuthorityCommitErrorCode {
167 ProvisionalVerificationIncomplete,
168 ApprovalUnauthorized,
169 ApprovalInvalid,
170 ApprovalStale,
171 AuthorityChanged,
172 RoutingChanged,
173 CandidateUnhealthy,
174 FinalStateInvalid,
175 ConcurrentStateChange,
176 PersistenceFailed,
177}
178
179#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
180#[serde(rename_all = "camelCase")]
181pub struct ExtractionAuthorityCommitError {
182 pub code: ExtractionAuthorityCommitErrorCode,
183 pub message: String,
184 pub next_actions: Vec<String>,
185 pub mutation_started: bool,
186}
187
188impl fmt::Display for ExtractionAuthorityCommitError {
189 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
190 formatter.write_str(&self.message)
191 }
192}
193
194impl std::error::Error for ExtractionAuthorityCommitError {}
195
196#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
197#[serde(rename_all = "camelCase")]
198pub struct ExtractionAuthorityCommitReceipt {
199 pub receipt_id: String,
200 pub receipt_digest: String,
201 pub expected_authority_revision: String,
202 pub expected_routing_revision: String,
203 pub expected_system_graph_revision: String,
204 pub committed_authority_revision: String,
205 pub committed_routing_revision: String,
206 pub committed_system_graph_revision: String,
207 pub candidate_service_id: String,
208 pub outcome: String,
209}
210
211#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
212#[serde(rename_all = "camelCase")]
213pub struct ExtractionAuthorityCommitResult {
214 pub protocol: String,
215 pub commit_id: String,
216 pub commit_digest: String,
217 pub status: ExtractionAuthorityCommitStatus,
218 pub plan_digest: String,
219 pub approval: ExtractionApproval,
220 pub authority_revision: String,
221 pub routing_revision: String,
222 pub system_graph_revision: String,
223 pub candidate_service_id: String,
224 pub candidate_authoritative: bool,
225 pub linked_authoritative: bool,
226 pub candidate_mutations_open: bool,
227 pub linked_recovery_read_only: bool,
228 pub source_cleanup_performed: bool,
229 #[serde(default)]
230 pub autonomous_mutation_ids: Vec<String>,
231 pub fast_rollback_blocked: bool,
232 pub business_execution_requires_runtime_console: bool,
233 pub business_execution_requires_system_plane: bool,
234 pub commit_receipts: Vec<ExtractionAuthorityCommitReceipt>,
235}
236
237#[derive(
238 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
239)]
240#[serde(rename_all = "snake_case")]
241pub enum ExtractionFastRollbackIssueCode {
242 ReverseMigrationEvidenceRequired,
243}
244
245#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
246#[serde(rename_all = "camelCase")]
247pub struct ExtractionFastRollbackError {
248 pub code: ExtractionFastRollbackIssueCode,
249 pub message: String,
250 pub next_actions: Vec<String>,
251}
252
253impl fmt::Display for ExtractionFastRollbackError {
254 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
255 formatter.write_str(&self.message)
256 }
257}
258
259impl std::error::Error for ExtractionFastRollbackError {}
260
261#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
262#[serde(rename_all = "camelCase")]
263pub struct ExtractionReverseMigrationEvidence {
264 pub plan_digest: String,
265 pub reconciliation_digest: String,
266 pub reviewed_by: String,
267 pub approved: bool,
268}
269
270pub fn commit_extraction_authority(
271 inputs: ExtractionAuthorityCommitInputs,
272) -> Result<ExtractionAuthorityCommitResult, ExtractionAuthorityCommitError> {
273 let cutover = &inputs.cutover;
274 if !extraction_provisional_cutover_integrity_is_valid(cutover)
275 || cutover.status != ExtractionProvisionalCutoverStatus::Verified
276 || !cutover.external_mutations_paused
277 || !cutover.linked_authoritative
278 || cutover.candidate_authoritative
279 {
280 return Err(error(
281 ExtractionAuthorityCommitErrorCode::ProvisionalVerificationIncomplete,
282 "Provisional Cutover is not verified in a single-authority paused state.",
283 "Repeat provisional verification before requesting approval.",
284 ));
285 }
286 if !inputs.approval.authorized || inputs.approval.approver.trim().is_empty() {
287 return Err(error(
288 ExtractionAuthorityCommitErrorCode::ApprovalUnauthorized,
289 "The Approval Boundary was not crossed by an authorized identity.",
290 "Request approval from an authorized operator.",
291 ));
292 }
293 if inputs.approval.approval_digest != approval_digest(&inputs.approval) {
294 return Err(error(
295 ExtractionAuthorityCommitErrorCode::ApprovalInvalid,
296 "Approval integrity validation failed.",
297 "Discard the changed approval and bind a new one.",
298 ));
299 }
300 if inputs.approval.cutover_id != cutover.cutover_id
301 || inputs.approval.cutover_digest != cutover.cutover_digest
302 || inputs.approval.plan_digest != cutover.plan_digest
303 || inputs.approval.authority_revision != cutover.authority_revision
304 || inputs.approval.destination_checkpoint != cutover.destination_checkpoint
305 || inputs.approval.verification_digest != cutover.verification_digest
306 || inputs.approval.quiescence_digest != cutover.quiescence_digest
307 || inputs.approval.candidate_service_id != cutover.candidate_service_id
308 {
309 return Err(error(
310 ExtractionAuthorityCommitErrorCode::ApprovalStale,
311 "Approval does not match the exact verified Cutover state.",
312 "Bind a fresh approval to the current Cutover evidence.",
313 ));
314 }
315 if inputs.current_authority_revision != cutover.authority_revision {
316 return Err(error(
317 ExtractionAuthorityCommitErrorCode::AuthorityChanged,
318 "Authority revision changed before commit.",
319 "Regenerate Cutover evidence from the current authority revision.",
320 ));
321 }
322 if inputs.current_routing_revision != cutover.routing_revision_current {
323 return Err(error(
324 ExtractionAuthorityCommitErrorCode::RoutingChanged,
325 "Routing revision changed before commit.",
326 "Restore the verified provisional route or restart Cutover.",
327 ));
328 }
329 let revalidation = &inputs.revalidation;
330 let drain_complete = revalidation.quiescence.drain.as_ref().is_some_and(|drain| {
331 drain.in_flight_requests == 0
332 && drain.outbox_messages == 0
333 && drain.inbox_messages == 0
334 && drain.scheduled_functions == 0
335 && drain.timers == 0
336 && drain.durable_workflows == 0
337 && drain.unresolved.is_empty()
338 && !drain.timed_out
339 });
340 if !extraction_reconciliation_integrity_is_valid(&revalidation.reconciliation)
341 || !extraction_verification_integrity_is_valid(&revalidation.verification)
342 || !extraction_quiescence_integrity_is_valid(&revalidation.quiescence)
343 || revalidation.reconciliation.plan_digest != cutover.plan_digest
344 || revalidation.reconciliation.destination_checkpoint != cutover.destination_checkpoint
345 || revalidation.verification.verification_digest != cutover.verification_digest
346 || revalidation.verification.reconciliation_id
347 != revalidation.reconciliation.reconciliation_id
348 || revalidation.verification.reconciliation_digest
349 != revalidation.reconciliation.reconciliation_digest
350 || revalidation.quiescence.plan_digest != cutover.plan_digest
351 || revalidation.quiescence.expected_authority_revision != cutover.authority_revision
352 || revalidation.quiescence.destination_checkpoint.as_deref()
353 != Some(cutover.destination_checkpoint.as_str())
354 || revalidation.quiescence.quiescence_digest != cutover.quiescence_digest
355 {
356 return Err(error(
357 ExtractionAuthorityCommitErrorCode::ApprovalStale,
358 "Final revalidation pins do not match the approved Cutover state.",
359 "Repeat final revalidation and bind a fresh approval.",
360 ));
361 }
362 if revalidation.reconciliation.status != ExtractionReconciliationStatus::Matched
363 || !revalidation.reconciliation.issues.is_empty()
364 || revalidation.verification.status != ExtractionVerificationStatus::Verified
365 || !revalidation.verification.issues.is_empty()
366 || revalidation
367 .verification
368 .compatibility
369 .iter()
370 .any(|evidence| !evidence.compatible)
371 || revalidation
372 .verification
373 .policy
374 .iter()
375 .any(|evidence| !evidence.passed)
376 || revalidation.quiescence.status != ExtractionQuiescenceStatus::Quiesced
377 || !revalidation.quiescence.linked_mutations_paused
378 || !drain_complete
379 {
380 return Err(error(
381 ExtractionAuthorityCommitErrorCode::FinalStateInvalid,
382 "Final Cutover safety revalidation failed before commit.",
383 "Restore quiescence, drain, reconciliation, compatibility, and policy evidence.",
384 ));
385 }
386 if !extraction_candidate_health_integrity_is_valid(&revalidation.candidate_health)
387 || revalidation.candidate_health.plan_id != cutover.plan_id
388 || revalidation.candidate_health.candidate_service_id != cutover.candidate_service_id
389 || !revalidation.candidate_health.healthy
390 || !cutover.candidate_healthy
391 {
392 return Err(error(
393 ExtractionAuthorityCommitErrorCode::CandidateUnhealthy,
394 "Candidate health changed before commit.",
395 "Restore candidate health and repeat verification.",
396 ));
397 }
398 if inputs.approval.candidate_health_digest != revalidation.candidate_health.evidence_digest {
399 return Err(error(
400 ExtractionAuthorityCommitErrorCode::ApprovalStale,
401 "Approval does not bind the persisted candidate health evidence.",
402 "Probe candidate health and bind a fresh approval to that exact evidence.",
403 ));
404 }
405 let commit_identity = digest(&(
406 inputs.approval.approval_digest.as_str(),
407 inputs.current_authority_revision.as_str(),
408 inputs.current_routing_revision.as_str(),
409 inputs.current_system_graph_revision.as_str(),
410 ));
411 let authority_revision = format!("autonomous-authority:{commit_identity}");
412 let routing_revision = format!("autonomous-routing:{commit_identity}");
413 let system_graph_revision = format!("autonomous-system-graph:{commit_identity}");
414 let mut receipt = ExtractionAuthorityCommitReceipt {
415 receipt_id: format!("extraction-authority-commit-receipt:{commit_identity}"),
416 receipt_digest: String::new(),
417 expected_authority_revision: inputs.current_authority_revision,
418 expected_routing_revision: inputs.current_routing_revision,
419 expected_system_graph_revision: inputs.current_system_graph_revision,
420 committed_authority_revision: authority_revision.clone(),
421 committed_routing_revision: routing_revision.clone(),
422 committed_system_graph_revision: system_graph_revision.clone(),
423 candidate_service_id: cutover.candidate_service_id.clone(),
424 outcome: "committed_single_compare_and_set".to_owned(),
425 };
426 receipt.receipt_digest = digest(&receipt_without_digest(&receipt));
427 let mut result = ExtractionAuthorityCommitResult {
428 protocol: EXTRACTION_AUTHORITY_COMMIT_PROTOCOL.to_owned(),
429 commit_id: format!("extraction-authority-commit:{commit_identity}"),
430 commit_digest: String::new(),
431 status: ExtractionAuthorityCommitStatus::Committed,
432 plan_digest: cutover.plan_digest.clone(),
433 approval: inputs.approval,
434 authority_revision,
435 routing_revision,
436 system_graph_revision,
437 candidate_service_id: cutover.candidate_service_id.clone(),
438 candidate_authoritative: true,
439 linked_authoritative: false,
440 candidate_mutations_open: true,
441 linked_recovery_read_only: true,
442 source_cleanup_performed: false,
443 autonomous_mutation_ids: Vec::new(),
444 fast_rollback_blocked: false,
445 business_execution_requires_runtime_console: false,
446 business_execution_requires_system_plane: false,
447 commit_receipts: vec![receipt],
448 };
449 refresh(&mut result);
450 Ok(result)
451}
452
453pub async fn commit_extraction_authority_postgres(
457 pool: &sqlx::PgPool,
458 mut inputs: ExtractionAuthorityCommitInputs,
459 verifier: &dyn ExtractionApprovalVerifier,
460) -> Result<ExtractionAuthorityCommitResult, ExtractionAuthorityCommitError> {
461 if !verifier.verify(&inputs.approval) {
462 return Err(error(
463 ExtractionAuthorityCommitErrorCode::ApprovalUnauthorized,
464 "The configured Approval Authority rejected this approval.",
465 "Obtain a fresh approval through the protected workflow.",
466 ));
467 }
468 let mut tx = pool.begin().await.map_err(persistence_error)?;
469 sqlx::query("select pg_advisory_xact_lock(hashtext('lenso-extraction-topology'))")
470 .execute(&mut *tx)
471 .await
472 .map_err(persistence_error)?;
473 inputs.revalidation.reconciliation = load_commit_artifact(
474 &mut tx,
475 &inputs.cutover.plan_id,
476 "lenso.extraction-reconciliation.v1",
477 )
478 .await?;
479 inputs.revalidation.verification = load_commit_artifact(
480 &mut tx,
481 &inputs.cutover.plan_id,
482 "lenso.extraction-verification.v1",
483 )
484 .await?;
485 inputs.revalidation.quiescence = load_commit_artifact(
486 &mut tx,
487 &inputs.cutover.plan_id,
488 "lenso.extraction-quiescence.v1",
489 )
490 .await?;
491 inputs.revalidation.candidate_health = load_commit_artifact(
492 &mut tx,
493 &inputs.cutover.plan_id,
494 EXTRACTION_CANDIDATE_HEALTH_PROTOCOL,
495 )
496 .await?;
497 let result = commit_extraction_authority(inputs.clone())?;
498 sqlx::raw_sql(
499 r#"
500 create schema if not exists lenso_extraction;
501 create table if not exists lenso_extraction.authority_states (
502 state_id text primary key,
503 authority_revision text not null,
504 routing_revision text not null,
505 system_graph_revision text not null,
506 authority_kind text not null,
507 owner_id text not null,
508 updated_at timestamptz not null default now()
509 );
510 create table if not exists lenso_extraction.authority_commits (
511 approval_digest text primary key,
512 plan_digest text not null,
513 result_json jsonb not null,
514 committed_at timestamptz not null default now()
515 );
516 "#,
517 )
518 .execute(&mut *tx)
519 .await
520 .map_err(persistence_error)?;
521 if let Some(value) = sqlx::query_scalar::<_, serde_json::Value>(
522 "select result_json from lenso_extraction.authority_commits where approval_digest = $1",
523 )
524 .bind(&inputs.approval.approval_digest)
525 .fetch_optional(&mut *tx)
526 .await
527 .map_err(persistence_error)?
528 {
529 let _: ExtractionAuthorityCommitResult =
530 serde_json::from_value(value).map_err(|source| {
531 error(
532 ExtractionAuthorityCommitErrorCode::PersistenceFailed,
533 format!("Stored authority commit is unreadable: {source}"),
534 "Repair or restore the last valid authority commit receipt.",
535 )
536 })?;
537 return Err(error(
538 ExtractionAuthorityCommitErrorCode::ApprovalStale,
539 "This approval has already been committed and cannot be replayed.",
540 "Inspect the persisted commit receipt and current topology state.",
541 ));
542 }
543 let persisted = sqlx::query_as::<_, (String, String, String, String)>(
544 "select authority_revision, routing_revision, system_graph_revision, authority_kind from lenso_extraction.authority_states where state_id = 'system' for update",
545 )
546 .fetch_optional(&mut *tx)
547 .await
548 .map_err(persistence_error)?
549 .ok_or_else(|| error(
550 ExtractionAuthorityCommitErrorCode::FinalStateInvalid,
551 "Authoritative topology state has not been initialized by the composition root.",
552 "Install the linked authority, routing, and System graph state before Cutover.",
553 ))?;
554 if persisted
555 != (
556 inputs.current_authority_revision.clone(),
557 inputs.current_routing_revision.clone(),
558 inputs.current_system_graph_revision.clone(),
559 "linked".to_owned(),
560 )
561 {
562 return Err(error(
563 ExtractionAuthorityCommitErrorCode::ConcurrentStateChange,
564 "Persisted authority, routing, or System graph state changed before commit.",
565 "Reload persisted topology state and repeat Cutover verification.",
566 ));
567 }
568 let changed = sqlx::query(
569 r#"
570 update lenso_extraction.authority_states
571 set authority_revision = $2, routing_revision = $3, system_graph_revision = $4,
572 authority_kind = 'autonomous', owner_id = $5, updated_at = now()
573 where state_id = 'system' and authority_revision = $6 and routing_revision = $7
574 and system_graph_revision = $8 and authority_kind = 'linked'
575 "#,
576 )
577 .bind("system")
578 .bind(&result.authority_revision)
579 .bind(&result.routing_revision)
580 .bind(&result.system_graph_revision)
581 .bind(&result.candidate_service_id)
582 .bind(&result.commit_receipts[0].expected_authority_revision)
583 .bind(&result.commit_receipts[0].expected_routing_revision)
584 .bind(&result.commit_receipts[0].expected_system_graph_revision)
585 .execute(&mut *tx)
586 .await
587 .map_err(persistence_error)?;
588 if changed.rows_affected() != 1 {
589 return Err(error(
590 ExtractionAuthorityCommitErrorCode::ConcurrentStateChange,
591 "Atomic authority compare-and-set lost a concurrent race.",
592 "Reload topology state; do not retry with stale approval evidence.",
593 ));
594 }
595 sqlx::query(
596 "insert into lenso_extraction.authority_commits (approval_digest, plan_digest, result_json) values ($1, $2, $3)",
597 )
598 .bind(&result.approval.approval_digest)
599 .bind(&result.plan_digest)
600 .bind(serde_json::to_value(&result).map_err(|source| {
601 error(
602 ExtractionAuthorityCommitErrorCode::PersistenceFailed,
603 format!("Authority commit could not serialize: {source}"),
604 "Abort before committing topology state.",
605 )
606 })?)
607 .execute(&mut *tx)
608 .await
609 .map_err(persistence_error)?;
610 tx.commit().await.map_err(persistence_error)?;
611 Ok(result)
612}
613
614async fn load_commit_artifact<T: serde::de::DeserializeOwned>(
615 tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
616 plan_id: &str,
617 protocol: &str,
618) -> Result<T, ExtractionAuthorityCommitError> {
619 let value = sqlx::query_scalar::<_, serde_json::Value>(
620 "select artifact_json from platform.extraction_artifacts where plan_id = $1 and protocol = $2 order by recorded_at desc, artifact_id desc limit 1",
621 )
622 .bind(plan_id)
623 .bind(protocol)
624 .fetch_optional(&mut **tx)
625 .await
626 .map_err(persistence_error)?
627 .ok_or_else(|| {
628 error(
629 ExtractionAuthorityCommitErrorCode::FinalStateInvalid,
630 format!("Required persisted final artifact `{protocol}` is missing."),
631 "Persist fresh reconciliation, verification, and quiescence artifacts before approval commit.",
632 )
633 })?;
634 serde_json::from_value(value).map_err(|source| {
635 error(
636 ExtractionAuthorityCommitErrorCode::FinalStateInvalid,
637 format!("Persisted final artifact `{protocol}` is unreadable: {source}"),
638 "Regenerate and persist an integrity-valid final artifact.",
639 )
640 })
641}
642
643pub async fn initialize_extraction_topology_state(
644 pool: &sqlx::PgPool,
645 state: &ExtractionTopologyState,
646) -> Result<(), ExtractionAuthorityCommitError> {
647 sqlx::raw_sql(
648 r#"
649 create schema if not exists lenso_extraction;
650 create table if not exists lenso_extraction.authority_states (
651 state_id text primary key,
652 authority_revision text not null,
653 routing_revision text not null,
654 system_graph_revision text not null,
655 authority_kind text not null,
656 owner_id text not null,
657 updated_at timestamptz not null default now()
658 );
659 "#,
660 )
661 .execute(pool)
662 .await
663 .map_err(persistence_error)?;
664 sqlx::query(
665 "insert into lenso_extraction.authority_states (state_id, authority_revision, routing_revision, system_graph_revision, authority_kind, owner_id) values ('system',$1,$2,$3,$4,$5) on conflict (state_id) do nothing",
666 )
667 .bind(&state.authority_revision)
668 .bind(&state.routing_revision)
669 .bind(&state.system_graph_revision)
670 .bind(&state.authority_kind)
671 .bind(&state.owner_id)
672 .execute(pool)
673 .await
674 .map_err(persistence_error)?;
675 Ok(())
676}
677
678fn persistence_error(source: sqlx::Error) -> ExtractionAuthorityCommitError {
679 error(
680 ExtractionAuthorityCommitErrorCode::PersistenceFailed,
681 format!("PostgreSQL authority commit failed: {source}"),
682 "Restore PostgreSQL and reload the persisted topology state.",
683 )
684}
685
686#[must_use]
687pub fn record_autonomous_mutation(
688 mut result: ExtractionAuthorityCommitResult,
689 mutation_id: impl Into<String>,
690) -> ExtractionAuthorityCommitResult {
691 let mutation_id = mutation_id.into();
692 if !result.autonomous_mutation_ids.contains(&mutation_id) {
693 result.autonomous_mutation_ids.push(mutation_id);
694 result.autonomous_mutation_ids.sort();
695 }
696 result.fast_rollback_blocked = !result.autonomous_mutation_ids.is_empty();
697 refresh(&mut result);
698 result
699}
700
701pub fn request_fast_extraction_rollback(
702 result: &ExtractionAuthorityCommitResult,
703 evidence: Option<&ExtractionReverseMigrationEvidence>,
704) -> Result<(), ExtractionFastRollbackError> {
705 let reviewed = evidence.is_some_and(|evidence| {
706 evidence.approved
707 && !evidence.reviewed_by.trim().is_empty()
708 && evidence.plan_digest == result.plan_digest
709 && !evidence.reconciliation_digest.trim().is_empty()
710 });
711 if result.fast_rollback_blocked && !reviewed {
712 return Err(ExtractionFastRollbackError {
713 code: ExtractionFastRollbackIssueCode::ReverseMigrationEvidenceRequired,
714 message: "Fast rollback is blocked after Autonomous writes began.".to_owned(),
715 next_actions: vec!["Review a reverse-migration and reconciliation plan before changing authority again.".to_owned()],
716 });
717 }
718 Ok(())
719}
720
721fn error(
722 code: ExtractionAuthorityCommitErrorCode,
723 message: impl Into<String>,
724 next_action: impl Into<String>,
725) -> ExtractionAuthorityCommitError {
726 ExtractionAuthorityCommitError {
727 code,
728 message: message.into(),
729 next_actions: vec![next_action.into()],
730 mutation_started: false,
731 }
732}
733
734fn approval_digest(approval: &ExtractionApproval) -> String {
735 let mut value = approval.clone();
736 value.approval_digest.clear();
737 digest(&value)
738}
739
740fn candidate_health_digest(evidence: &ExtractionCandidateHealthEvidence) -> String {
741 let mut value = evidence.clone();
742 value.evidence_digest.clear();
743 digest(&value)
744}
745
746fn receipt_without_digest(
747 receipt: &ExtractionAuthorityCommitReceipt,
748) -> ExtractionAuthorityCommitReceipt {
749 let mut value = receipt.clone();
750 value.receipt_digest.clear();
751 value
752}
753
754fn refresh(result: &mut ExtractionAuthorityCommitResult) {
755 result.commit_digest.clear();
756 result.commit_digest = digest(result);
757}
758
759fn digest(value: &impl Serialize) -> String {
760 extraction_input_digest(&serde_json::to_vec(value).expect("Authority commit values serialize"))
761}