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;