Skip to main content

chio_store_sqlite/
admission_operation_store.rs

1use std::sync::{Arc, Mutex, MutexGuard};
2use std::time::{SystemTime, UNIX_EPOCH};
3
4use chio_core::canonical::canonical_json_bytes;
5use chio_core::economic_continuity::VerifiedEconomicStateBatchAdvance;
6use chio_core::receipt::{body::ChioReceipt, lineage::ChildRequestReceipt};
7use chio_core::{sha256_hex, StoreMutationFence};
8#[cfg(test)]
9use chio_credit::obligation::CreditExposureReservationRequest;
10use chio_credit::obligation::{
11    CreditExposureReservationRecordV1, ObligationAtomV1, ObligationDispositionRecordV1,
12    ObligationSettlementLifecycleV1,
13};
14use chio_kernel::admission_operation::{
15    AdmissionAttachment, AdmissionBeginResult, AdmissionCaptureError, AdmissionCommandResult,
16    AdmissionDigest, AdmissionIdentifier, AdmissionOperationCommand, AdmissionOperationError,
17    AdmissionOperationId, AdmissionOperationKind, AdmissionOperationState, AdmissionOperationStore,
18    AdmissionOperationStoreError, AdmissionOperationV1, AdmissionProjectionCapabilities,
19    AdmissionProjectionContext, AdmissionProjectionManifestV1, AdmissionProjectionRecordKind,
20    AdmissionRecoveryLease, AdmissionReplayClassification, AdmissionReplayKey, AdmissionTerminal,
21    AdmissionTerminalProjection, AdmissionTerminalReplay, CanonicalAdmissionProjectionRecord,
22    CanonicalAdmissionTerminalProjection, PersistedAdmissionOperationV1,
23    QualifiedAdmissionOperationStore, SideEffectClass, SignedAdmissionTerminalProjectionV1,
24    UntrustedAdmissionRecoveryClaim, VerifiedAdmissionTerminalProjectionRecordV1,
25    VerifiedAdmissionTerminalProjectionV1,
26};
27use chio_kernel::budget_store::{
28    BudgetAuthorizeHoldDecision, BudgetAuthorizeHoldRequest, BudgetCaptureInvocationRequest,
29    BudgetReconcileHoldRequest, BudgetStoreError,
30};
31use chio_kernel::payment::{PaymentJournalRecord, PaymentJournalTransition};
32use chio_kernel::receipt_store::{
33    AdmissionPaymentJournalAdvance, AdmissionPaymentJournalError, AdmissionPaymentSettlement,
34    AdmissionPaymentSettlementBegin, AuthorizationReceiptConsumption, PendingSettlementObservation,
35    ReceiptStore, ReceiptStoreError,
36};
37use rusqlite::{params, Connection, OptionalExtension, Row, Transaction, TransactionBehavior};
38use serde::{Deserialize, Serialize};
39
40use crate::serving_owner::{SqliteServingOwner, SqliteServingOwnerError};
41
42mod commit_chain;
43mod credit_exposure;
44mod errors;
45mod factor_assignment;
46mod obligation;
47mod participant;
48mod projection;
49mod schema;
50mod store;
51mod threshold_approval;
52
53use commit_chain::append_operation_commit;
54pub(crate) use commit_chain::{
55    append_operation_commit_with_participant, load_admission_commit_head,
56    verify_admission_commit_chain, verify_admission_commit_suffix, AdmissionCommitHead,
57    GENESIS_CHAIN_DIGEST,
58};
59pub use credit_exposure::CreditExposureAccountSnapshot;
60pub(crate) use credit_exposure::{
61    apply_credit_exposure_terminal_tx, load_credit_exposure_reservation_tx,
62    reserve_credit_exposure_tx,
63};
64use errors::*;
65pub use factor_assignment::{
66    DurableFactorAssignmentResultV1, FactorAssignmentAuthorityRegistryV1,
67    FactorAssignmentAuthoritySetHeadV1, FactorAssignmentCommitV1,
68    FactorAssignmentSigningAuthorityV1, FactorAssignmentVerificationAuthorityV1,
69    SqliteFactorAssignmentStore, StoredFactorAssignmentResultV1,
70};
71use obligation::load_durable_obligation;
72pub(crate) use participant::{
73    advance_budget_authorization_tx, advance_budget_capture_tx, advance_tool_outcome_tx,
74    append_participant_update_tx, finalize_channel_reservation_operation_tx,
75    verify_budget_authorization_replay_tx, verify_participant_recovery_tx,
76    BudgetAuthorizationAdvance,
77};
78use participant::{
79    ensure_no_reserved_terminal_stage, qualify_generic_channel_command,
80    validate_payment_reconcile_binding, verify_payment_terminal_source,
81    verify_payment_write_context,
82};
83use projection::{
84    ensure_projection_absent, full_projection_capabilities, insert_terminal_projection,
85    insert_verified_terminal_projection, projected_terminal_state, terminal_from_operation,
86    validate_canonical_projection_size, verify_exact_signed_terminal_replay,
87    verify_exact_terminal_replay, verify_stored_terminal_projection,
88};
89use schema::{coordinator_lease_id_for_epoch, recovery_claim_digest, verify_latest_commit};
90pub(crate) use schema::{
91    initialize_admission_operation_schema, validate_trusted_time, verify_active_owner,
92    verify_admission_operation_invariants, verify_trusted_time,
93};
94
95const ADMISSION_OPERATION_SCHEMA_KEY: &str = "admission_operation";
96pub(crate) const ADMISSION_OPERATION_SUPPORTED_SCHEMA_VERSION: i32 = 9;
97const ADMISSION_OPERATION_SCHEMA_ANCHORS: &[&str] = &[
98    "admission_operations",
99    "admission_operation_commits",
100    "threshold_approval_proposals",
101    "chio_serving_owner",
102    "capability_grant_budgets",
103];
104const MAX_PERSISTED_OPERATION_BYTES: usize = 256 * 1024;
105const MAX_TERMINAL_PROJECTION_BYTES: usize = 4 * 1024 * 1024;
106const MAX_TERMINAL_MANIFEST_BYTES: usize = 256 * 1024;
107const MAX_TERMINAL_RECORD_BYTES: usize = 1024 * 1024;
108const MAX_TERMINAL_RECORDS: usize = 32;
109const MAX_RECOVERY_BATCH: usize = 256;
110const MAX_TRUSTED_UNIX_MS: u64 = (1_u64 << 53) - 1;
111const MAX_TRUSTED_CLOCK_SKEW_MS: u64 = 5 * 60 * 1_000;
112const MAX_RECOVERY_LEASE_DURATION_MS: u64 = 5 * 60 * 1_000;
113const COMBINED_CAPTURE_OPERATION_MUTATION_KIND: &str = "compare_and_swap";
114
115const ADMISSION_OPERATION_SCHEMA: &str = include_str!("admission_operation_store.sql");
116
117#[derive(Clone)]
118pub struct SqliteAdmissionOperationStore {
119    connection: Arc<Mutex<Connection>>,
120    serving_owner: Arc<SqliteServingOwner>,
121}
122
123#[derive(Debug, Clone, PartialEq, Eq)]
124pub struct DurableObligationV1 {
125    atom: ObligationAtomV1,
126    disposition: ObligationDispositionRecordV1,
127    settlement_lifecycle: ObligationSettlementLifecycleV1,
128    head_sequence: u64,
129    head_digest: String,
130    snapshot_version: u64,
131    resource_fence: u64,
132}
133
134impl DurableObligationV1 {
135    #[must_use]
136    pub const fn atom(&self) -> &ObligationAtomV1 {
137        &self.atom
138    }
139
140    #[must_use]
141    pub const fn disposition(&self) -> &ObligationDispositionRecordV1 {
142        &self.disposition
143    }
144
145    #[must_use]
146    pub const fn settlement_lifecycle(&self) -> &ObligationSettlementLifecycleV1 {
147        &self.settlement_lifecycle
148    }
149
150    #[must_use]
151    pub const fn head_sequence(&self) -> u64 {
152        self.head_sequence
153    }
154
155    #[must_use]
156    pub fn head_digest(&self) -> &str {
157        &self.head_digest
158    }
159
160    #[must_use]
161    pub const fn snapshot_version(&self) -> u64 {
162        self.snapshot_version
163    }
164
165    #[must_use]
166    pub const fn resource_fence(&self) -> u64 {
167        self.resource_fence
168    }
169}
170
171impl SqliteAdmissionOperationStore {
172    pub(crate) fn open_alongside(
173        connection: Arc<Mutex<Connection>>,
174        serving_owner: Arc<SqliteServingOwner>,
175    ) -> Self {
176        Self {
177            connection,
178            serving_owner,
179        }
180    }
181
182    fn connection(&self) -> Result<MutexGuard<'_, Connection>, AdmissionOperationStoreError> {
183        self.connection.lock().map_err(|_| {
184            AdmissionOperationStoreError::Invariant(
185                "sqlite admission operation lock poisoned".to_string(),
186            )
187        })
188    }
189
190    fn begin_read<'a>(
191        &self,
192        connection: &'a mut Connection,
193    ) -> Result<Transaction<'a>, AdmissionOperationStoreError> {
194        let transaction = connection
195            .transaction_with_behavior(TransactionBehavior::Deferred)
196            .map_err(sqlite_error)?;
197        verify_active_owner(&transaction, &self.serving_owner, None)?;
198        self.serving_owner
199            .verify_authority_anchor(&transaction)
200            .map_err(map_owner_error)?;
201        Ok(transaction)
202    }
203
204    fn begin_write<'a>(
205        &self,
206        connection: &'a mut Connection,
207        fence: Option<&StoreMutationFence>,
208    ) -> Result<Transaction<'a>, AdmissionOperationStoreError> {
209        let transaction = connection
210            .transaction_with_behavior(TransactionBehavior::Immediate)
211            .map_err(sqlite_error)?;
212        verify_active_owner(&transaction, &self.serving_owner, fence)?;
213        self.serving_owner
214            .verify_authority_anchor(&transaction)
215            .map_err(map_owner_error)?;
216        Ok(transaction)
217    }
218
219    fn sync_after_write(
220        &self,
221        connection: &Connection,
222    ) -> Result<(), AdmissionOperationStoreError> {
223        self.serving_owner
224            .sync_authority_anchor(connection)
225            .map_err(map_owner_error)
226    }
227
228    fn commit_write(
229        &self,
230        transaction: Transaction<'_>,
231    ) -> Result<(), AdmissionOperationStoreError> {
232        transaction.commit().map_err(|error| {
233            map_owner_error(self.serving_owner.outcome_unknown(format!(
234                "sqlite admission operation commit outcome is unknown: {error}"
235            )))
236        })
237    }
238
239    pub fn load_obligation(
240        &self,
241        obligation_id: &str,
242    ) -> Result<Option<DurableObligationV1>, AdmissionOperationStoreError> {
243        let mut connection = self.connection()?;
244        let transaction = self.begin_read(&mut connection)?;
245        load_durable_obligation(&transaction, obligation_id)
246    }
247
248    #[cfg(test)]
249    pub(crate) fn provision_credit_exposure_account(
250        &self,
251        request: &CreditExposureReservationRequest,
252        active_fence: &StoreMutationFence,
253        trusted_now_unix_ms: u64,
254    ) -> Result<CreditExposureAccountSnapshot, AdmissionOperationStoreError> {
255        if active_fence != &self.serving_owner.fence {
256            return Err(AdmissionOperationStoreError::Fenced);
257        }
258        request
259            .validate()
260            .map_err(|error| invariant(error.to_string()))?;
261        request
262            .authorities
263            .ensure_current_at(trusted_now_unix_ms / 1_000)
264            .map_err(|error| invariant(error.to_string()))?;
265        request
266            .credit_facility_bind
267            .ensure_current_at(trusted_now_unix_ms)
268            .map_err(|error| invariant(error.to_string()))?;
269        let bind = request.credit_facility_bind.body();
270        let reservation = CreditExposureReservationRecordV1::prepare_reserved(
271            request,
272            bind.expected_exposure_version()
273                .checked_add(1)
274                .ok_or_else(|| invariant("credit exposure account version overflowed"))?,
275            bind.expected_exposure_fence()
276                .checked_add(1)
277                .ok_or_else(|| invariant("credit exposure resource fence overflowed"))?,
278        )
279        .map_err(|error| invariant(error.to_string()))?;
280        let mut connection = self.connection()?;
281        let transaction = self.begin_write(&mut connection, Some(active_fence))?;
282        let snapshot = credit_exposure::initialize_credit_exposure_account_tx(
283            &transaction,
284            &reservation,
285            0,
286            0,
287            active_fence,
288            trusted_now_unix_ms,
289        )?;
290        self.commit_write(transaction)?;
291        self.sync_after_write(&connection)?;
292        Ok(snapshot)
293    }
294
295    pub fn load_credit_exposure_account(
296        &self,
297        debtor_id: &str,
298        scope_digest: &str,
299        currency: &str,
300    ) -> Result<Option<CreditExposureAccountSnapshot>, AdmissionOperationStoreError> {
301        AdmissionIdentifier::try_new("credit_exposure_debtor_id", debtor_id.to_owned())?;
302        AdmissionDigest::try_new("credit_exposure_scope_digest", scope_digest.to_owned())?;
303        if currency.len() != 3 || !currency.bytes().all(|byte| byte.is_ascii_uppercase()) {
304            return Err(invariant("credit exposure currency is invalid"));
305        }
306        let mut connection = self.connection()?;
307        let transaction = self.begin_read(&mut connection)?;
308        credit_exposure::load_credit_exposure_account_tx(
309            &transaction,
310            debtor_id,
311            scope_digest,
312            currency,
313        )
314    }
315
316    pub fn load_credit_exposure_reservation(
317        &self,
318        operation_id: &str,
319    ) -> Result<Option<CreditExposureReservationRecordV1>, AdmissionOperationStoreError> {
320        AdmissionDigest::try_new("credit_exposure_operation_id", operation_id.to_owned())?;
321        let mut connection = self.connection()?;
322        let transaction = self.begin_read(&mut connection)?;
323        load_credit_exposure_reservation_tx(&transaction, operation_id)
324    }
325
326    pub fn capture_invocation_and_commit_dispatch(
327        &self,
328        operation: &AdmissionOperationV1,
329        recovery_lease: &AdmissionRecoveryLease,
330        request: BudgetCaptureInvocationRequest,
331        active_fence: &StoreMutationFence,
332        trusted_now_unix_ms: u64,
333    ) -> Result<
334        (
335            chio_kernel::budget_store::BudgetInvocationCaptureDecision,
336            AdmissionOperationV1,
337        ),
338        AdmissionCaptureError,
339    > {
340        if active_fence != &self.serving_owner.fence
341            || recovery_lease.store_fence() != active_fence
342            || operation.state() != AdmissionOperationState::CapturePending
343            || operation.binding().capability_id().as_str() != request.capability_id
344            || operation
345                .budget_hold_id()
346                .is_none_or(|hold_id| hold_id.as_str() != request.hold_id)
347        {
348            return Err(AdmissionCaptureError::Fenced);
349        }
350        let budget = crate::budget_store::SqliteBudgetStore::open_alongside(
351            self.connection.clone(),
352            self.serving_owner.clone(),
353        );
354        budget
355            .capture_composite_invocation_and_commit_dispatch(
356                request,
357                crate::budget_store::AdmissionCaptureBinding {
358                    operation,
359                    recovery_lease,
360                    trusted_now_unix_ms,
361                },
362            )
363            .map_err(map_budget_capture_error)
364    }
365
366    #[allow(clippy::too_many_arguments)]
367    pub fn authorize_budget_and_commit_admission(
368        &self,
369        operation: &AdmissionOperationV1,
370        recovery_lease: &AdmissionRecoveryLease,
371        request: BudgetAuthorizeHoldRequest,
372        payment_journal: Option<PaymentJournalRecord>,
373        credit_exposure: Option<chio_credit::obligation::CreditExposureReservationRequest>,
374        active_fence: &StoreMutationFence,
375        trusted_now_unix_ms: u64,
376    ) -> Result<(BudgetAuthorizeHoldDecision, AdmissionOperationV1), AdmissionCaptureError> {
377        if active_fence != &self.serving_owner.fence || recovery_lease.store_fence() != active_fence
378        {
379            return Err(AdmissionCaptureError::Fenced);
380        }
381        let budget = crate::budget_store::SqliteBudgetStore::open_alongside(
382            self.connection.clone(),
383            self.serving_owner.clone(),
384        );
385        budget
386            .authorize_composite_hold_and_commit_admission(
387                request,
388                crate::budget_store::AdmissionAuthorizationBinding {
389                    operation,
390                    recovery_lease,
391                    payment_journal: payment_journal.as_ref(),
392                    credit_exposure: credit_exposure.as_ref(),
393                    trusted_now_unix_ms,
394                },
395            )
396            .map_err(map_budget_capture_error)
397    }
398
399    pub fn load_payment_journal(
400        &self,
401        operation_id: &str,
402        active_fence: &StoreMutationFence,
403    ) -> Result<Option<PaymentJournalRecord>, AdmissionPaymentJournalError> {
404        if active_fence != &self.serving_owner.fence {
405            return Err(AdmissionPaymentJournalError::Fenced);
406        }
407        let mut connection = self.connection().map_err(map_payment_operation_error)?;
408        let transaction = self
409            .begin_read(&mut connection)
410            .map_err(map_payment_operation_error)?;
411        let journal = crate::budget_store::load_payment_journal(&transaction, operation_id)
412            .map_err(map_payment_budget_error)?;
413        transaction
414            .commit()
415            .map_err(|error| AdmissionPaymentJournalError::Invariant(error.to_string()))?;
416        Ok(journal)
417    }
418
419    pub fn advance_payment_journal(
420        &self,
421        advance: AdmissionPaymentJournalAdvance<'_>,
422    ) -> Result<PaymentJournalRecord, AdmissionPaymentJournalError> {
423        let AdmissionPaymentJournalAdvance {
424            operation,
425            recovery_lease,
426            expected,
427            transition,
428            release_evidence,
429            active_fence,
430            trusted_now_unix_ms,
431        } = advance;
432        if active_fence != &self.serving_owner.fence
433            || recovery_lease.store_fence() != active_fence
434            || expected.operation_id != operation.binding().operation_id().as_str()
435        {
436            return Err(AdmissionPaymentJournalError::Fenced);
437        }
438        let mut connection = self.connection().map_err(map_payment_operation_error)?;
439        let transaction = self
440            .begin_write(&mut connection, Some(active_fence))
441            .map_err(map_payment_operation_error)?;
442        verify_payment_write_context(
443            &transaction,
444            &self.serving_owner,
445            operation,
446            recovery_lease,
447            active_fence,
448            trusted_now_unix_ms,
449        )?;
450        let (updated, changed) = crate::budget_store::advance_payment_journal(
451            &transaction,
452            expected,
453            transition,
454            release_evidence,
455            trusted_now_unix_ms,
456        )
457        .map_err(map_payment_budget_error)?;
458        if changed {
459            self.serving_owner
460                .append_global_commit(
461                    &transaction,
462                    "payment_journal_transition",
463                    "payment",
464                    &updated.operation_id,
465                    updated.journal_version,
466                )
467                .map_err(map_payment_owner_error)?;
468        }
469        self.commit_write(transaction)
470            .map_err(map_payment_operation_error)?;
471        if changed {
472            self.sync_after_write(&connection)
473                .map_err(map_payment_operation_error)?;
474        }
475        Ok(updated)
476    }
477
478    pub fn begin_payment_settlement(
479        &self,
480        begin: AdmissionPaymentSettlementBegin<'_>,
481    ) -> Result<AdmissionPaymentSettlement, AdmissionPaymentJournalError> {
482        let AdmissionPaymentSettlementBegin {
483            operation,
484            recovery_lease,
485            expected,
486            transition,
487            release_evidence,
488            budget_reconcile,
489            active_fence,
490            trusted_now_unix_ms,
491        } = begin;
492        if active_fence != &self.serving_owner.fence
493            || recovery_lease.store_fence() != active_fence
494            || expected.operation_id != operation.binding().operation_id().as_str()
495        {
496            return Err(AdmissionPaymentJournalError::Fenced);
497        }
498        validate_payment_reconcile_binding(expected, transition, &budget_reconcile)?;
499        let mut connection = self.connection().map_err(map_payment_operation_error)?;
500        let transaction = self
501            .begin_write(&mut connection, Some(active_fence))
502            .map_err(map_payment_operation_error)?;
503        verify_payment_write_context(
504            &transaction,
505            &self.serving_owner,
506            operation,
507            recovery_lease,
508            active_fence,
509            trusted_now_unix_ms,
510        )?;
511        let budget = crate::budget_store::SqliteBudgetStore::open_alongside(
512            self.connection.clone(),
513            self.serving_owner.clone(),
514        );
515        let (budget, budget_changed) = budget
516            .reconcile_composite_hold_in_transaction(&transaction, &budget_reconcile)
517            .map_err(map_payment_budget_error)?;
518        let (journal, payment_changed) = match transition {
519            Some(transition) => crate::budget_store::advance_payment_journal(
520                &transaction,
521                expected,
522                transition,
523                release_evidence,
524                trusted_now_unix_ms,
525            )
526            .map_err(map_payment_budget_error)?,
527            None => {
528                if release_evidence.is_some() {
529                    return Err(AdmissionPaymentJournalError::Invariant(
530                        "payment release evidence requires a journal transition".to_owned(),
531                    ));
532                }
533                let stored = crate::budget_store::load_payment_journal(
534                    &transaction,
535                    operation.binding().operation_id().as_str(),
536                )
537                .map_err(map_payment_budget_error)?
538                .ok_or_else(|| {
539                    AdmissionPaymentJournalError::Invariant(
540                        "payment settlement journal is absent".to_owned(),
541                    )
542                })?;
543                if stored != *expected {
544                    return Err(AdmissionPaymentJournalError::Conflict(
545                        "payment settlement journal changed".to_owned(),
546                    ));
547                }
548                (stored, false)
549            }
550        };
551        if payment_changed {
552            self.serving_owner
553                .append_global_commit(
554                    &transaction,
555                    "payment_settlement_intent",
556                    "payment",
557                    &journal.operation_id,
558                    journal.journal_version,
559                )
560                .map_err(map_payment_owner_error)?;
561        }
562        self.commit_write(transaction)
563            .map_err(map_payment_operation_error)?;
564        if budget_changed || payment_changed {
565            self.sync_after_write(&connection)
566                .map_err(map_payment_operation_error)?;
567        }
568        Ok(AdmissionPaymentSettlement {
569            journal,
570            budget,
571            budget_already_reconciled: !budget_changed,
572        })
573    }
574
575    pub fn commit_terminal_projection(
576        &self,
577        projection: &AdmissionTerminalProjection,
578    ) -> Result<AdmissionTerminal, AdmissionOperationStoreError> {
579        if projection.requires_anchored_economic_commit() {
580            return Err(invariant(
581                "terminal projection requires an advanced economic anchor",
582            ));
583        }
584        projection.context().validate()?;
585        let mut connection = self.connection()?;
586        let transaction =
587            self.begin_write(&mut connection, Some(&projection.context().store_fence))?;
588        let (terminal, changed) =
589            self.commit_terminal_projection_in_transaction(&transaction, projection)?;
590        self.commit_write(transaction)?;
591        if changed {
592            self.sync_after_write(&connection)?;
593        }
594        Ok(terminal)
595    }
596
597    pub(super) fn commit_terminal_projection_in_transaction(
598        &self,
599        transaction: &Transaction<'_>,
600        projection: &AdmissionTerminalProjection,
601    ) -> Result<(AdmissionTerminal, bool), AdmissionOperationStoreError> {
602        let stored = load_by_operation_id_tx(transaction, &projection.context().operation_id)?
603            .ok_or(AdmissionOperationStoreError::NotFound)?;
604        self.commit_terminal_projection_from_source_in_transaction(transaction, projection, &stored)
605    }
606
607    fn commit_terminal_projection_from_source_in_transaction(
608        &self,
609        transaction: &Transaction<'_>,
610        projection: &AdmissionTerminalProjection,
611        stored: &StoredOperation,
612    ) -> Result<(AdmissionTerminal, bool), AdmissionOperationStoreError> {
613        if projection.requires_anchored_economic_commit() {
614            return Err(invariant(
615                "terminal projection requires an advanced economic anchor",
616            ));
617        }
618        let canonical = projection.canonical_projection()?;
619        validate_canonical_projection_size(&canonical)?;
620        let context = projection.context();
621        context.validate()?;
622        verify_trusted_time(transaction, context.trusted_time_unix_ms)?;
623        verify_payment_terminal_source(
624            transaction,
625            &stored.operation,
626            context,
627            projected_terminal_state(projection),
628            canonical.records().iter().filter_map(|record| {
629                (record.commitment().kind() == AdmissionProjectionRecordKind::PaymentTerminal)
630                    .then_some(record.canonical_bytes())
631            }),
632        )?;
633
634        if stored.operation.state().is_terminal() {
635            let terminal = verify_exact_terminal_replay(
636                transaction,
637                &stored.operation,
638                projection,
639                &canonical,
640            )?;
641            apply_credit_exposure_terminal_tx(
642                transaction,
643                &stored.operation,
644                canonical.projection_digest(),
645                projection.pre_dispatch_release_proof(),
646                &context.store_fence,
647                context.trusted_time_unix_ms,
648            )?;
649            return Ok((terminal, false));
650        }
651        if context.request_id != stored.operation.replay_key().request_id
652            || context.expected_operation_version != stored.operation.version()
653            || context.coordinator_lease_epoch != stored.operation.coordinator_lease_epoch()
654            || context.trusted_time_unix_ms < stored.updated_at_unix_ms
655        {
656            return Err(AdmissionOperationError::TerminalProjectionBindingMismatch.into());
657        }
658        let recovery_claim = stored
659            .recovery_claim
660            .as_ref()
661            .ok_or(AdmissionOperationStoreError::Fenced)?;
662        if recovery_claim.coordinator_lease_id() != &context.coordinator_lease_id
663            || recovery_claim.coordinator_lease_epoch() != context.coordinator_lease_epoch
664            || recovery_claim.store_fence() != &context.store_fence
665        {
666            return Err(AdmissionOperationStoreError::Fenced);
667        }
668        verify_stored_recovery_claim(
669            transaction,
670            &self.serving_owner,
671            stored,
672            recovery_claim,
673            context.trusted_time_unix_ms,
674            &context.store_fence,
675        )?;
676        let capabilities = full_projection_capabilities();
677        let updated = stored
678            .operation
679            .apply_terminal_projection(projection, &capabilities)?;
680        if updated
681            .terminal_replay()
682            .is_none_or(|replay| replay.projection_digest() != canonical.projection_digest())
683        {
684            return Err(invariant(
685                "terminal operation does not retain its exact projection digest",
686            ));
687        }
688        let encoded = encode_operation(&updated)?;
689        let changed = transaction
690            .execute(
691                r#"
692                UPDATE admission_operations
693                SET operation_json = ?1, state = ?2, terminal = 1,
694                    coordinator_lease_epoch = ?3, version = ?4,
695                    updated_at_unix_ms = ?5
696                WHERE operation_id = ?6 AND version = ?7 AND terminal = 0
697                "#,
698                params![
699                    &encoded,
700                    state_name(updated.state()),
701                    sqlite_i64(updated.coordinator_lease_epoch(), "coordinator_lease_epoch")?,
702                    sqlite_i64(updated.version(), "terminal_operation_version")?,
703                    sqlite_i64(context.trusted_time_unix_ms, "trusted_now_unix_ms")?,
704                    context.operation_id.as_str(),
705                    sqlite_i64(stored.operation.version(), "expected_operation_version")?,
706                ],
707            )
708            .map_err(sqlite_error)?;
709        if changed != 1 {
710            return Err(AdmissionOperationStoreError::Fenced);
711        }
712        insert_terminal_projection(transaction, projection, &canonical, &updated)?;
713        apply_credit_exposure_terminal_tx(
714            transaction,
715            &updated,
716            canonical.projection_digest(),
717            projection.pre_dispatch_release_proof(),
718            &context.store_fence,
719            context.trusted_time_unix_ms,
720        )?;
721        append_operation_commit(
722            transaction,
723            &updated,
724            &encoded,
725            Some(recovery_claim),
726            "compare_and_swap",
727            &self.serving_owner,
728            context.trusted_time_unix_ms,
729        )?;
730        terminal_from_operation(&updated).map(|terminal| (terminal, true))
731    }
732
733    pub fn commit_signed_terminal_projection(
734        &self,
735        envelope: &SignedAdmissionTerminalProjectionV1,
736    ) -> Result<AdmissionTerminal, AdmissionOperationStoreError> {
737        let verified = envelope.verify()?;
738        let context = verified.context();
739        let mut connection = self.connection()?;
740        let transaction = self.begin_write(&mut connection, Some(&context.store_fence))?;
741        let terminal = self.commit_verified_signed_terminal_projection_in_transaction(
742            &transaction,
743            &verified,
744            context.trusted_time_unix_ms,
745            None,
746        )?;
747        self.commit_write(transaction)?;
748        self.sync_after_write(&connection)?;
749        Ok(terminal)
750    }
751}
752
753fn verify_stored_recovery_claim(
754    transaction: &Transaction<'_>,
755    owner: &SqliteServingOwner,
756    stored: &StoredOperation,
757    claim: &UntrustedAdmissionRecoveryClaim,
758    trusted_now_unix_ms: u64,
759    current_store_fence: &StoreMutationFence,
760) -> Result<(), AdmissionOperationStoreError> {
761    if trusted_now_unix_ms < stored.updated_at_unix_ms {
762        return Err(invariant("trusted operation time regressed"));
763    }
764    if stored.recovery_claim.as_ref() != Some(claim)
765        || stored.operation.binding().operation_id() != claim.operation_id()
766        || stored.operation.version() != claim.claimed_version()
767        || stored.operation.coordinator_lease_epoch() != claim.coordinator_lease_epoch()
768        || claim.store_fence() != current_store_fence
769        || current_store_fence != &owner.fence
770    {
771        return Err(AdmissionOperationStoreError::Fenced);
772    }
773    let historical_lease_id = coordinator_lease_id_for_epoch(
774        transaction,
775        owner,
776        stored.operation.coordinator_lease_epoch(),
777    )?;
778    if &historical_lease_id != claim.coordinator_lease_id() {
779        return Err(AdmissionOperationStoreError::Fenced);
780    }
781    if trusted_now_unix_ms >= claim.expires_at_unix_ms() {
782        return Err(AdmissionOperationError::LeaseExpired.into());
783    }
784    Ok(())
785}
786
787struct StoredOperation {
788    operation: AdmissionOperationV1,
789    recovery_claim: Option<UntrustedAdmissionRecoveryClaim>,
790    updated_at_unix_ms: u64,
791}
792
793pub(crate) enum PreparedAdmissionBeginTxResult {
794    Created {
795        encoded: Vec<u8>,
796    },
797    ExactReplay {
798        operation: Box<AdmissionOperationV1>,
799        terminal_replay: Option<AdmissionTerminalReplay>,
800    },
801    Conflict {
802        existing_operation_id: AdmissionOperationId,
803    },
804}
805
806struct RawOperationRow {
807    operation_id: String,
808    request_namespace_digest: String,
809    request_id: String,
810    operation_json: Vec<u8>,
811    state: String,
812    terminal: i64,
813    coordinator_lease_epoch: i64,
814    version: i64,
815    created_at_unix_ms: i64,
816    updated_at_unix_ms: i64,
817    recovery_claimant_id: Option<String>,
818    recovery_coordinator_lease_id: Option<String>,
819    recovery_coordinator_lease_epoch: Option<i64>,
820    recovery_claimed_version: Option<i64>,
821    recovery_expires_at_unix_ms: Option<i64>,
822    recovery_store_uuid: Option<String>,
823    recovery_store_lease_id: Option<String>,
824    recovery_store_owner_epoch: Option<i64>,
825}
826
827fn read_raw_row(row: &Row<'_>) -> rusqlite::Result<RawOperationRow> {
828    Ok(RawOperationRow {
829        operation_id: row.get(0)?,
830        request_namespace_digest: row.get(1)?,
831        request_id: row.get(2)?,
832        operation_json: row.get(3)?,
833        state: row.get(4)?,
834        terminal: row.get(5)?,
835        coordinator_lease_epoch: row.get(6)?,
836        version: row.get(7)?,
837        created_at_unix_ms: row.get(8)?,
838        updated_at_unix_ms: row.get(9)?,
839        recovery_claimant_id: row.get(10)?,
840        recovery_coordinator_lease_id: row.get(11)?,
841        recovery_coordinator_lease_epoch: row.get(12)?,
842        recovery_claimed_version: row.get(13)?,
843        recovery_expires_at_unix_ms: row.get(14)?,
844        recovery_store_uuid: row.get(15)?,
845        recovery_store_lease_id: row.get(16)?,
846        recovery_store_owner_epoch: row.get(17)?,
847    })
848}
849
850fn decode_row(raw: RawOperationRow) -> Result<StoredOperation, AdmissionOperationStoreError> {
851    if raw.operation_json.is_empty() || raw.operation_json.len() > MAX_PERSISTED_OPERATION_BYTES {
852        return Err(invariant("persisted admission operation size is invalid"));
853    }
854    let persisted: PersistedAdmissionOperationV1 = serde_json::from_slice(&raw.operation_json)
855        .map_err(|error| invariant(format!("persisted admission operation is invalid: {error}")))?;
856    let operation = AdmissionOperationV1::from_persisted(persisted)?;
857    let canonical = encode_operation(&operation)?;
858    if canonical != raw.operation_json {
859        return Err(invariant(
860            "persisted admission operation encoding is not canonical",
861        ));
862    }
863    let replay_key = operation.replay_key();
864    if operation.binding().operation_id().as_str() != raw.operation_id
865        || replay_key.request_namespace_digest.as_str() != raw.request_namespace_digest
866        || replay_key.request_id.as_str() != raw.request_id
867        || state_name(operation.state()) != raw.state
868        || i64::from(operation.state().is_terminal()) != raw.terminal
869        || operation.coordinator_lease_epoch()
870            != stored_u64(raw.coordinator_lease_epoch, "coordinator_lease_epoch")?
871        || operation.version() != stored_u64(raw.version, "version")?
872    {
873        return Err(invariant(
874            "admission operation columns do not match the checked record",
875        ));
876    }
877    let created_at = stored_u64(raw.created_at_unix_ms, "created_at_unix_ms")?;
878    let updated_at = stored_u64(raw.updated_at_unix_ms, "updated_at_unix_ms")?;
879    validate_trusted_time(created_at, "created_at_unix_ms")?;
880    validate_trusted_time(updated_at, "updated_at_unix_ms")?;
881    if updated_at < created_at {
882        return Err(invariant("admission operation timestamp regressed"));
883    }
884
885    let recovery_claim = match (
886        raw.recovery_claimant_id,
887        raw.recovery_coordinator_lease_id,
888        raw.recovery_coordinator_lease_epoch,
889        raw.recovery_claimed_version,
890        raw.recovery_expires_at_unix_ms,
891        raw.recovery_store_uuid,
892        raw.recovery_store_lease_id,
893        raw.recovery_store_owner_epoch,
894    ) {
895        (None, None, None, None, None, None, None, None) => None,
896        (
897            Some(claimant_id),
898            Some(coordinator_lease_id),
899            Some(coordinator_lease_epoch),
900            Some(claimed_version),
901            Some(expires_at_unix_ms),
902            Some(store_uuid),
903            Some(store_lease_id),
904            Some(store_owner_epoch),
905        ) => {
906            let claimed_version = stored_u64(claimed_version, "recovery_claimed_version")?;
907            if claimed_version > operation.version()
908                || claimed_version
909                    .checked_add(1)
910                    .is_some_and(|next| next < operation.version())
911            {
912                return Err(invariant(
913                    "recovery claim is not for the current or immediately preceding version",
914                ));
915            }
916            let expires_at_unix_ms = stored_u64(expires_at_unix_ms, "recovery_expires_at_unix_ms")?;
917            validate_trusted_time(expires_at_unix_ms, "recovery_expires_at_unix_ms")?;
918            Some(UntrustedAdmissionRecoveryClaim::new(
919                operation.binding().operation_id().clone(),
920                AdmissionIdentifier::try_new("recovery_claimant_id", claimant_id)?,
921                AdmissionIdentifier::try_new(
922                    "recovery_coordinator_lease_id",
923                    coordinator_lease_id,
924                )?,
925                stored_u64(coordinator_lease_epoch, "recovery_coordinator_lease_epoch")?,
926                claimed_version,
927                expires_at_unix_ms,
928                StoreMutationFence {
929                    store_uuid,
930                    lease_id: store_lease_id,
931                    owner_epoch: stored_u64(store_owner_epoch, "recovery_store_owner_epoch")?,
932                },
933            )?)
934        }
935        _ => return Err(invariant("recovery claim tuple is partial")),
936    };
937    Ok(StoredOperation {
938        operation,
939        recovery_claim,
940        updated_at_unix_ms: updated_at,
941    })
942}
943
944fn load_by_operation_id_tx(
945    transaction: &Transaction<'_>,
946    operation_id: &AdmissionOperationId,
947) -> Result<Option<StoredOperation>, AdmissionOperationStoreError> {
948    let raw = transaction
949        .query_row(
950            r#"
951            SELECT operation_id, request_namespace_digest, request_id,
952                   operation_json, state, terminal, coordinator_lease_epoch,
953                   version, created_at_unix_ms, updated_at_unix_ms,
954                   recovery_claimant_id, recovery_coordinator_lease_id,
955                   recovery_coordinator_lease_epoch, recovery_claimed_version,
956                   recovery_expires_at_unix_ms, recovery_store_uuid,
957                   recovery_store_lease_id, recovery_store_owner_epoch
958            FROM admission_operations WHERE operation_id = ?1
959            "#,
960            [operation_id.as_str()],
961            read_raw_row,
962        )
963        .optional()
964        .map_err(sqlite_error)?;
965    let stored = raw.map(decode_row).transpose()?;
966    if let Some(stored) = &stored {
967        verify_latest_commit(transaction, stored)?;
968        verify_stored_terminal_projection(transaction, stored)?;
969    }
970    Ok(stored)
971}
972
973pub(crate) fn begin_prepared_operation_tx(
974    transaction: &Transaction<'_>,
975    operation: &AdmissionOperationV1,
976    fence: &StoreMutationFence,
977    trusted_now_unix_ms: u64,
978) -> Result<PreparedAdmissionBeginTxResult, AdmissionOperationStoreError> {
979    operation.validate()?;
980    if operation.state() != AdmissionOperationState::Prepared || operation.version() != 1 {
981        return Err(invariant("begin requires a version-one Prepared operation"));
982    }
983    if operation.coordinator_lease_epoch() != fence.owner_epoch {
984        return Err(AdmissionOperationStoreError::Fenced);
985    }
986    verify_trusted_time(transaction, trusted_now_unix_ms)?;
987    let encoded = encode_operation(operation)?;
988    let replay_key = operation.replay_key();
989    if let Some(existing) = load_by_replay_key_tx(transaction, &replay_key)? {
990        return Ok(match existing.operation.classify_replay(operation) {
991            AdmissionReplayClassification::Exact { terminal_replay } => {
992                PreparedAdmissionBeginTxResult::ExactReplay {
993                    operation: Box::new(existing.operation),
994                    terminal_replay,
995                }
996            }
997            AdmissionReplayClassification::Conflict => PreparedAdmissionBeginTxResult::Conflict {
998                existing_operation_id: existing.operation.binding().operation_id().clone(),
999            },
1000        });
1001    }
1002    if load_by_operation_id_tx(transaction, operation.binding().operation_id())?.is_some() {
1003        return Err(invariant(
1004            "operation id is already bound to a different replay key",
1005        ));
1006    }
1007    let changed = transaction
1008        .execute(
1009            r#"
1010            INSERT INTO admission_operations (
1011                operation_id, request_namespace_digest, request_id,
1012                operation_json, state, terminal, coordinator_lease_epoch,
1013                version, created_at_unix_ms, updated_at_unix_ms
1014            ) VALUES (?1, ?2, ?3, ?4, ?5, 0, ?6, ?7, ?8, ?8)
1015            "#,
1016            params![
1017                operation.binding().operation_id().as_str(),
1018                replay_key.request_namespace_digest.as_str(),
1019                replay_key.request_id.as_str(),
1020                encoded,
1021                state_name(operation.state()),
1022                sqlite_i64(
1023                    operation.coordinator_lease_epoch(),
1024                    "coordinator_lease_epoch"
1025                )?,
1026                sqlite_i64(operation.version(), "version")?,
1027                sqlite_i64(trusted_now_unix_ms, "created_at_unix_ms")?,
1028            ],
1029        )
1030        .map_err(sqlite_error)?;
1031    if changed != 1 {
1032        return Err(invariant("begin did not insert exactly one operation"));
1033    }
1034    Ok(PreparedAdmissionBeginTxResult::Created { encoded })
1035}
1036
1037pub(crate) fn load_operation_for_participant_tx(
1038    transaction: &Transaction<'_>,
1039    operation_id: &AdmissionOperationId,
1040) -> Result<Option<AdmissionOperationV1>, AdmissionOperationStoreError> {
1041    load_by_operation_id_tx(transaction, operation_id)
1042        .map(|stored| stored.map(|stored| stored.operation))
1043}
1044
1045fn load_by_replay_key_tx(
1046    transaction: &Transaction<'_>,
1047    replay_key: &AdmissionReplayKey,
1048) -> Result<Option<StoredOperation>, AdmissionOperationStoreError> {
1049    let raw = transaction
1050        .query_row(
1051            r#"
1052            SELECT operation_id, request_namespace_digest, request_id,
1053                   operation_json, state, terminal, coordinator_lease_epoch,
1054                   version, created_at_unix_ms, updated_at_unix_ms,
1055                   recovery_claimant_id, recovery_coordinator_lease_id,
1056                   recovery_coordinator_lease_epoch, recovery_claimed_version,
1057                   recovery_expires_at_unix_ms, recovery_store_uuid,
1058                   recovery_store_lease_id, recovery_store_owner_epoch
1059            FROM admission_operations
1060            WHERE request_namespace_digest = ?1 AND request_id = ?2
1061            "#,
1062            params![
1063                replay_key.request_namespace_digest.as_str(),
1064                replay_key.request_id.as_str(),
1065            ],
1066            read_raw_row,
1067        )
1068        .optional()
1069        .map_err(sqlite_error)?;
1070    let stored = raw.map(decode_row).transpose()?;
1071    if let Some(stored) = &stored {
1072        verify_latest_commit(transaction, stored)?;
1073        verify_stored_terminal_projection(transaction, stored)?;
1074    }
1075    Ok(stored)
1076}
1077
1078fn encode_operation(
1079    operation: &AdmissionOperationV1,
1080) -> Result<Vec<u8>, AdmissionOperationStoreError> {
1081    operation.validate()?;
1082    let encoded = canonical_json_bytes(&operation.to_persisted())
1083        .map_err(|error| invariant(format!("admission operation encoding failed: {error}")))?;
1084    if encoded.is_empty() || encoded.len() > MAX_PERSISTED_OPERATION_BYTES {
1085        return Err(invariant(
1086            "persisted admission operation exceeds its size limit",
1087        ));
1088    }
1089    Ok(encoded)
1090}
1091
1092fn state_name(state: AdmissionOperationState) -> &'static str {
1093    match state {
1094        AdmissionOperationState::Prepared => "prepared",
1095        AdmissionOperationState::BrokerAttemptRegistered => "broker_attempt_registered",
1096        AdmissionOperationState::ApprovalRequired => "approval_required",
1097        AdmissionOperationState::BudgetAuthorized => "budget_authorized",
1098        AdmissionOperationState::ApprovalReserved => "approval_reserved",
1099        AdmissionOperationState::ReadyToDispatch => "ready_to_dispatch",
1100        AdmissionOperationState::CapturePending => "capture_pending",
1101        AdmissionOperationState::DispatchCommitted => "dispatch_committed",
1102        AdmissionOperationState::Finalizing => "finalizing",
1103        AdmissionOperationState::Completed => "completed",
1104        AdmissionOperationState::CompensatedBeforeDispatch => "compensated_before_dispatch",
1105        AdmissionOperationState::NotAcceptedAfterDispatchCommit => {
1106            "not_accepted_after_dispatch_commit"
1107        }
1108        AdmissionOperationState::OutcomeUnknownAfterDispatch => "outcome_unknown_after_dispatch",
1109        AdmissionOperationState::MutationReady => "mutation_ready",
1110        AdmissionOperationState::MutationSubmitted => "mutation_submitted",
1111        AdmissionOperationState::EconomicMutationApplied => "economic_mutation_applied",
1112        AdmissionOperationState::EconomicMutationNotApplied => "economic_mutation_not_applied",
1113    }
1114}
1115
1116fn sqlite_i64(value: u64, field: &'static str) -> Result<i64, AdmissionOperationStoreError> {
1117    i64::try_from(value).map_err(|_| invariant(format!("{field} exceeds SQLite integer range")))
1118}
1119
1120fn stored_u64(value: i64, field: &'static str) -> Result<u64, AdmissionOperationStoreError> {
1121    u64::try_from(value).map_err(|_| invariant(format!("{field} is negative")))
1122}
1123
1124fn invariant(detail: impl Into<String>) -> AdmissionOperationStoreError {
1125    AdmissionOperationStoreError::Invariant(detail.into())
1126}
1127
1128pub(crate) fn receipt_projection_error(error: AdmissionOperationStoreError) -> ReceiptStoreError {
1129    match error {
1130        AdmissionOperationStoreError::Unavailable(detail) => ReceiptStoreError::Pool(detail),
1131        AdmissionOperationStoreError::Fenced => ReceiptStoreError::Fenced,
1132        AdmissionOperationStoreError::NotFound => {
1133            ReceiptStoreError::NotFound("admission operation".to_string())
1134        }
1135        AdmissionOperationStoreError::Invariant(detail) => ReceiptStoreError::Conflict(detail),
1136        AdmissionOperationStoreError::OutcomeUnknown(detail) => {
1137            ReceiptStoreError::OutcomeUnknown(detail)
1138        }
1139        AdmissionOperationStoreError::Operation(error) => {
1140            ReceiptStoreError::Conflict(error.to_string())
1141        }
1142    }
1143}
1144
1145fn decode_projection_receipt(bytes: Vec<u8>) -> Result<ChioReceipt, ReceiptStoreError> {
1146    let receipt: ChioReceipt = serde_json::from_slice(&bytes)?;
1147    if canonical_json_bytes(&receipt)
1148        .map_err(|error| ReceiptStoreError::Canonical(error.to_string()))?
1149        != bytes
1150        || !receipt
1151            .verify_signature()
1152            .map_err(|error| ReceiptStoreError::CryptoDecode(error.to_string()))?
1153    {
1154        return Err(ReceiptStoreError::Conflict(
1155            "persisted admission receipt is invalid".to_string(),
1156        ));
1157    }
1158    Ok(receipt)
1159}
1160
1161fn map_owner_error(error: SqliteServingOwnerError) -> AdmissionOperationStoreError {
1162    match error {
1163        SqliteServingOwnerError::OutcomeUnknown(detail) => {
1164            AdmissionOperationStoreError::OutcomeUnknown(detail)
1165        }
1166        error => invariant(error.to_string()),
1167    }
1168}
1169
1170fn map_economic_cache_error(
1171    error: crate::economic_state_cache::EconomicStateCacheError,
1172) -> AdmissionOperationStoreError {
1173    match error {
1174        crate::economic_state_cache::EconomicStateCacheError::Unavailable(detail) => {
1175            AdmissionOperationStoreError::Unavailable(detail)
1176        }
1177        crate::economic_state_cache::EconomicStateCacheError::Fenced => {
1178            AdmissionOperationStoreError::Fenced
1179        }
1180        crate::economic_state_cache::EconomicStateCacheError::NotFound => {
1181            AdmissionOperationStoreError::NotFound
1182        }
1183        crate::economic_state_cache::EconomicStateCacheError::OutcomeUnknown(detail) => {
1184            AdmissionOperationStoreError::OutcomeUnknown(detail)
1185        }
1186        error => invariant(error.to_string()),
1187    }
1188}
1189
1190fn sqlite_error(error: rusqlite::Error) -> AdmissionOperationStoreError {
1191    match error {
1192        rusqlite::Error::FromSqlConversionFailure(..)
1193        | rusqlite::Error::IntegralValueOutOfRange(..)
1194        | rusqlite::Error::InvalidColumnType(..)
1195        | rusqlite::Error::Utf8Error(..) => invariant(error.to_string()),
1196        other => AdmissionOperationStoreError::Unavailable(other.to_string()),
1197    }
1198}
1199
1200impl From<AdmissionOperationStoreError> for SqliteServingOwnerError {
1201    fn from(error: AdmissionOperationStoreError) -> Self {
1202        Self::Invalid(error.to_string())
1203    }
1204}
1205
1206#[cfg(test)]
1207#[path = "admission_operation_store_tests.rs"]
1208#[allow(clippy::expect_used, clippy::unwrap_used)]
1209mod tests;