Skip to main content

chio_store_sqlite/frost_store/
rotation.rs

1use chio_core::canonical::canonical_json_bytes;
2use chio_core::{sha256_hex, StoreMutationFence};
3use chio_federation::frost::{
4    FrostArtifactTrustStore, FrostEpochAnchor, FrostEpochAnchorWriter, FrostEpochCheckpointV1,
5    FrostRosterV1, FrostSessionBurnSummaryV1, VerifiedFrostEpochAdvance,
6};
7use rusqlite::{params, Connection, OptionalExtension, Row, Transaction};
8use serde::Serialize;
9
10use super::commit::{append_projection_commit, prefixed_digest, ProjectionMutation};
11use super::rotation_validation::{
12    verify_active_predecessor, verify_advance_checkpoint, verify_checkpoint_roster,
13    verify_completed_ceremony, verify_local_burn_summary, verify_stored_rotation_matches_advance,
14    verify_stored_successor,
15};
16use super::{
17    FrostActiveRosterRecord, FrostRotationRecord, FrostRotationState, FrostStoreError,
18    SqliteFrostStore, StagedFrostRotation,
19};
20
21const ROTATION_ID_PREFIX: &[u8] = b"chio.frost.rotation.id.v1\0";
22const ROTATION_RECORD_PREFIX: &[u8] = b"chio.frost.rotation-record.digest.v1\0";
23const ROSTER_HISTORY_RECORD_PREFIX: &[u8] = b"chio.frost.roster-history.digest.v1\0";
24
25#[derive(Debug, Clone)]
26pub(super) struct StoredRotation {
27    pub(super) rotation_id: String,
28    pub(super) scope_id: String,
29    pub(super) state: FrostRotationState,
30    pub(super) state_version: u64,
31    pub(super) predecessor_checkpoint_digest: String,
32    pub(super) predecessor_checkpoint_json: Vec<u8>,
33    pub(super) target_roster_digest: String,
34    pub(super) target_roster_json: Vec<u8>,
35    pub(super) target_key_epoch: u64,
36    pub(super) rotation_authorization_digest: String,
37    pub(super) burn_summary_json: Vec<u8>,
38    pub(super) old_session_burn_root: String,
39    pub(super) expected_checkpoint_sequence: u64,
40    pub(super) activation_fence: u64,
41    pub(super) clock_high_water: u64,
42    pub(super) anchored_checkpoint_digest: Option<String>,
43    pub(super) anchored_checkpoint_json: Option<Vec<u8>>,
44    pub(super) source_fence: StoreMutationFence,
45    pub(super) created_at_unix_ms: u64,
46    pub(super) updated_at_unix_ms: u64,
47    pub(super) record_digest: String,
48}
49
50#[derive(Serialize)]
51#[serde(rename_all = "camelCase")]
52struct RotationIdPreimage<'a> {
53    predecessor_checkpoint_digest: &'a str,
54    target_roster_digest: &'a str,
55    rotation_authorization_digest: &'a str,
56    old_session_burn_root: &'a str,
57    activation_fence: u64,
58    clock_high_water: u64,
59}
60
61#[derive(Serialize)]
62#[serde(rename_all = "camelCase")]
63struct RotationRecordPreimage<'a> {
64    rotation_id: &'a str,
65    scope_id: &'a str,
66    state: FrostRotationState,
67    state_version: u64,
68    predecessor_checkpoint_digest: &'a str,
69    predecessor_checkpoint_json_digest: String,
70    target_roster_digest: &'a str,
71    target_roster_json_digest: String,
72    target_key_epoch: u64,
73    rotation_authorization_digest: &'a str,
74    burn_summary_json_digest: String,
75    old_session_burn_root: &'a str,
76    expected_checkpoint_sequence: u64,
77    activation_fence: u64,
78    clock_high_water: u64,
79    anchored_checkpoint_digest: Option<&'a str>,
80    anchored_checkpoint_json_digest: Option<String>,
81    source_store_uuid: &'a str,
82    source_lease_id: &'a str,
83    source_owner_epoch: u64,
84    created_at_unix_ms: u64,
85    updated_at_unix_ms: u64,
86}
87
88#[derive(Serialize)]
89#[serde(rename_all = "camelCase")]
90struct RosterHistoryRecordPreimage<'a> {
91    scope_id: &'a str,
92    key_epoch: u64,
93    roster_id: &'a str,
94    roster_digest: &'a str,
95    roster_json_digest: String,
96    checkpoint_sequence: u64,
97    checkpoint_digest: &'a str,
98    checkpoint_json_digest: String,
99    activation_fence: u64,
100    clock_high_water: u64,
101    activated_at_unix_ms: u64,
102    source_store_uuid: &'a str,
103    source_lease_id: &'a str,
104    source_owner_epoch: u64,
105}
106
107impl SqliteFrostStore {
108    pub fn import_active_roster(
109        &self,
110        roster: &FrostRosterV1,
111        checkpoint: &FrostEpochCheckpointV1,
112        artifact_trust: &FrostArtifactTrustStore,
113        fence: &StoreMutationFence,
114        trusted_now_unix_ms: u64,
115    ) -> Result<FrostActiveRosterRecord, FrostStoreError> {
116        validate_trusted_time(trusted_now_unix_ms)?;
117        verify_checkpoint_roster(roster, checkpoint, artifact_trust, trusted_now_unix_ms)?;
118        let roster_json = canonical_json_bytes(roster).map_err(canonical_error)?;
119        let checkpoint_json = canonical_json_bytes(checkpoint).map_err(canonical_error)?;
120        let mut connection = self.connection()?;
121        let transaction = self.begin_write(&mut connection, fence)?;
122        if let Some(active) = load_active_record(&transaction, &roster.scope_id)? {
123            if active.key_epoch == roster.key_epoch
124                && active.roster_digest == roster.roster_digest
125                && active.checkpoint_digest == checkpoint.checkpoint_digest
126            {
127                transaction.commit().map_err(super::sqlite_error)?;
128                return Ok(active);
129            }
130            return Err(FrostStoreError::Conflict(
131                "active roster already exists and must advance through rotation",
132            ));
133        }
134        insert_roster_history(
135            &transaction,
136            roster,
137            &roster_json,
138            checkpoint,
139            &checkpoint_json,
140            fence,
141            trusted_now_unix_ms,
142        )?;
143        insert_active_pointer(&transaction, roster, checkpoint, fence)?;
144        let history_digest = roster_history_record_digest(
145            roster,
146            &roster_json,
147            checkpoint,
148            &checkpoint_json,
149            fence,
150            trusted_now_unix_ms,
151        )?;
152        append_projection_commit(
153            &transaction,
154            self,
155            &roster_projection_key(&roster.scope_id, roster.key_epoch),
156            ProjectionMutation {
157                sequence: 1,
158                projection_type: "rotation",
159                mutation_kind: "frost.roster.import",
160                record_digest: &history_digest,
161            },
162            fence,
163        )?;
164        self.commit_write(transaction)?;
165        self.sync_after_write(&connection)?;
166        Ok(active_record(roster, checkpoint))
167    }
168
169    pub fn load_active_roster(
170        &self,
171        scope_id: &str,
172    ) -> Result<Option<(FrostRosterV1, FrostEpochCheckpointV1)>, FrostStoreError> {
173        let mut connection = self.connection()?;
174        let transaction = self.begin_read(&mut connection, None)?;
175        let result = load_active_artifacts(&transaction, scope_id)?;
176        transaction.commit().map_err(super::sqlite_error)?;
177        Ok(result)
178    }
179
180    pub fn stage_rotation(
181        &self,
182        advance: VerifiedFrostEpochAdvance,
183        fence: &StoreMutationFence,
184        trusted_now_unix_ms: u64,
185    ) -> Result<StagedFrostRotation, FrostStoreError> {
186        validate_trusted_time(trusted_now_unix_ms)?;
187        let rotation_id = rotation_id(&advance)?;
188        let predecessor_json =
189            canonical_json_bytes(advance.predecessor()).map_err(canonical_error)?;
190        let target_roster_json =
191            canonical_json_bytes(advance.target_roster()).map_err(canonical_error)?;
192        let burn_summary_json =
193            canonical_json_bytes(advance.burn_summary()).map_err(canonical_error)?;
194        let mut connection = self.connection()?;
195        let transaction = self.begin_write(&mut connection, fence)?;
196        verify_active_predecessor(&transaction, &advance)?;
197        verify_completed_ceremony(&transaction, advance.target_roster())?;
198        verify_local_burn_summary(&transaction, advance.burn_summary())?;
199        if let Some(stored) = load_rotation_query(&transaction, &rotation_id)? {
200            verify_stored_rotation_matches_advance(&stored, &advance)?;
201            if stored.state == FrostRotationState::Discarded {
202                return Err(FrostStoreError::Conflict(
203                    "the exact rotation stage was discarded",
204                ));
205            }
206            transaction.commit().map_err(super::sqlite_error)?;
207            return Ok(StagedFrostRotation {
208                rotation_id,
209                advance,
210            });
211        }
212        let mut stored = StoredRotation {
213            rotation_id: rotation_id.clone(),
214            scope_id: advance.target_roster().scope_id.clone(),
215            state: FrostRotationState::Staged,
216            state_version: 1,
217            predecessor_checkpoint_digest: advance.predecessor().checkpoint_digest.clone(),
218            predecessor_checkpoint_json: predecessor_json,
219            target_roster_digest: advance.target_roster().roster_digest.clone(),
220            target_roster_json,
221            target_key_epoch: advance.target_roster().key_epoch,
222            rotation_authorization_digest: advance.rotation_authorization_digest().to_string(),
223            burn_summary_json,
224            old_session_burn_root: advance.burn_summary().burn_root.clone(),
225            expected_checkpoint_sequence: advance.expected_checkpoint_sequence(),
226            activation_fence: advance.activation_fence(),
227            clock_high_water: advance.clock_high_water(),
228            anchored_checkpoint_digest: None,
229            anchored_checkpoint_json: None,
230            source_fence: fence.clone(),
231            created_at_unix_ms: trusted_now_unix_ms,
232            updated_at_unix_ms: trusted_now_unix_ms,
233            record_digest: String::new(),
234        };
235        stored.record_digest = rotation_record_digest(&stored)?;
236        insert_rotation(&transaction, &stored)?;
237        append_projection_commit(
238            &transaction,
239            self,
240            &rotation_projection_key(&rotation_id),
241            ProjectionMutation {
242                sequence: 1,
243                projection_type: "rotation",
244                mutation_kind: "frost.rotation.stage",
245                record_digest: &stored.record_digest,
246            },
247            fence,
248        )?;
249        self.commit_write(transaction)?;
250        self.sync_after_write(&connection)?;
251        Ok(StagedFrostRotation {
252            rotation_id,
253            advance,
254        })
255    }
256
257    pub fn advance_rotation_anchor(
258        &self,
259        staged: &StagedFrostRotation,
260        anchor: &dyn FrostEpochAnchorWriter,
261        artifact_trust: &FrostArtifactTrustStore,
262        fence: &StoreMutationFence,
263        trusted_now_unix_ms: u64,
264    ) -> Result<FrostEpochCheckpointV1, FrostStoreError> {
265        validate_trusted_time(trusted_now_unix_ms)?;
266        if let Some(checkpoint) = self.persisted_rotation_checkpoint(staged.rotation_id())? {
267            verify_advance_checkpoint(&checkpoint, staged.advance(), artifact_trust)?;
268            return Ok(checkpoint);
269        }
270        let checkpoint = match anchor.compare_and_swap_epoch(
271            &staged.advance().predecessor().checkpoint_digest,
272            staged.advance(),
273        ) {
274            Ok(checkpoint) => checkpoint,
275            Err(error) => {
276                let reconciled = anchor
277                    .resolve_epoch_checkpoint(&staged.advance().target_roster().scope_id)
278                    .map_err(|reconcile_error| {
279                        FrostStoreError::Unavailable(format!(
280                            "FROST epoch CAS failed ({error}); reconciliation failed ({reconcile_error})"
281                        ))
282                    })?;
283                if verify_advance_checkpoint(&reconciled, staged.advance(), artifact_trust).is_err()
284                {
285                    return Err(FrostStoreError::Unavailable(format!(
286                        "FROST epoch CAS failed without an exact anchored successor: {error}"
287                    )));
288                }
289                reconciled
290            }
291        };
292        verify_advance_checkpoint(&checkpoint, staged.advance(), artifact_trust)?;
293        self.persist_anchor_advanced(
294            staged.rotation_id(),
295            &checkpoint,
296            fence,
297            trusted_now_unix_ms,
298        )?;
299        Ok(checkpoint)
300    }
301
302    pub fn activate_rotation(
303        &self,
304        rotation_id: &str,
305        anchor: &dyn FrostEpochAnchor,
306        artifact_trust: &FrostArtifactTrustStore,
307        fence: &StoreMutationFence,
308        trusted_now_unix_ms: u64,
309    ) -> Result<FrostRotationRecord, FrostStoreError> {
310        validate_trusted_time(trusted_now_unix_ms)?;
311        let stored = self
312            .load_stored_rotation(rotation_id)?
313            .ok_or(FrostStoreError::Conflict("rotation stage is absent"))?;
314        let checkpoint = anchor
315            .resolve_epoch_checkpoint(&stored.scope_id)
316            .map_err(anchor_error)?;
317        verify_stored_successor(&stored, &checkpoint, artifact_trust)?;
318        self.activate_with_checkpoint(rotation_id, &checkpoint, fence, trusted_now_unix_ms)
319    }
320
321    pub fn recover_rotation(
322        &self,
323        scope_id: &str,
324        anchor: &dyn FrostEpochAnchor,
325        artifact_trust: &FrostArtifactTrustStore,
326        fence: &StoreMutationFence,
327        trusted_now_unix_ms: u64,
328    ) -> Result<Option<FrostRotationRecord>, FrostStoreError> {
329        validate_trusted_time(trusted_now_unix_ms)?;
330        let stored = self.load_live_rotation(scope_id)?;
331        let Some(stored) = stored else {
332            return Ok(None);
333        };
334        let checkpoint = anchor
335            .resolve_epoch_checkpoint(scope_id)
336            .map_err(anchor_error)?;
337        artifact_trust
338            .verify_epoch_checkpoint(&checkpoint)
339            .map_err(trust_error)?;
340        if checkpoint.checkpoint_digest == stored.predecessor_checkpoint_digest {
341            if stored.state != FrostRotationState::Staged {
342                return Err(FrostStoreError::Conflict(
343                    "anchor regressed behind an acknowledged rotation",
344                ));
345            }
346            return self
347                .discard_rotation(&stored.rotation_id, fence, trusted_now_unix_ms)
348                .map(Some);
349        }
350        verify_stored_successor(&stored, &checkpoint, artifact_trust)?;
351        if stored.state == FrostRotationState::Staged {
352            self.persist_anchor_advanced(
353                &stored.rotation_id,
354                &checkpoint,
355                fence,
356                trusted_now_unix_ms,
357            )?;
358        }
359        self.activate_with_checkpoint(&stored.rotation_id, &checkpoint, fence, trusted_now_unix_ms)
360            .map(Some)
361    }
362
363    pub fn load_rotation(
364        &self,
365        rotation_id: &str,
366    ) -> Result<Option<FrostRotationRecord>, FrostStoreError> {
367        self.load_stored_rotation(rotation_id)
368            .map(|stored| stored.map(|value| public_rotation_record(&value)))
369    }
370
371    fn persisted_rotation_checkpoint(
372        &self,
373        rotation_id: &str,
374    ) -> Result<Option<FrostEpochCheckpointV1>, FrostStoreError> {
375        let stored = self.load_stored_rotation(rotation_id)?;
376        let Some(stored) = stored else {
377            return Err(FrostStoreError::Conflict("rotation stage is absent"));
378        };
379        if stored.state == FrostRotationState::Discarded {
380            return Err(FrostStoreError::Conflict("rotation stage was discarded"));
381        }
382        stored
383            .anchored_checkpoint_json
384            .map(|json| serde_json::from_slice(&json).map_err(canonical_error))
385            .transpose()
386    }
387
388    fn persist_anchor_advanced(
389        &self,
390        rotation_id: &str,
391        checkpoint: &FrostEpochCheckpointV1,
392        fence: &StoreMutationFence,
393        trusted_now_unix_ms: u64,
394    ) -> Result<(), FrostStoreError> {
395        let checkpoint_json = canonical_json_bytes(checkpoint).map_err(canonical_error)?;
396        let mut connection = self.connection()?;
397        let transaction = self.begin_write(&mut connection, fence)?;
398        let mut stored = load_rotation_query(&transaction, rotation_id)?
399            .ok_or(FrostStoreError::Conflict("rotation stage is absent"))?;
400        if stored.state == FrostRotationState::AnchorAdvanced
401            || stored.state == FrostRotationState::Active
402        {
403            if stored.anchored_checkpoint_digest.as_deref()
404                != Some(checkpoint.checkpoint_digest.as_str())
405            {
406                return Err(FrostStoreError::Conflict(
407                    "rotation retained another anchored checkpoint",
408                ));
409            }
410            transaction.commit().map_err(super::sqlite_error)?;
411            return Ok(());
412        }
413        if stored.state != FrostRotationState::Staged {
414            return Err(FrostStoreError::Conflict("rotation stage is not live"));
415        }
416        stored.state = FrostRotationState::AnchorAdvanced;
417        stored.state_version = 2;
418        stored.anchored_checkpoint_digest = Some(checkpoint.checkpoint_digest.clone());
419        stored.anchored_checkpoint_json = Some(checkpoint_json);
420        stored.source_fence = fence.clone();
421        stored.updated_at_unix_ms = trusted_now_unix_ms;
422        stored.record_digest = rotation_record_digest(&stored)?;
423        update_rotation(&transaction, &stored, FrostRotationState::Staged)?;
424        append_projection_commit(
425            &transaction,
426            self,
427            &rotation_projection_key(rotation_id),
428            ProjectionMutation {
429                sequence: 2,
430                projection_type: "rotation",
431                mutation_kind: "frost.rotation.anchor_advanced",
432                record_digest: &stored.record_digest,
433            },
434            fence,
435        )?;
436        self.commit_write(transaction)?;
437        self.sync_after_write(&connection)
438    }
439
440    fn activate_with_checkpoint(
441        &self,
442        rotation_id: &str,
443        checkpoint: &FrostEpochCheckpointV1,
444        fence: &StoreMutationFence,
445        trusted_now_unix_ms: u64,
446    ) -> Result<FrostRotationRecord, FrostStoreError> {
447        let checkpoint_json = canonical_json_bytes(checkpoint).map_err(canonical_error)?;
448        let mut connection = self.connection()?;
449        let transaction = self.begin_write(&mut connection, fence)?;
450        let mut stored = load_rotation_query(&transaction, rotation_id)?
451            .ok_or(FrostStoreError::Conflict("rotation stage is absent"))?;
452        if stored.state == FrostRotationState::Active {
453            if stored.anchored_checkpoint_digest.as_deref()
454                != Some(checkpoint.checkpoint_digest.as_str())
455            {
456                return Err(FrostStoreError::Conflict(
457                    "active rotation retained another checkpoint",
458                ));
459            }
460            transaction.commit().map_err(super::sqlite_error)?;
461            return Ok(public_rotation_record(&stored));
462        }
463        if stored.state != FrostRotationState::AnchorAdvanced {
464            return Err(FrostStoreError::Conflict(
465                "rotation anchor has not been durably acknowledged",
466            ));
467        }
468        let target: FrostRosterV1 =
469            serde_json::from_slice(&stored.target_roster_json).map_err(canonical_error)?;
470        let burn: FrostSessionBurnSummaryV1 =
471            serde_json::from_slice(&stored.burn_summary_json).map_err(canonical_error)?;
472        verify_local_burn_summary(&transaction, &burn)?;
473        let predecessor_key_epoch =
474            target
475                .key_epoch
476                .checked_sub(1)
477                .ok_or(FrostStoreError::Conflict(
478                    "target roster key epoch has no predecessor",
479                ))?;
480        insert_roster_history(
481            &transaction,
482            &target,
483            &stored.target_roster_json,
484            checkpoint,
485            &checkpoint_json,
486            fence,
487            trusted_now_unix_ms,
488        )?;
489        let changed = transaction
490            .execute(
491                r#"
492                UPDATE frost_active_rosters
493                SET key_epoch = ?1, roster_digest = ?2,
494                    checkpoint_sequence = ?3, checkpoint_digest = ?4,
495                    activation_fence = ?5, clock_high_water = ?6,
496                    source_store_uuid = ?7, source_lease_id = ?8,
497                    source_owner_epoch = ?9
498                WHERE scope_id = ?10 AND key_epoch = ?11
499                  AND checkpoint_digest = ?12
500                "#,
501                params![
502                    sqlite_u64(target.key_epoch, "target key epoch")?,
503                    &target.roster_digest,
504                    sqlite_u64(checkpoint.checkpoint_sequence, "checkpoint sequence")?,
505                    &checkpoint.checkpoint_digest,
506                    sqlite_u64(checkpoint.activation_fence, "activation fence")?,
507                    sqlite_u64(checkpoint.clock_high_water, "clock high-water")?,
508                    &fence.store_uuid,
509                    &fence.lease_id,
510                    sqlite_u64(fence.owner_epoch, "source owner epoch")?,
511                    &target.scope_id,
512                    sqlite_u64(predecessor_key_epoch, "predecessor key epoch")?,
513                    &stored.predecessor_checkpoint_digest,
514                ],
515            )
516            .map_err(super::sqlite_error)?;
517        if changed != 1 {
518            return Err(FrostStoreError::Conflict(
519                "active roster changed after rotation staging",
520            ));
521        }
522        stored.state = FrostRotationState::Active;
523        stored.state_version = 3;
524        stored.anchored_checkpoint_digest = Some(checkpoint.checkpoint_digest.clone());
525        stored.anchored_checkpoint_json = Some(checkpoint_json);
526        stored.source_fence = fence.clone();
527        stored.updated_at_unix_ms = trusted_now_unix_ms;
528        stored.record_digest = rotation_record_digest(&stored)?;
529        update_rotation(&transaction, &stored, FrostRotationState::AnchorAdvanced)?;
530        append_projection_commit(
531            &transaction,
532            self,
533            &rotation_projection_key(rotation_id),
534            ProjectionMutation {
535                sequence: 3,
536                projection_type: "rotation",
537                mutation_kind: "frost.rotation.activate",
538                record_digest: &stored.record_digest,
539            },
540            fence,
541        )?;
542        self.commit_write(transaction)?;
543        self.sync_after_write(&connection)?;
544        Ok(public_rotation_record(&stored))
545    }
546
547    fn discard_rotation(
548        &self,
549        rotation_id: &str,
550        fence: &StoreMutationFence,
551        trusted_now_unix_ms: u64,
552    ) -> Result<FrostRotationRecord, FrostStoreError> {
553        let mut connection = self.connection()?;
554        let transaction = self.begin_write(&mut connection, fence)?;
555        let mut stored = load_rotation_query(&transaction, rotation_id)?
556            .ok_or(FrostStoreError::Conflict("rotation stage is absent"))?;
557        if stored.state == FrostRotationState::Discarded {
558            transaction.commit().map_err(super::sqlite_error)?;
559            return Ok(public_rotation_record(&stored));
560        }
561        if stored.state != FrostRotationState::Staged {
562            return Err(FrostStoreError::Conflict(
563                "only an unanchored stage may be discarded",
564            ));
565        }
566        stored.state = FrostRotationState::Discarded;
567        stored.state_version = 2;
568        stored.source_fence = fence.clone();
569        stored.updated_at_unix_ms = trusted_now_unix_ms;
570        stored.record_digest = rotation_record_digest(&stored)?;
571        update_rotation(&transaction, &stored, FrostRotationState::Staged)?;
572        append_projection_commit(
573            &transaction,
574            self,
575            &rotation_projection_key(rotation_id),
576            ProjectionMutation {
577                sequence: 2,
578                projection_type: "rotation",
579                mutation_kind: "frost.rotation.discard",
580                record_digest: &stored.record_digest,
581            },
582            fence,
583        )?;
584        self.commit_write(transaction)?;
585        self.sync_after_write(&connection)?;
586        Ok(public_rotation_record(&stored))
587    }
588
589    fn load_stored_rotation(
590        &self,
591        rotation_id: &str,
592    ) -> Result<Option<StoredRotation>, FrostStoreError> {
593        let mut connection = self.connection()?;
594        let transaction = self.begin_read(&mut connection, None)?;
595        let stored = load_rotation_query(&transaction, rotation_id)?;
596        transaction.commit().map_err(super::sqlite_error)?;
597        Ok(stored)
598    }
599
600    fn load_live_rotation(
601        &self,
602        scope_id: &str,
603    ) -> Result<Option<StoredRotation>, FrostStoreError> {
604        let mut connection = self.connection()?;
605        let transaction = self.begin_read(&mut connection, None)?;
606        let rotation_id = transaction
607            .query_row(
608                r#"
609                SELECT rotation_id FROM frost_roster_rotations
610                WHERE scope_id = ?1 AND state IN ('staged', 'anchor_advanced')
611                "#,
612                [scope_id],
613                |row| row.get::<_, String>(0),
614            )
615            .optional()
616            .map_err(super::sqlite_error)?;
617        let stored = rotation_id
618            .as_deref()
619            .map(|id| load_rotation_query(&transaction, id))
620            .transpose()?
621            .flatten();
622        transaction.commit().map_err(super::sqlite_error)?;
623        Ok(stored)
624    }
625}
626
627pub(super) fn verify_rotation_invariants(connection: &Connection) -> Result<(), FrostStoreError> {
628    let mut statement = connection
629        .prepare("SELECT rotation_id FROM frost_roster_rotations ORDER BY rotation_id")
630        .map_err(super::sqlite_error)?;
631    let ids = statement
632        .query_map([], |row| row.get::<_, String>(0))
633        .map_err(super::sqlite_error)?
634        .collect::<Result<Vec<_>, _>>()
635        .map_err(super::sqlite_error)?;
636    drop(statement);
637    for id in ids {
638        let stored = load_rotation_query(connection, &id)?
639            .ok_or_else(|| invalid("rotation disappeared during verification"))?;
640        if rotation_record_digest(&stored)? != stored.record_digest {
641            return Err(invalid("FROST rotation record digest is invalid"));
642        }
643        let latest = connection
644            .query_row(
645                r#"
646                SELECT projection_sequence, record_digest
647                FROM frost_projection_commits
648                WHERE projection_key = ?1
649                ORDER BY projection_sequence DESC LIMIT 1
650                "#,
651                [rotation_projection_key(&id)],
652                |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)),
653            )
654            .optional()
655            .map_err(super::sqlite_error)?
656            .ok_or_else(|| invalid("rotation has no projection commit"))?;
657        if read_u64(latest.0, "rotation projection sequence")? != stored.state_version
658            || latest.1 != stored.record_digest
659        {
660            return Err(invalid("rotation does not match its projection head"));
661        }
662    }
663    verify_active_roster_invariants(connection)
664}
665
666fn verify_active_roster_invariants(connection: &Connection) -> Result<(), FrostStoreError> {
667    let mut statement = connection
668        .prepare("SELECT scope_id FROM frost_active_rosters ORDER BY scope_id")
669        .map_err(super::sqlite_error)?;
670    let scopes = statement
671        .query_map([], |row| row.get::<_, String>(0))
672        .map_err(super::sqlite_error)?
673        .collect::<Result<Vec<_>, _>>()
674        .map_err(super::sqlite_error)?;
675    drop(statement);
676    for scope in scopes {
677        let (roster, checkpoint) = load_active_artifacts(connection, &scope)?
678            .ok_or_else(|| invalid("active roster history is absent"))?;
679        let active = load_active_record(connection, &scope)?
680            .ok_or_else(|| invalid("active roster pointer disappeared"))?;
681        if active != active_record(&roster, &checkpoint) {
682            return Err(invalid("active roster pointer diverges from history"));
683        }
684    }
685    Ok(())
686}
687
688fn insert_roster_history(
689    transaction: &Transaction<'_>,
690    roster: &FrostRosterV1,
691    roster_json: &[u8],
692    checkpoint: &FrostEpochCheckpointV1,
693    checkpoint_json: &[u8],
694    fence: &StoreMutationFence,
695    activated_at_unix_ms: u64,
696) -> Result<(), FrostStoreError> {
697    let digest = roster_history_record_digest(
698        roster,
699        roster_json,
700        checkpoint,
701        checkpoint_json,
702        fence,
703        activated_at_unix_ms,
704    )?;
705    transaction
706        .execute(
707            r#"
708            INSERT INTO frost_roster_history (
709                scope_id, key_epoch, roster_id, roster_digest, roster_json,
710                checkpoint_sequence, checkpoint_digest, checkpoint_json,
711                activation_fence, clock_high_water, activated_at_unix_ms,
712                source_store_uuid, source_lease_id, source_owner_epoch,
713                record_digest
714            ) VALUES (
715                ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10,
716                ?11, ?12, ?13, ?14, ?15
717            )
718            "#,
719            params![
720                &roster.scope_id,
721                sqlite_u64(roster.key_epoch, "roster key epoch")?,
722                &roster.roster_id,
723                &roster.roster_digest,
724                roster_json,
725                sqlite_u64(checkpoint.checkpoint_sequence, "checkpoint sequence")?,
726                &checkpoint.checkpoint_digest,
727                checkpoint_json,
728                sqlite_u64(checkpoint.activation_fence, "activation fence")?,
729                sqlite_u64(checkpoint.clock_high_water, "clock high-water")?,
730                sqlite_u64(activated_at_unix_ms, "activation time")?,
731                &fence.store_uuid,
732                &fence.lease_id,
733                sqlite_u64(fence.owner_epoch, "source owner epoch")?,
734                digest,
735            ],
736        )
737        .map_err(super::sqlite_error)?;
738    Ok(())
739}
740
741fn insert_active_pointer(
742    transaction: &Transaction<'_>,
743    roster: &FrostRosterV1,
744    checkpoint: &FrostEpochCheckpointV1,
745    fence: &StoreMutationFence,
746) -> Result<(), FrostStoreError> {
747    transaction
748        .execute(
749            r#"
750            INSERT INTO frost_active_rosters (
751                scope_id, key_epoch, roster_digest, checkpoint_sequence,
752                checkpoint_digest, activation_fence, clock_high_water,
753                source_store_uuid, source_lease_id, source_owner_epoch
754            ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)
755            "#,
756            params![
757                &roster.scope_id,
758                sqlite_u64(roster.key_epoch, "roster key epoch")?,
759                &roster.roster_digest,
760                sqlite_u64(checkpoint.checkpoint_sequence, "checkpoint sequence")?,
761                &checkpoint.checkpoint_digest,
762                sqlite_u64(checkpoint.activation_fence, "activation fence")?,
763                sqlite_u64(checkpoint.clock_high_water, "clock high-water")?,
764                &fence.store_uuid,
765                &fence.lease_id,
766                sqlite_u64(fence.owner_epoch, "source owner epoch")?,
767            ],
768        )
769        .map_err(super::sqlite_error)?;
770    Ok(())
771}
772
773fn insert_rotation(
774    transaction: &Transaction<'_>,
775    stored: &StoredRotation,
776) -> Result<(), FrostStoreError> {
777    transaction
778        .execute(
779            r#"
780            INSERT INTO frost_roster_rotations (
781                rotation_id, scope_id, state, state_version,
782                predecessor_checkpoint_digest, predecessor_checkpoint_json,
783                target_roster_digest, target_roster_json, target_key_epoch,
784                rotation_authorization_digest, burn_summary_json,
785                old_session_burn_root, expected_checkpoint_sequence,
786                activation_fence, clock_high_water,
787                anchored_checkpoint_digest, anchored_checkpoint_json,
788                source_store_uuid, source_lease_id, source_owner_epoch,
789                created_at_unix_ms, updated_at_unix_ms, record_digest
790            ) VALUES (
791                ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12,
792                ?13, ?14, ?15, ?16, ?17, ?18, ?19, ?20, ?21, ?22, ?23
793            )
794            "#,
795            rusqlite::params_from_iter(rotation_params(stored)?),
796        )
797        .map_err(super::sqlite_error)?;
798    Ok(())
799}
800
801fn update_rotation(
802    transaction: &Transaction<'_>,
803    stored: &StoredRotation,
804    expected_state: FrostRotationState,
805) -> Result<(), FrostStoreError> {
806    let changed = transaction
807        .execute(
808            r#"
809            UPDATE frost_roster_rotations
810            SET state = ?1, state_version = ?2,
811                anchored_checkpoint_digest = ?3, anchored_checkpoint_json = ?4,
812                source_store_uuid = ?5, source_lease_id = ?6,
813                source_owner_epoch = ?7, updated_at_unix_ms = ?8,
814                record_digest = ?9
815            WHERE rotation_id = ?10 AND state = ?11
816            "#,
817            params![
818                stored.state.as_str(),
819                sqlite_u64(stored.state_version, "rotation state version")?,
820                stored.anchored_checkpoint_digest.as_deref(),
821                stored.anchored_checkpoint_json.as_deref(),
822                &stored.source_fence.store_uuid,
823                &stored.source_fence.lease_id,
824                sqlite_u64(stored.source_fence.owner_epoch, "source owner epoch")?,
825                sqlite_u64(stored.updated_at_unix_ms, "rotation updated time")?,
826                &stored.record_digest,
827                &stored.rotation_id,
828                expected_state.as_str(),
829            ],
830        )
831        .map_err(super::sqlite_error)?;
832    if changed != 1 {
833        return Err(FrostStoreError::Conflict(
834            "rotation state changed concurrently",
835        ));
836    }
837    Ok(())
838}
839
840fn rotation_params(
841    stored: &StoredRotation,
842) -> Result<Vec<rusqlite::types::Value>, FrostStoreError> {
843    use rusqlite::types::Value;
844    Ok(vec![
845        Value::Text(stored.rotation_id.clone()),
846        Value::Text(stored.scope_id.clone()),
847        Value::Text(stored.state.as_str().to_string()),
848        Value::Integer(sqlite_u64(stored.state_version, "rotation state version")?),
849        Value::Text(stored.predecessor_checkpoint_digest.clone()),
850        Value::Blob(stored.predecessor_checkpoint_json.clone()),
851        Value::Text(stored.target_roster_digest.clone()),
852        Value::Blob(stored.target_roster_json.clone()),
853        Value::Integer(sqlite_u64(stored.target_key_epoch, "target key epoch")?),
854        Value::Text(stored.rotation_authorization_digest.clone()),
855        Value::Blob(stored.burn_summary_json.clone()),
856        Value::Text(stored.old_session_burn_root.clone()),
857        Value::Integer(sqlite_u64(
858            stored.expected_checkpoint_sequence,
859            "expected checkpoint sequence",
860        )?),
861        Value::Integer(sqlite_u64(stored.activation_fence, "activation fence")?),
862        Value::Integer(sqlite_u64(stored.clock_high_water, "clock high-water")?),
863        stored
864            .anchored_checkpoint_digest
865            .clone()
866            .map_or(Value::Null, Value::Text),
867        stored
868            .anchored_checkpoint_json
869            .clone()
870            .map_or(Value::Null, Value::Blob),
871        Value::Text(stored.source_fence.store_uuid.clone()),
872        Value::Text(stored.source_fence.lease_id.clone()),
873        Value::Integer(sqlite_u64(
874            stored.source_fence.owner_epoch,
875            "source owner epoch",
876        )?),
877        Value::Integer(sqlite_u64(
878            stored.created_at_unix_ms,
879            "rotation created time",
880        )?),
881        Value::Integer(sqlite_u64(
882            stored.updated_at_unix_ms,
883            "rotation updated time",
884        )?),
885        Value::Text(stored.record_digest.clone()),
886    ])
887}
888
889fn load_rotation_query(
890    connection: &Connection,
891    rotation_id: &str,
892) -> Result<Option<StoredRotation>, FrostStoreError> {
893    connection
894        .query_row(
895            r#"
896            SELECT rotation_id, scope_id, state, state_version,
897                   predecessor_checkpoint_digest, predecessor_checkpoint_json,
898                   target_roster_digest, target_roster_json, target_key_epoch,
899                   rotation_authorization_digest, burn_summary_json,
900                   old_session_burn_root, expected_checkpoint_sequence,
901                   activation_fence, clock_high_water,
902                   anchored_checkpoint_digest, anchored_checkpoint_json,
903                   source_store_uuid, source_lease_id, source_owner_epoch,
904                   created_at_unix_ms, updated_at_unix_ms, record_digest
905            FROM frost_roster_rotations WHERE rotation_id = ?1
906            "#,
907            [rotation_id],
908            read_rotation_row,
909        )
910        .optional()
911        .map_err(super::sqlite_error)
912}
913
914fn read_rotation_row(row: &Row<'_>) -> Result<StoredRotation, rusqlite::Error> {
915    let state_value: String = row.get(2)?;
916    let state = FrostRotationState::parse(&state_value).map_err(|error| {
917        rusqlite::Error::FromSqlConversionFailure(2, rusqlite::types::Type::Text, error.into())
918    })?;
919    Ok(StoredRotation {
920        rotation_id: row.get(0)?,
921        scope_id: row.get(1)?,
922        state,
923        state_version: sqlite_read_u64(row, 3)?,
924        predecessor_checkpoint_digest: row.get(4)?,
925        predecessor_checkpoint_json: row.get(5)?,
926        target_roster_digest: row.get(6)?,
927        target_roster_json: row.get(7)?,
928        target_key_epoch: sqlite_read_u64(row, 8)?,
929        rotation_authorization_digest: row.get(9)?,
930        burn_summary_json: row.get(10)?,
931        old_session_burn_root: row.get(11)?,
932        expected_checkpoint_sequence: sqlite_read_u64(row, 12)?,
933        activation_fence: sqlite_read_u64(row, 13)?,
934        clock_high_water: sqlite_read_u64(row, 14)?,
935        anchored_checkpoint_digest: row.get(15)?,
936        anchored_checkpoint_json: row.get(16)?,
937        source_fence: StoreMutationFence {
938            store_uuid: row.get(17)?,
939            lease_id: row.get(18)?,
940            owner_epoch: sqlite_read_u64(row, 19)?,
941        },
942        created_at_unix_ms: sqlite_read_u64(row, 20)?,
943        updated_at_unix_ms: sqlite_read_u64(row, 21)?,
944        record_digest: row.get(22)?,
945    })
946}
947
948pub(super) fn load_active_record(
949    connection: &Connection,
950    scope_id: &str,
951) -> Result<Option<FrostActiveRosterRecord>, FrostStoreError> {
952    connection
953        .query_row(
954            r#"
955            SELECT scope_id, key_epoch, roster_digest, checkpoint_sequence,
956                   checkpoint_digest, activation_fence, clock_high_water
957            FROM frost_active_rosters WHERE scope_id = ?1
958            "#,
959            [scope_id],
960            |row| {
961                Ok(FrostActiveRosterRecord {
962                    scope_id: row.get(0)?,
963                    key_epoch: sqlite_read_u64(row, 1)?,
964                    roster_digest: row.get(2)?,
965                    checkpoint_sequence: sqlite_read_u64(row, 3)?,
966                    checkpoint_digest: row.get(4)?,
967                    activation_fence: sqlite_read_u64(row, 5)?,
968                    clock_high_water: sqlite_read_u64(row, 6)?,
969                })
970            },
971        )
972        .optional()
973        .map_err(super::sqlite_error)
974}
975
976fn load_active_artifacts(
977    connection: &Connection,
978    scope_id: &str,
979) -> Result<Option<(FrostRosterV1, FrostEpochCheckpointV1)>, FrostStoreError> {
980    let values = connection
981        .query_row(
982            r#"
983            SELECT history.roster_json, history.checkpoint_json
984            FROM frost_active_rosters AS active
985            JOIN frost_roster_history AS history
986              ON history.scope_id = active.scope_id
987             AND history.key_epoch = active.key_epoch
988             AND history.roster_digest = active.roster_digest
989             AND history.checkpoint_digest = active.checkpoint_digest
990            WHERE active.scope_id = ?1
991            "#,
992            [scope_id],
993            |row| Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?)),
994        )
995        .optional()
996        .map_err(super::sqlite_error)?;
997    values
998        .map(|(roster, checkpoint)| {
999            Ok((
1000                serde_json::from_slice(&roster).map_err(canonical_error)?,
1001                serde_json::from_slice(&checkpoint).map_err(canonical_error)?,
1002            ))
1003        })
1004        .transpose()
1005}
1006
1007fn rotation_id(advance: &VerifiedFrostEpochAdvance) -> Result<String, FrostStoreError> {
1008    prefixed_digest(
1009        ROTATION_ID_PREFIX,
1010        &RotationIdPreimage {
1011            predecessor_checkpoint_digest: &advance.predecessor().checkpoint_digest,
1012            target_roster_digest: &advance.target_roster().roster_digest,
1013            rotation_authorization_digest: advance.rotation_authorization_digest(),
1014            old_session_burn_root: &advance.burn_summary().burn_root,
1015            activation_fence: advance.activation_fence(),
1016            clock_high_water: advance.clock_high_water(),
1017        },
1018    )
1019}
1020
1021fn rotation_record_digest(stored: &StoredRotation) -> Result<String, FrostStoreError> {
1022    prefixed_digest(
1023        ROTATION_RECORD_PREFIX,
1024        &RotationRecordPreimage {
1025            rotation_id: &stored.rotation_id,
1026            scope_id: &stored.scope_id,
1027            state: stored.state,
1028            state_version: stored.state_version,
1029            predecessor_checkpoint_digest: &stored.predecessor_checkpoint_digest,
1030            predecessor_checkpoint_json_digest: sha256_hex(&stored.predecessor_checkpoint_json),
1031            target_roster_digest: &stored.target_roster_digest,
1032            target_roster_json_digest: sha256_hex(&stored.target_roster_json),
1033            target_key_epoch: stored.target_key_epoch,
1034            rotation_authorization_digest: &stored.rotation_authorization_digest,
1035            burn_summary_json_digest: sha256_hex(&stored.burn_summary_json),
1036            old_session_burn_root: &stored.old_session_burn_root,
1037            expected_checkpoint_sequence: stored.expected_checkpoint_sequence,
1038            activation_fence: stored.activation_fence,
1039            clock_high_water: stored.clock_high_water,
1040            anchored_checkpoint_digest: stored.anchored_checkpoint_digest.as_deref(),
1041            anchored_checkpoint_json_digest: stored
1042                .anchored_checkpoint_json
1043                .as_deref()
1044                .map(sha256_hex),
1045            source_store_uuid: &stored.source_fence.store_uuid,
1046            source_lease_id: &stored.source_fence.lease_id,
1047            source_owner_epoch: stored.source_fence.owner_epoch,
1048            created_at_unix_ms: stored.created_at_unix_ms,
1049            updated_at_unix_ms: stored.updated_at_unix_ms,
1050        },
1051    )
1052}
1053
1054fn roster_history_record_digest(
1055    roster: &FrostRosterV1,
1056    roster_json: &[u8],
1057    checkpoint: &FrostEpochCheckpointV1,
1058    checkpoint_json: &[u8],
1059    fence: &StoreMutationFence,
1060    activated_at_unix_ms: u64,
1061) -> Result<String, FrostStoreError> {
1062    prefixed_digest(
1063        ROSTER_HISTORY_RECORD_PREFIX,
1064        &RosterHistoryRecordPreimage {
1065            scope_id: &roster.scope_id,
1066            key_epoch: roster.key_epoch,
1067            roster_id: &roster.roster_id,
1068            roster_digest: &roster.roster_digest,
1069            roster_json_digest: sha256_hex(roster_json),
1070            checkpoint_sequence: checkpoint.checkpoint_sequence,
1071            checkpoint_digest: &checkpoint.checkpoint_digest,
1072            checkpoint_json_digest: sha256_hex(checkpoint_json),
1073            activation_fence: checkpoint.activation_fence,
1074            clock_high_water: checkpoint.clock_high_water,
1075            activated_at_unix_ms,
1076            source_store_uuid: &fence.store_uuid,
1077            source_lease_id: &fence.lease_id,
1078            source_owner_epoch: fence.owner_epoch,
1079        },
1080    )
1081}
1082
1083fn active_record(
1084    roster: &FrostRosterV1,
1085    checkpoint: &FrostEpochCheckpointV1,
1086) -> FrostActiveRosterRecord {
1087    FrostActiveRosterRecord {
1088        scope_id: roster.scope_id.clone(),
1089        key_epoch: roster.key_epoch,
1090        roster_digest: roster.roster_digest.clone(),
1091        checkpoint_sequence: checkpoint.checkpoint_sequence,
1092        checkpoint_digest: checkpoint.checkpoint_digest.clone(),
1093        activation_fence: checkpoint.activation_fence,
1094        clock_high_water: checkpoint.clock_high_water,
1095    }
1096}
1097
1098fn public_rotation_record(stored: &StoredRotation) -> FrostRotationRecord {
1099    FrostRotationRecord {
1100        rotation_id: stored.rotation_id.clone(),
1101        scope_id: stored.scope_id.clone(),
1102        state: stored.state,
1103        state_version: stored.state_version,
1104        predecessor_checkpoint_digest: stored.predecessor_checkpoint_digest.clone(),
1105        target_roster_digest: stored.target_roster_digest.clone(),
1106        target_key_epoch: stored.target_key_epoch,
1107        anchored_checkpoint_digest: stored.anchored_checkpoint_digest.clone(),
1108    }
1109}
1110
1111fn roster_projection_key(scope_id: &str, key_epoch: u64) -> String {
1112    format!("roster/{scope_id}/{key_epoch}")
1113}
1114
1115fn rotation_projection_key(rotation_id: &str) -> String {
1116    format!("rotation/{rotation_id}")
1117}
1118
1119fn validate_trusted_time(value: u64) -> Result<(), FrostStoreError> {
1120    if value == 0 || i64::try_from(value).is_err() {
1121        return Err(FrostStoreError::Conflict(
1122            "trusted time is outside the SQLite range",
1123        ));
1124    }
1125    Ok(())
1126}
1127
1128fn sqlite_read_u64(row: &Row<'_>, index: usize) -> Result<u64, rusqlite::Error> {
1129    let value = row.get::<_, i64>(index)?;
1130    u64::try_from(value).map_err(|_| rusqlite::Error::IntegralValueOutOfRange(index, value))
1131}
1132
1133pub(super) fn sqlite_u64(value: u64, field: &'static str) -> Result<i64, FrostStoreError> {
1134    i64::try_from(value)
1135        .map_err(|_| FrostStoreError::InvalidState(format!("{field} exceeds SQLite range")))
1136}
1137
1138fn read_u64(value: i64, field: &'static str) -> Result<u64, FrostStoreError> {
1139    u64::try_from(value).map_err(|_| FrostStoreError::InvalidState(format!("{field} is negative")))
1140}
1141
1142pub(super) fn canonical_error(error: impl std::fmt::Display) -> FrostStoreError {
1143    FrostStoreError::InvalidState(error.to_string())
1144}
1145
1146pub(super) fn trust_error(error: impl std::fmt::Display) -> FrostStoreError {
1147    FrostStoreError::InvalidState(format!("FROST artifact trust failed: {error}"))
1148}
1149
1150fn anchor_error(error: impl std::fmt::Display) -> FrostStoreError {
1151    FrostStoreError::Unavailable(format!("FROST epoch anchor failed: {error}"))
1152}
1153
1154pub(super) fn invalid(detail: impl Into<String>) -> FrostStoreError {
1155    FrostStoreError::InvalidState(detail.into())
1156}