Skip to main content

lenso_service/
extraction_authority_commit.rs

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
453/// Atomically compares and changes the persisted authority, routing, and
454/// System graph revisions. The receipt and committed state share one database
455/// transaction, so a partial topology transfer cannot escape.
456pub 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}