1use std::sync::{Arc, Mutex, MutexGuard};
2
3use chio_core::canonical::canonical_json_bytes;
4use chio_core::economic_continuity::{
5 verify_economic_state_batch_advance, verify_economic_state_batch_commit,
6 verify_economic_state_view, EconomicEffectStateV1, VerifiedEconomicStateBatchAdvance,
7 VerifiedEconomicStateView,
8};
9use chio_core::{sha256_hex, StoreMutationFence};
10use chio_kernel::admission_operation::{
11 expected_dispatch_committed_version, AdmissionAttachment, AdmissionDigest,
12 AdmissionOperationCommand, AdmissionOperationId, AdmissionOperationState,
13 AdmissionOperationStoreError, AdmissionOperationV1, AdmissionRecoveryLease,
14};
15use chio_settle::channel::{
16 derive_channel_reservation_id, ChannelEscrowReservationStatusV1, ChannelLifecycleStatusV1,
17 ChannelPreparedReservationV1, ChannelTransitionReplayAuthorityPinsV1,
18 ChannelTransitionReplayKindV1, ChannelTransitionReplayVerifierV1, RetainedChannelStateV1,
19 SignedChannelReservationV1, VerifiedAdmittedChannelReservationV1,
20 VerifiedChannelPreparedReservationV1, CHANNEL_SERVICE_DISPATCH_EFFECT_KIND,
21 CHANNEL_TRANSITION_REPLAY_FORMAT,
22};
23use rusqlite::{
24 params, Connection, ErrorCode, OptionalExtension, Row, Transaction, TransactionBehavior,
25};
26use serde::Serialize;
27
28use crate::serving_owner::{SqliteServingOwner, SqliteServingOwnerError};
29use crate::{
30 EconomicOperationStageContext, EconomicStateCacheError, EconomicStateStageDescriptor,
31 EconomicStateStageRecord, EconomicStateStageStatus,
32};
33
34mod prepared;
35mod reservation;
36mod schema;
37mod terminal;
38
39use prepared::*;
40use reservation::*;
41pub(crate) use schema::{initialize_channel_lifecycle_schema, verify_channel_lifecycle_invariants};
42pub(crate) use terminal::{
43 consume_channel_terminal_projection_tx, verify_consumed_channel_terminal_projection_tx,
44};
45
46const CHANNEL_LIFECYCLE_SCHEMA_KEY: &str = "channel_lifecycle";
47pub(crate) const CHANNEL_LIFECYCLE_SUPPORTED_SCHEMA_VERSION: i32 = 1;
48const CHANNEL_LIFECYCLE_SCHEMA_ANCHORS: &[&str] =
49 &["channel_lifecycle_records", "admission_operations"];
50const CHANNEL_LIFECYCLE_SCHEMA: &str = include_str!("channel_lifecycle_store.sql");
51const MAX_CHANNEL_ARTIFACT_BYTES: usize = 1024 * 1024;
52const MAX_CHANNEL_PREPARED_PLAN_BYTES: usize = 4 * 1024 * 1024;
53
54#[derive(Debug, thiserror::Error)]
55pub enum ChannelLifecycleStoreError {
56 #[error("channel lifecycle store is unavailable: {0}")]
57 Unavailable(String),
58 #[error("channel lifecycle store mutation was fenced")]
59 Fenced,
60 #[error("channel lifecycle record was not found")]
61 NotFound,
62 #[error("channel lifecycle record conflicts with retained state")]
63 Conflict,
64 #[error("channel lifecycle invariant failed: {0}")]
65 Invalid(String),
66 #[error("channel lifecycle durable outcome is unknown: {0}")]
67 OutcomeUnknown(String),
68}
69
70#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct ChannelPreparedAdmissionRecordV1 {
72 operation: AdmissionOperationV1,
73 plan: ChannelPreparedReservationV1,
74 plan_digest: String,
75 store_fence: StoreMutationFence,
76 created_at_unix_ms: u64,
77}
78
79impl ChannelPreparedAdmissionRecordV1 {
80 #[must_use]
81 pub const fn operation(&self) -> &AdmissionOperationV1 {
82 &self.operation
83 }
84
85 #[must_use]
86 pub const fn plan(&self) -> &ChannelPreparedReservationV1 {
87 &self.plan
88 }
89
90 #[must_use]
91 pub fn plan_digest(&self) -> &str {
92 &self.plan_digest
93 }
94
95 #[must_use]
96 pub const fn store_fence(&self) -> &StoreMutationFence {
97 &self.store_fence
98 }
99
100 #[must_use]
101 pub const fn created_at_unix_ms(&self) -> u64 {
102 self.created_at_unix_ms
103 }
104}
105
106#[derive(Debug, Clone, PartialEq, Eq)]
107pub enum ChannelPreparedBeginResult {
108 Created(ChannelPreparedAdmissionRecordV1),
109 ExactReplay(ChannelPreparedAdmissionRecordV1),
110 Conflict {
111 existing_operation_id: AdmissionOperationId,
112 },
113}
114
115#[derive(Debug, Clone, Copy, PartialEq, Eq)]
116pub enum ChannelReservationDispositionV1 {
117 PendingAnchor,
118 Live,
119 Consumed,
120 Cancelled,
121 Incident,
122}
123
124impl ChannelReservationDispositionV1 {
125 fn parse(value: &str) -> Result<Self, ChannelLifecycleStoreError> {
126 match value {
127 "pending_anchor" => Ok(Self::PendingAnchor),
128 "live" => Ok(Self::Live),
129 "consumed" => Ok(Self::Consumed),
130 "cancelled" => Ok(Self::Cancelled),
131 "incident" => Ok(Self::Incident),
132 _ => Err(invalid(
133 "retained channel reservation disposition is invalid",
134 )),
135 }
136 }
137}
138
139#[derive(Debug, Clone)]
140pub struct ChannelReservationStageRecordV1 {
141 operation: AdmissionOperationV1,
142 reservation: SignedChannelReservationV1,
143 authority_pins: ChannelTransitionReplayAuthorityPinsV1,
144 replay_bytes: Vec<u8>,
145 economic_stage: EconomicStateStageRecord,
146 disposition: ChannelReservationDispositionV1,
147 record_version: u64,
148 updated_at_unix_ms: u64,
149}
150
151impl ChannelReservationStageRecordV1 {
152 #[must_use]
153 pub const fn operation(&self) -> &AdmissionOperationV1 {
154 &self.operation
155 }
156
157 #[must_use]
158 pub const fn reservation(&self) -> &SignedChannelReservationV1 {
159 &self.reservation
160 }
161
162 #[must_use]
163 pub const fn authority_pins(&self) -> &ChannelTransitionReplayAuthorityPinsV1 {
164 &self.authority_pins
165 }
166
167 #[must_use]
168 pub fn replay_bytes(&self) -> &[u8] {
169 &self.replay_bytes
170 }
171
172 #[must_use]
173 pub const fn economic_stage(&self) -> &EconomicStateStageRecord {
174 &self.economic_stage
175 }
176
177 #[must_use]
178 pub const fn disposition(&self) -> ChannelReservationDispositionV1 {
179 self.disposition
180 }
181
182 #[must_use]
183 pub const fn record_version(&self) -> u64 {
184 self.record_version
185 }
186
187 #[must_use]
188 pub const fn updated_at_unix_ms(&self) -> u64 {
189 self.updated_at_unix_ms
190 }
191}
192
193#[derive(Clone)]
194pub struct SqliteChannelLifecycleStore {
195 connection: Arc<Mutex<Connection>>,
196 serving_owner: Arc<SqliteServingOwner>,
197}
198
199struct EncodedPreparedPlan {
200 plan_digest: String,
201 plan_json: Vec<u8>,
202 open_intent_digest: String,
203 open_intent_json: Vec<u8>,
204 open_digest: String,
205 open_json: Vec<u8>,
206 prior_state_kind: &'static str,
207 prior_state_digest: String,
208 prior_sequence: u64,
209 prior_state_json: Vec<u8>,
210 reservation_proposal_digest: String,
211 lifecycle_json: Vec<u8>,
212 escrow_json: Vec<u8>,
213}
214
215struct StoredPreparedPlan {
216 request_id: String,
217 request_namespace_digest: String,
218 request_binding_digest: String,
219 provider_binding_digest: String,
220 reservation_id: String,
221 channel_id: String,
222 open_digest: String,
223 prior_state_digest: String,
224 prior_sequence: u64,
225 reservation_proposal_digest: String,
226 lifecycle_state: String,
227 state_version: u64,
228 lifecycle_fence: u64,
229 live_reservation_id: Option<String>,
230 lifecycle_operation_id: Option<String>,
231 channel_head_digest: String,
232 escrow_head_digest: String,
233 checkpoint_sequence: u64,
234 checkpoint_digest: String,
235 plan_digest: String,
236 plan_json: Vec<u8>,
237 store_fence: StoreMutationFence,
238 created_at_unix_ms: u64,
239}
240
241impl SqliteChannelLifecycleStore {
242 pub(crate) fn open_alongside(
243 connection: Arc<Mutex<Connection>>,
244 serving_owner: Arc<SqliteServingOwner>,
245 ) -> Self {
246 Self {
247 connection,
248 serving_owner,
249 }
250 }
251
252 #[must_use]
253 pub fn mutation_fence(&self) -> StoreMutationFence {
254 self.serving_owner.fence.clone()
255 }
256
257 pub fn verify_invariants(&self) -> Result<(), ChannelLifecycleStoreError> {
258 let connection = self.connection()?;
259 verify_channel_lifecycle_invariants(&connection)
260 .map_err(|error| ChannelLifecycleStoreError::Unavailable(error.to_string()))
261 }
262
263 pub fn begin_channel_prepared(
264 &self,
265 operation: &AdmissionOperationV1,
266 prepared: &VerifiedChannelPreparedReservationV1,
267 fence: &StoreMutationFence,
268 trusted_now_unix_ms: u64,
269 ) -> Result<ChannelPreparedBeginResult, ChannelLifecycleStoreError> {
270 self.begin_channel_prepared_inner(
271 operation,
272 prepared.prepared(),
273 fence,
274 trusted_now_unix_ms,
275 )
276 }
277
278 pub fn load_channel_prepared(
279 &self,
280 operation_id: &AdmissionOperationId,
281 ) -> Result<Option<ChannelPreparedAdmissionRecordV1>, ChannelLifecycleStoreError> {
282 let mut connection = self.connection()?;
283 let transaction = connection
284 .transaction_with_behavior(TransactionBehavior::Deferred)
285 .map_err(sqlite_error)?;
286 crate::admission_operation_store::verify_active_owner(
287 &transaction,
288 &self.serving_owner,
289 None,
290 )
291 .map_err(admission_error)?;
292 self.serving_owner
293 .verify_authority_anchor(&transaction)
294 .map_err(owner_error)?;
295 let operation = crate::admission_operation_store::load_operation_for_participant_tx(
296 &transaction,
297 operation_id,
298 )
299 .map_err(admission_error)?;
300 let record = match operation {
301 Some(operation) => {
302 let requires_channel = operation.binding().participant_requirements().channel;
303 let require_base_lifecycle = operation.channel_reservation_digest().is_none();
304 let record = load_prepared_record(&transaction, operation, require_base_lifecycle)?;
305 if requires_channel && record.is_none() {
306 return Err(ChannelLifecycleStoreError::NotFound);
307 }
308 record
309 }
310 None => None,
311 };
312 transaction.commit().map_err(sqlite_error)?;
313 Ok(record)
314 }
315
316 #[allow(clippy::too_many_arguments)]
317 pub fn stage_channel_reservation(
318 &self,
319 advance: &VerifiedEconomicStateBatchAdvance,
320 operation: &AdmissionOperationV1,
321 recovery_lease: &AdmissionRecoveryLease,
322 replay_bytes: &[u8],
323 expected_authority_pins: &ChannelTransitionReplayAuthorityPinsV1,
324 fence: &StoreMutationFence,
325 trusted_now_unix_ms: u64,
326 ) -> Result<ChannelReservationStageRecordV1, ChannelLifecycleStoreError> {
327 if fence != &self.serving_owner.fence || recovery_lease.store_fence() != fence {
328 return Err(ChannelLifecycleStoreError::Fenced);
329 }
330 let authority_pins_json = expected_authority_pins
331 .canonical_bytes()
332 .map_err(channel_error)?;
333 let authority_pins_digest = expected_authority_pins.digest().map_err(channel_error)?;
334 let replay = ChannelTransitionReplayVerifierV1::from_canonical_bytes(
335 replay_bytes,
336 expected_authority_pins,
337 )
338 .map_err(channel_error)?;
339 if replay.descriptor().kind() != ChannelTransitionReplayKindV1::Reservation {
340 return Err(invalid(
341 "channel reservation stage requires a reservation replay",
342 ));
343 }
344 verify_economic_state_batch_advance(
345 advance.current(),
346 advance.batch().clone(),
347 &expected_authority_pins.anchor_pins(),
348 &replay,
349 )
350 .map_err(|error| invalid(error.to_string()))?;
351 let proposal = replay.verified_reservation_proposal();
352 let reservation = proposal.artifact();
353 let reservation_body = &reservation.body;
354 let reservation_digest = reservation.digest().map_err(channel_error)?;
355 let reservation_json = encode(
356 reservation,
357 MAX_CHANNEL_ARTIFACT_BYTES,
358 "signed channel reservation",
359 )?;
360 let replay_protocol_digest = replay.descriptor().digest().to_owned();
361 let replay_content_digest = sha256_hex(replay_bytes);
362 let descriptor = EconomicStateStageDescriptor::new(
363 CHANNEL_TRANSITION_REPLAY_FORMAT,
364 replay.descriptor().key(),
365 replay.descriptor(),
366 )
367 .map_err(economic_error)?;
368 if descriptor.digest() != replay_content_digest
369 || replay.descriptor().expected_batch_digest() != advance.batch().checkpoint_digest
370 || replay.descriptor().request().request_id != reservation_body.request_id
371 || reservation_body.operation_id != operation.binding().operation_id().as_str()
372 || trusted_now_unix_ms < proposal.accepted_at_unix_ms()
373 || trusted_now_unix_ms >= reservation_body.expires_at_unix_ms
374 {
375 return Err(invalid(
376 "channel reservation replay binding is inconsistent",
377 ));
378 }
379 let ready_effect_head_digest = exact_ready_effect_head_digest(
380 advance,
381 &reservation_body.operation_id,
382 &reservation_digest,
383 )?;
384 let mut connection = self.connection()?;
385 let transaction = connection
386 .transaction_with_behavior(TransactionBehavior::Immediate)
387 .map_err(sqlite_error)?;
388 crate::admission_operation_store::verify_active_owner(
389 &transaction,
390 &self.serving_owner,
391 Some(fence),
392 )
393 .map_err(admission_error)?;
394 self.serving_owner
395 .verify_authority_anchor(&transaction)
396 .map_err(owner_error)?;
397 if let Some(existing) = load_channel_reservation_tx(
398 &transaction,
399 operation.binding().operation_id(),
400 expected_authority_pins,
401 )? {
402 qualify_exact_staged_replay(
403 &existing,
404 operation,
405 advance,
406 replay_bytes,
407 &reservation_digest,
408 )?;
409 qualify_exact_recovery_authority(&existing, recovery_lease, trusted_now_unix_ms)?;
410 crate::admission_operation_store::verify_participant_recovery_tx(
411 &transaction,
412 &self.serving_owner,
413 operation,
414 recovery_lease,
415 trusted_now_unix_ms,
416 )
417 .map_err(admission_error)?;
418 transaction.commit().map_err(sqlite_error)?;
419 return Ok(existing);
420 }
421 let stored_operation = crate::admission_operation_store::load_operation_for_participant_tx(
422 &transaction,
423 operation.binding().operation_id(),
424 )
425 .map_err(admission_error)?
426 .ok_or(ChannelLifecycleStoreError::NotFound)?;
427 if stored_operation != *operation
428 || !matches!(
429 operation.state(),
430 AdmissionOperationState::BudgetAuthorized
431 | AdmissionOperationState::ApprovalReserved
432 )
433 || operation.channel_reservation_digest().is_some()
434 {
435 return Err(ChannelLifecycleStoreError::Fenced);
436 }
437 let prepared = load_prepared_record(&transaction, stored_operation.clone(), true)?
438 .ok_or(ChannelLifecycleStoreError::NotFound)?;
439 qualify_reservation_against_prepared(
440 &prepared,
441 reservation,
442 replay.descriptor().request(),
443 )?;
444 let operation_context = EconomicOperationStageContext::new(operation, recovery_lease)
445 .with_not_after_unix_ms(reservation_body.expires_at_unix_ms)
446 .map_err(economic_error)?;
447 let economic_stage = crate::economic_state_cache::stage_channel_batch_in_transaction(
448 &transaction,
449 advance,
450 operation_context,
451 descriptor,
452 fence,
453 trusted_now_unix_ms,
454 &self.serving_owner,
455 )
456 .map_err(economic_error)?;
457 let participant_digest = channel_reservation_participant_digest(
458 &prepared,
459 &reservation_digest,
460 &authority_pins_digest,
461 &replay_protocol_digest,
462 &replay_content_digest,
463 &economic_stage,
464 &ready_effect_head_digest,
465 )?;
466 insert_pending_reservation_tx(
467 &transaction,
468 &prepared,
469 reservation,
470 &reservation_digest,
471 &reservation_json,
472 &authority_pins_digest,
473 &authority_pins_json,
474 replay_bytes,
475 &replay_protocol_digest,
476 &replay_content_digest,
477 &economic_stage,
478 &ready_effect_head_digest,
479 fence,
480 trusted_now_unix_ms,
481 )?;
482 crate::admission_operation_store::append_participant_update_tx(
483 &transaction,
484 &self.serving_owner,
485 operation,
486 recovery_lease,
487 &participant_digest,
488 trusted_now_unix_ms,
489 )
490 .map_err(admission_error)?;
491 transaction.commit().map_err(|error| {
492 owner_error(self.serving_owner.outcome_unknown(format!(
493 "sqlite channel reservation stage commit outcome is unknown: {error}"
494 )))
495 })?;
496 self.serving_owner
497 .sync_authority_anchor(&connection)
498 .map_err(owner_error)?;
499 Ok(ChannelReservationStageRecordV1 {
500 operation: operation.clone(),
501 reservation: reservation.clone(),
502 authority_pins: expected_authority_pins.clone(),
503 replay_bytes: replay_bytes.to_vec(),
504 economic_stage,
505 disposition: ChannelReservationDispositionV1::PendingAnchor,
506 record_version: 1,
507 updated_at_unix_ms: trusted_now_unix_ms,
508 })
509 }
510
511 pub fn load_channel_reservation(
512 &self,
513 operation_id: &AdmissionOperationId,
514 expected_authority_pins: &ChannelTransitionReplayAuthorityPinsV1,
515 ) -> Result<Option<ChannelReservationStageRecordV1>, ChannelLifecycleStoreError> {
516 let mut connection = self.connection()?;
517 let transaction = connection
518 .transaction_with_behavior(TransactionBehavior::Deferred)
519 .map_err(sqlite_error)?;
520 crate::admission_operation_store::verify_active_owner(
521 &transaction,
522 &self.serving_owner,
523 None,
524 )
525 .map_err(admission_error)?;
526 self.serving_owner
527 .verify_authority_anchor(&transaction)
528 .map_err(owner_error)?;
529 let record =
530 load_channel_reservation_tx(&transaction, operation_id, expected_authority_pins)?;
531 transaction.commit().map_err(sqlite_error)?;
532 Ok(record)
533 }
534
535 #[allow(clippy::too_many_arguments)]
536 pub fn record_channel_anchor_advanced(
537 &self,
538 operation_id: &AdmissionOperationId,
539 advance: &VerifiedEconomicStateBatchAdvance,
540 committed: &VerifiedEconomicStateView,
541 expected_authority_pins: &ChannelTransitionReplayAuthorityPinsV1,
542 fence: &StoreMutationFence,
543 trusted_now_unix_ms: u64,
544 ) -> Result<ChannelReservationStageRecordV1, ChannelLifecycleStoreError> {
545 if fence != &self.serving_owner.fence {
546 return Err(ChannelLifecycleStoreError::Fenced);
547 }
548 verify_economic_state_batch_commit(
549 advance,
550 committed,
551 &expected_authority_pins.anchor_pins(),
552 )
553 .map_err(|error| invalid(error.to_string()))?;
554 let mut connection = self.connection()?;
555 let transaction = connection
556 .transaction_with_behavior(TransactionBehavior::Immediate)
557 .map_err(sqlite_error)?;
558 crate::admission_operation_store::verify_active_owner(
559 &transaction,
560 &self.serving_owner,
561 Some(fence),
562 )
563 .map_err(admission_error)?;
564 self.serving_owner
565 .verify_authority_anchor(&transaction)
566 .map_err(owner_error)?;
567 let record =
568 load_channel_reservation_tx(&transaction, operation_id, expected_authority_pins)?
569 .ok_or(ChannelLifecycleStoreError::NotFound)?;
570 if record.disposition() != ChannelReservationDispositionV1::PendingAnchor
571 || record.economic_stage().base_view() != advance.current().view()
572 || record.economic_stage().batch() != advance.batch()
573 {
574 return Err(ChannelLifecycleStoreError::Conflict);
575 }
576 let replay = ChannelTransitionReplayVerifierV1::from_canonical_bytes(
577 record.replay_bytes(),
578 expected_authority_pins,
579 )
580 .map_err(channel_error)?;
581 let admitted = replay
582 .verify_committed_reservation(committed)
583 .map_err(channel_error)?;
584 if admitted.artifact() != record.reservation() {
585 return Err(invalid(
586 "committed channel reservation differs from staged evidence",
587 ));
588 }
589 let already_advanced =
590 record.economic_stage().status() == EconomicStateStageStatus::EconomicAnchorAdvanced;
591 let economic_stage =
592 crate::economic_state_cache::record_channel_anchor_advanced_in_transaction(
593 &transaction,
594 advance,
595 committed,
596 trusted_now_unix_ms,
597 &self.serving_owner,
598 )
599 .map_err(economic_error)?;
600 transaction.commit().map_err(|error| {
601 owner_error(self.serving_owner.outcome_unknown(format!(
602 "sqlite channel anchor record commit outcome is unknown: {error}"
603 )))
604 })?;
605 if !already_advanced {
606 self.serving_owner
607 .sync_authority_anchor(&connection)
608 .map_err(owner_error)?;
609 }
610 Ok(ChannelReservationStageRecordV1 {
611 economic_stage,
612 ..record
613 })
614 }
615
616 pub fn finalize_channel_reservation(
617 &self,
618 operation_id: &AdmissionOperationId,
619 recovery_lease: &AdmissionRecoveryLease,
620 expected_authority_pins: &ChannelTransitionReplayAuthorityPinsV1,
621 fence: &StoreMutationFence,
622 trusted_now_unix_ms: u64,
623 ) -> Result<ChannelReservationStageRecordV1, ChannelLifecycleStoreError> {
624 if fence != &self.serving_owner.fence || recovery_lease.store_fence() != fence {
625 return Err(ChannelLifecycleStoreError::Fenced);
626 }
627 let mut connection = self.connection()?;
628 let transaction = connection
629 .transaction_with_behavior(TransactionBehavior::Immediate)
630 .map_err(sqlite_error)?;
631 crate::admission_operation_store::verify_active_owner(
632 &transaction,
633 &self.serving_owner,
634 Some(fence),
635 )
636 .map_err(admission_error)?;
637 self.serving_owner
638 .verify_authority_anchor(&transaction)
639 .map_err(owner_error)?;
640 let current_operation =
641 crate::admission_operation_store::load_operation_for_participant_tx(
642 &transaction,
643 operation_id,
644 )
645 .map_err(admission_error)?
646 .ok_or(ChannelLifecycleStoreError::NotFound)?;
647 crate::admission_operation_store::verify_participant_recovery_tx(
648 &transaction,
649 &self.serving_owner,
650 ¤t_operation,
651 recovery_lease,
652 trusted_now_unix_ms,
653 )
654 .map_err(admission_error)?;
655 let record =
656 load_channel_reservation_tx(&transaction, operation_id, expected_authority_pins)?
657 .ok_or(ChannelLifecycleStoreError::NotFound)?;
658 let committed = record
659 .economic_stage()
660 .committed_view()
661 .cloned()
662 .ok_or_else(|| invalid("channel reservation anchor is not retained"))?;
663 let committed =
664 verify_economic_state_view(committed, &expected_authority_pins.anchor_pins())
665 .map_err(|error| invalid(error.to_string()))?;
666 let replay = ChannelTransitionReplayVerifierV1::from_canonical_bytes(
667 record.replay_bytes(),
668 expected_authority_pins,
669 )
670 .map_err(channel_error)?;
671 let admitted = replay
672 .verify_committed_reservation(&committed)
673 .map_err(channel_error)?;
674 if admitted.artifact() != record.reservation() {
675 return Err(invalid(
676 "anchored channel reservation differs from retained evidence",
677 ));
678 }
679 if record.disposition() == ChannelReservationDispositionV1::Live {
680 qualify_finalized_channel_tx(&transaction, &record, &admitted)?;
681 transaction.commit().map_err(sqlite_error)?;
682 return Ok(record);
683 }
684 if record.disposition() != ChannelReservationDispositionV1::PendingAnchor
685 || record.economic_stage().status() != EconomicStateStageStatus::EconomicAnchorAdvanced
686 {
687 return Err(ChannelLifecycleStoreError::Conflict);
688 }
689 let prepared = load_prepared_record(&transaction, record.operation().clone(), true)?
690 .ok_or(ChannelLifecycleStoreError::NotFound)?;
691 let reservation_digest = record.reservation().digest().map_err(channel_error)?;
692 let participant_digest = channel_reservation_participant_digest(
693 &prepared,
694 &reservation_digest,
695 &expected_authority_pins.digest().map_err(channel_error)?,
696 replay.descriptor().digest(),
697 &sha256_hex(record.replay_bytes()),
698 record.economic_stage(),
699 admitted.ready_effect_head_digest(),
700 )?;
701 let command = AdmissionOperationCommand::new(
702 operation_id.clone(),
703 record.operation().version(),
704 recovery_lease.clone(),
705 vec![AdmissionAttachment::ChannelReservationDigest(
706 AdmissionDigest::try_new("channel_reservation_digest", reservation_digest)
707 .map_err(|error| invalid(error.to_string()))?,
708 )],
709 Some(AdmissionOperationState::ReadyToDispatch),
710 None,
711 None,
712 )
713 .map_err(|error| invalid(error.to_string()))?;
714 let updated_operation =
715 crate::admission_operation_store::finalize_channel_reservation_operation_tx(
716 &transaction,
717 &self.serving_owner,
718 record.operation(),
719 &command,
720 &participant_digest,
721 trusted_now_unix_ms,
722 )
723 .map_err(admission_error)?;
724 publish_live_lifecycle_tx(
725 &transaction,
726 &prepared,
727 &admitted,
728 fence,
729 trusted_now_unix_ms,
730 )?;
731 let changed = transaction
732 .execute(
733 r#"
734 UPDATE channel_reservation_records
735 SET disposition = 'live', record_version = record_version + 1,
736 store_uuid = ?1, store_lease_id = ?2, store_owner_epoch = ?3,
737 updated_at_unix_ms = ?4
738 WHERE operation_id = ?5 AND reservation_id = ?6
739 AND disposition = 'pending_anchor' AND record_version = ?7
740 AND stage_batch_id = ?8
741 "#,
742 params![
743 &fence.store_uuid,
744 &fence.lease_id,
745 sqlite_i64(fence.owner_epoch, "store_owner_epoch")?,
746 sqlite_i64(trusted_now_unix_ms, "updated_at_unix_ms")?,
747 operation_id.as_str(),
748 &record.reservation().body.reservation_id,
749 sqlite_i64(record.record_version(), "record_version")?,
750 &record.economic_stage().batch().batch_id,
751 ],
752 )
753 .map_err(sqlite_error)?;
754 if changed != 1 {
755 return Err(ChannelLifecycleStoreError::Fenced);
756 }
757 let economic_stage = crate::economic_state_cache::finalize_stage_in_transaction(
758 &transaction,
759 &record.economic_stage().batch().batch_id,
760 &self.serving_owner,
761 trusted_now_unix_ms,
762 )
763 .map_err(economic_error)?;
764 transaction.commit().map_err(|error| {
765 owner_error(self.serving_owner.outcome_unknown(format!(
766 "sqlite channel reservation finalization commit outcome is unknown: {error}"
767 )))
768 })?;
769 self.serving_owner
770 .sync_authority_anchor(&connection)
771 .map_err(owner_error)?;
772 Ok(ChannelReservationStageRecordV1 {
773 operation: updated_operation,
774 economic_stage,
775 disposition: ChannelReservationDispositionV1::Live,
776 record_version: record
777 .record_version()
778 .checked_add(1)
779 .ok_or_else(|| invalid("channel reservation version overflowed"))?,
780 updated_at_unix_ms: trusted_now_unix_ms,
781 ..record
782 })
783 }
784
785 fn begin_channel_prepared_inner(
786 &self,
787 operation: &AdmissionOperationV1,
788 prepared: &ChannelPreparedReservationV1,
789 fence: &StoreMutationFence,
790 trusted_now_unix_ms: u64,
791 ) -> Result<ChannelPreparedBeginResult, ChannelLifecycleStoreError> {
792 if fence != &self.serving_owner.fence {
793 return Err(ChannelLifecycleStoreError::Fenced);
794 }
795 if trusted_now_unix_ms < prepared.observed_at_unix_ms
796 || trusted_now_unix_ms >= prepared.reservation.expires_at_unix_ms
797 {
798 return Err(invalid(
799 "trusted time is outside the authenticated channel plan window",
800 ));
801 }
802 let proposal_digest = prepared
803 .reservation
804 .proposal_digest()
805 .map_err(channel_error)?;
806 let operation = operation
807 .clone()
808 .with_initial_channel_reservation_proposal_digest(
809 AdmissionDigest::try_new("channel_reservation_proposal_digest", proposal_digest)
810 .map_err(|error| invalid(error.to_string()))?,
811 )
812 .map_err(|error| invalid(error.to_string()))?;
813 let encoded = encode_prepared_plan(&operation, prepared, fence)?;
814 let mut connection = self.connection()?;
815 let transaction = connection
816 .transaction_with_behavior(TransactionBehavior::Immediate)
817 .map_err(sqlite_error)?;
818 crate::admission_operation_store::verify_active_owner(
819 &transaction,
820 &self.serving_owner,
821 Some(fence),
822 )
823 .map_err(admission_error)?;
824 self.serving_owner
825 .verify_authority_anchor(&transaction)
826 .map_err(owner_error)?;
827 match crate::admission_operation_store::begin_prepared_operation_tx(
828 &transaction,
829 &operation,
830 fence,
831 trusted_now_unix_ms,
832 )
833 .map_err(admission_error)?
834 {
835 crate::admission_operation_store::PreparedAdmissionBeginTxResult::Conflict {
836 existing_operation_id,
837 } => {
838 transaction.commit().map_err(sqlite_error)?;
839 return Ok(ChannelPreparedBeginResult::Conflict {
840 existing_operation_id,
841 });
842 }
843 crate::admission_operation_store::PreparedAdmissionBeginTxResult::ExactReplay {
844 operation: stored_operation,
845 ..
846 } => {
847 let existing = load_prepared_record(&transaction, *stored_operation, true)?;
848 let result = match existing {
849 Some(record)
850 if record.plan_digest == encoded.plan_digest
851 && record.plan == *prepared =>
852 {
853 ChannelPreparedBeginResult::ExactReplay(record)
854 }
855 Some(record) => ChannelPreparedBeginResult::Conflict {
856 existing_operation_id: record.operation.binding().operation_id().clone(),
857 },
858 None => ChannelPreparedBeginResult::Conflict {
859 existing_operation_id: operation.binding().operation_id().clone(),
860 },
861 };
862 transaction.commit().map_err(sqlite_error)?;
863 return Ok(result);
864 }
865 crate::admission_operation_store::PreparedAdmissionBeginTxResult::Created {
866 encoded: encoded_operation,
867 } => {
868 insert_or_verify_state(
869 &transaction,
870 prepared,
871 &encoded,
872 fence,
873 trusted_now_unix_ms,
874 )?;
875 insert_or_verify_lifecycle(
876 &transaction,
877 prepared,
878 &encoded,
879 fence,
880 trusted_now_unix_ms,
881 )?;
882 insert_prepared_plan(
883 &transaction,
884 &operation,
885 prepared,
886 &encoded,
887 fence,
888 trusted_now_unix_ms,
889 )?;
890 crate::admission_operation_store::append_operation_commit_with_participant(
891 &transaction,
892 &operation,
893 &encoded_operation,
894 None,
895 "begin",
896 Some(&encoded.plan_digest),
897 &self.serving_owner,
898 trusted_now_unix_ms,
899 )
900 .map_err(admission_error)?;
901 }
902 }
903 transaction.commit().map_err(|error| {
904 owner_error(self.serving_owner.outcome_unknown(format!(
905 "sqlite channel prepared commit outcome is unknown: {error}"
906 )))
907 })?;
908 self.serving_owner
909 .sync_authority_anchor(&connection)
910 .map_err(owner_error)?;
911 Ok(ChannelPreparedBeginResult::Created(
912 ChannelPreparedAdmissionRecordV1 {
913 operation,
914 plan: prepared.clone(),
915 plan_digest: encoded.plan_digest,
916 store_fence: fence.clone(),
917 created_at_unix_ms: trusted_now_unix_ms,
918 },
919 ))
920 }
921
922 fn connection(&self) -> Result<MutexGuard<'_, Connection>, ChannelLifecycleStoreError> {
923 self.connection.lock().map_err(|_| {
924 ChannelLifecycleStoreError::Unavailable(
925 "sqlite channel lifecycle lock poisoned".to_owned(),
926 )
927 })
928 }
929}
930
931fn encode(
932 value: &impl Serialize,
933 maximum: usize,
934 label: &'static str,
935) -> Result<Vec<u8>, ChannelLifecycleStoreError> {
936 let encoded = canonical_json_bytes(value)
937 .map_err(|error| invalid(format!("{label} encoding failed: {error}")))?;
938 if encoded.is_empty() || encoded.len() > maximum {
939 return Err(invalid(format!("{label} exceeds its size limit")));
940 }
941 Ok(encoded)
942}
943
944fn sqlite_i64(value: u64, field: &'static str) -> Result<i64, ChannelLifecycleStoreError> {
945 i64::try_from(value).map_err(|_| invalid(format!("{field} exceeds SQLite integer range")))
946}
947
948fn stored_u64(value: i64, field: &'static str) -> Result<u64, ChannelLifecycleStoreError> {
949 u64::try_from(value).map_err(|_| invalid(format!("{field} is negative")))
950}
951
952fn stored_u64_sql(value: i64, field: &'static str) -> rusqlite::Result<u64> {
953 u64::try_from(value).map_err(|_| {
954 rusqlite::Error::FromSqlConversionFailure(
955 0,
956 rusqlite::types::Type::Integer,
957 std::io::Error::new(
958 std::io::ErrorKind::InvalidData,
959 format!("{field} is negative"),
960 )
961 .into(),
962 )
963 })
964}
965
966fn sqlite_error(error: rusqlite::Error) -> ChannelLifecycleStoreError {
967 if error.sqlite_error_code() == Some(ErrorCode::ConstraintViolation) {
968 ChannelLifecycleStoreError::Conflict
969 } else {
970 ChannelLifecycleStoreError::Unavailable(error.to_string())
971 }
972}
973
974fn admission_error(error: AdmissionOperationStoreError) -> ChannelLifecycleStoreError {
975 match error {
976 AdmissionOperationStoreError::Unavailable(detail) => {
977 ChannelLifecycleStoreError::Unavailable(detail)
978 }
979 AdmissionOperationStoreError::Fenced => ChannelLifecycleStoreError::Fenced,
980 AdmissionOperationStoreError::NotFound => ChannelLifecycleStoreError::NotFound,
981 AdmissionOperationStoreError::Invariant(detail) => invalid(detail),
982 AdmissionOperationStoreError::OutcomeUnknown(detail) => {
983 ChannelLifecycleStoreError::OutcomeUnknown(detail)
984 }
985 AdmissionOperationStoreError::Operation(error) => invalid(error.to_string()),
986 }
987}
988
989fn channel_error(error: chio_settle::channel::ChannelError) -> ChannelLifecycleStoreError {
990 invalid(error.to_string())
991}
992
993fn economic_error(error: EconomicStateCacheError) -> ChannelLifecycleStoreError {
994 match error {
995 EconomicStateCacheError::Unavailable(detail) => {
996 ChannelLifecycleStoreError::Unavailable(detail)
997 }
998 EconomicStateCacheError::Fenced => ChannelLifecycleStoreError::Fenced,
999 EconomicStateCacheError::NotFound => ChannelLifecycleStoreError::NotFound,
1000 EconomicStateCacheError::OutcomeUnknown(detail) => {
1001 ChannelLifecycleStoreError::OutcomeUnknown(detail)
1002 }
1003 error => invalid(error.to_string()),
1004 }
1005}
1006
1007fn owner_error(error: SqliteServingOwnerError) -> ChannelLifecycleStoreError {
1008 match error {
1009 SqliteServingOwnerError::OutcomeUnknown(detail) => {
1010 ChannelLifecycleStoreError::OutcomeUnknown(detail)
1011 }
1012 error => ChannelLifecycleStoreError::Unavailable(error.to_string()),
1013 }
1014}
1015
1016fn invalid(detail: impl Into<String>) -> ChannelLifecycleStoreError {
1017 ChannelLifecycleStoreError::Invalid(detail.into())
1018}
1019
1020fn retained_projection_error(error: ChannelLifecycleStoreError) -> ChannelLifecycleStoreError {
1021 match error {
1022 ChannelLifecycleStoreError::Conflict | ChannelLifecycleStoreError::NotFound => {
1023 invalid("retained channel prepared projection is inconsistent")
1024 }
1025 error => error,
1026 }
1027}
1028
1029#[cfg(test)]
1030#[path = "channel_lifecycle_store_tests.rs"]
1031mod tests;