Skip to main content

chio_store_sqlite/
channel_lifecycle_store.rs

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            &current_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;