1use std::collections::HashSet;
2use std::sync::{Arc, Mutex, MutexGuard};
3
4use chio_core::canonical::canonical_json_bytes;
5use chio_core::{sha256_hex, Keypair, StoreMutationFence};
6use chio_fiscal::{
7 fee_schedule::SignedOpenMarketFeeSchedule, FiscalActivationTarget, FiscalAdmissionAuthority,
8 FiscalAdmissionTrustRegistry, FiscalAuthorityState, FiscalCharterRegistry, FiscalDomain,
9 FiscalGenesisPolicy, FiscalParams, FiscalProposalAdmissionBuilder,
10 FiscalProposalAdmissionState, FiscalProposalAdmissionStatus, FiscalRuntimeAdapterRegistry,
11 FiscalScheduleHead, FiscalStagedTransition, SignedFiscalActivation, SignedFiscalProposal,
12 SignedFiscalSchedule, VerifiedFiscalActivation, VerifiedFiscalApproval, VerifiedFiscalCharter,
13 VerifiedFiscalContinuityAdvance, VerifiedFiscalContinuityCheckpoint, VerifiedFiscalProposal,
14 VerifiedFiscalProposalAdmission, VerifiedFiscalRuntimeReadiness, VerifiedFiscalSchedule,
15};
16use rusqlite::{params, Connection, OptionalExtension, Transaction, TransactionBehavior};
17use serde::Serialize;
18
19use crate::serving_owner::{SqliteServingOwner, SqliteServingOwnerError};
20
21const FISCAL_STORE_SCHEMA_KEY: &str = "fiscal";
22pub(crate) const FISCAL_STORE_SUPPORTED_SCHEMA_VERSION: i32 = 6;
23const FISCAL_STORE_SCHEMA: &str = include_str!("fiscal_store.sql");
24const FISCAL_GENESIS_PROJECTION_KEY: &str = "authority";
25const ZERO_DIGEST: &str = "0000000000000000000000000000000000000000000000000000000000000000";
26
27#[derive(Debug, thiserror::Error)]
28pub enum FiscalStoreError {
29 #[error("fiscal store is unavailable: {0}")]
30 Unavailable(String),
31 #[error("fiscal store mutation was fenced")]
32 Fenced,
33 #[error("fiscal store record was not found")]
34 NotFound,
35 #[error("fiscal store record conflicts with retained state")]
36 Conflict,
37 #[error("fiscal store invariant failed: {0}")]
38 Invariant(String),
39 #[error("fiscal store durable outcome is unknown: {0}")]
40 OutcomeUnknown(String),
41 #[error(transparent)]
42 Fiscal(#[from] chio_fiscal::FiscalError),
43}
44
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub enum FiscalStageStatus {
47 DbStaged,
48 FiscalAnchorAdvanced,
49 DbFinalized,
50 Discarded,
51}
52
53impl FiscalStageStatus {
54 const fn as_str(self) -> &'static str {
55 match self {
56 Self::DbStaged => "db_staged",
57 Self::FiscalAnchorAdvanced => "fiscal_anchor_advanced",
58 Self::DbFinalized => "db_finalized",
59 Self::Discarded => "discarded",
60 }
61 }
62
63 fn parse(value: &str) -> Result<Self, FiscalStoreError> {
64 match value {
65 "db_staged" => Ok(Self::DbStaged),
66 "fiscal_anchor_advanced" => Ok(Self::FiscalAnchorAdvanced),
67 "db_finalized" => Ok(Self::DbFinalized),
68 "discarded" => Ok(Self::Discarded),
69 _ => Err(invariant("fiscal stage status is invalid")),
70 }
71 }
72}
73
74#[derive(Debug, Clone, PartialEq, Eq)]
75pub struct FiscalStagedTransitionRecord {
76 pub transition_id: String,
77 pub current_checkpoint_digest: String,
78 pub next_checkpoint_digest: String,
79 pub proof_json: Vec<u8>,
80 pub status: FiscalStageStatus,
81 pub stage_version: u64,
82}
83
84#[derive(Debug, Clone, PartialEq, Eq)]
85pub struct FiscalLegacyFeeScheduleBindingRecord {
86 pub legacy_schedule_id: String,
87 pub fiscal_schedule_id: String,
88 pub fiscal_schedule_digest: String,
89 pub legacy_envelope_digest: String,
90}
91
92#[derive(Serialize)]
93#[serde(rename_all = "camelCase")]
94struct FiscalProjectionCommit<'a> {
95 format: &'static str,
96 projection_key: &'a str,
97 projection_sequence: u64,
98 mutation_kind: &'a str,
99 snapshot_digest: &'a str,
100 previous_commit_digest: &'a str,
101 store_uuid: &'a str,
102 store_lease_id: &'a str,
103 store_owner_epoch: u64,
104}
105
106#[derive(Serialize)]
107#[serde(rename_all = "camelCase")]
108struct FiscalGenesisSnapshot<'a> {
109 policy: &'a FiscalGenesisPolicy,
110 authority: &'a FiscalAuthorityState,
111 charter_digest: &'a str,
112 readiness_digest: &'a str,
113 checkpoint_digest: &'a str,
114}
115
116struct PreparedFiscalActivationMutation {
117 transition_id: String,
118 activation_id: String,
119 activation_digest: String,
120 admission_id: String,
121 admission_digest: String,
122 expected_admission_version: u64,
123 expected_admission_json: Vec<u8>,
124 activated_admission_version: u64,
125 activated_admission_json: Vec<u8>,
126 candidate_schedule_id: String,
127 candidate_schedule_digest: String,
128 predecessor_schedule_id: Option<String>,
129 predecessor_schedule_digest: Option<String>,
130}
131
132struct PreparedFiscalRotationScheduleMutation {
133 domain: FiscalDomain,
134 candidate_schedule_id: String,
135 candidate_schedule_digest: String,
136 predecessor_schedule_id: String,
137 predecessor_schedule_digest: String,
138}
139
140struct PreparedFiscalRotationMutation {
141 transition_id: String,
142 activation_id: String,
143 activation_digest: String,
144 admission_id: String,
145 admission_digest: String,
146 expected_admission_version: u64,
147 expected_admission_json: Vec<u8>,
148 activated_admission_version: u64,
149 activated_admission_json: Vec<u8>,
150 successor_charter_id: String,
151 successor_charter_digest: String,
152 predecessor_charter_id: String,
153 predecessor_charter_digest: String,
154 schedules: Vec<PreparedFiscalRotationScheduleMutation>,
155}
156
157#[derive(Clone)]
158pub struct SqliteFiscalStore {
159 connection: Arc<Mutex<Connection>>,
160 serving_owner: Arc<SqliteServingOwner>,
161}
162
163impl SqliteFiscalStore {
164 pub(crate) fn open_alongside(
165 connection: Arc<Mutex<Connection>>,
166 serving_owner: Arc<SqliteServingOwner>,
167 ) -> Self {
168 Self {
169 connection,
170 serving_owner,
171 }
172 }
173
174 fn connection(&self) -> Result<MutexGuard<'_, Connection>, FiscalStoreError> {
175 self.connection
176 .lock()
177 .map_err(|_| invariant("fiscal store lock is poisoned"))
178 }
179
180 fn begin_read<'a>(
181 &self,
182 connection: &'a mut Connection,
183 ) -> Result<Transaction<'a>, FiscalStoreError> {
184 let transaction = connection
185 .transaction_with_behavior(TransactionBehavior::Deferred)
186 .map_err(sqlite_error)?;
187 verify_owner(&transaction, &self.serving_owner, None)?;
188 self.serving_owner
189 .verify_authority_anchor(&transaction)
190 .map_err(map_owner_error)?;
191 verify_fiscal_sql_invariants(&transaction)?;
192 Ok(transaction)
193 }
194
195 fn begin_write<'a>(
196 &self,
197 connection: &'a mut Connection,
198 fence: &StoreMutationFence,
199 ) -> Result<Transaction<'a>, FiscalStoreError> {
200 let transaction = connection
201 .transaction_with_behavior(TransactionBehavior::Immediate)
202 .map_err(sqlite_error)?;
203 verify_owner(&transaction, &self.serving_owner, Some(fence))?;
204 self.serving_owner
205 .verify_authority_anchor(&transaction)
206 .map_err(map_owner_error)?;
207 verify_fiscal_sql_invariants(&transaction)?;
208 Ok(transaction)
209 }
210
211 fn commit_write(&self, transaction: Transaction<'_>) -> Result<(), FiscalStoreError> {
212 transaction.commit().map_err(|error| {
213 map_owner_error(self.serving_owner.outcome_unknown(format!(
214 "sqlite fiscal store commit outcome is unknown: {error}"
215 )))
216 })
217 }
218
219 fn sync_after_write(&self, connection: &Connection) -> Result<(), FiscalStoreError> {
220 self.serving_owner
221 .sync_authority_anchor(connection)
222 .map_err(map_owner_error)
223 }
224
225 pub fn initialize_genesis(
226 &self,
227 policy: &FiscalGenesisPolicy,
228 authority: &FiscalAuthorityState,
229 charter: &VerifiedFiscalCharter,
230 readiness: &VerifiedFiscalRuntimeReadiness,
231 checkpoint: &VerifiedFiscalContinuityCheckpoint,
232 fence: &StoreMutationFence,
233 ) -> Result<(), FiscalStoreError> {
234 policy.validate(charter)?;
235 authority.validate()?;
236 if checkpoint.body().continuity_sequence != 0
237 || authority.finalized_checkpoint_digest != checkpoint.digest()
238 || authority.genesis_policy_id != policy.policy_id
239 || authority.genesis_policy_digest != policy.digest()?
240 || authority.current_charter_id != charter.body().charter_id
241 || authority.current_charter_digest != charter.digest()
242 || checkpoint.body().runtime_readiness_digest != readiness.digest()
243 {
244 return Err(invariant("fiscal genesis bindings do not match"));
245 }
246 let policy_json = canonical_json_bytes(policy).map_err(canonical_error)?;
247 let authority_json = canonical_json_bytes(authority).map_err(canonical_error)?;
248 let charter_json = charter.canonical_bytes()?;
249 let readiness_json = readiness.canonical_bytes()?;
250 let registry_json = readiness.runtime_registry().canonical_bytes()?;
251 let checkpoint_json = checkpoint.canonical_bytes()?;
252 let snapshot_digest = canonical_digest(&FiscalGenesisSnapshot {
253 policy,
254 authority,
255 charter_digest: charter.digest(),
256 readiness_digest: readiness.digest(),
257 checkpoint_digest: checkpoint.digest(),
258 })?;
259
260 let mut connection = self.connection()?;
261 let transaction = self.begin_write(&mut connection, fence)?;
262 if transaction
263 .query_row(
264 "SELECT state_json FROM fiscal_authority_state WHERE singleton = 1",
265 [],
266 |row| row.get::<_, Vec<u8>>(0),
267 )
268 .optional()
269 .map_err(sqlite_error)?
270 .is_some()
271 {
272 verify_exact_genesis(
273 &transaction,
274 &policy_json,
275 &authority_json,
276 &charter_json,
277 &readiness_json,
278 ®istry_json,
279 &checkpoint_json,
280 policy,
281 charter,
282 readiness,
283 checkpoint,
284 )?;
285 transaction.commit().map_err(sqlite_error)?;
286 return Ok(());
287 }
288
289 transaction
290 .execute(
291 "INSERT INTO fiscal_genesis_policies (policy_id, policy_digest, policy_json) VALUES (?1, ?2, ?3)",
292 params![&policy.policy_id, policy.digest()?, &policy_json],
293 )
294 .map_err(sqlite_error)?;
295 transaction
296 .execute(
297 "INSERT INTO fiscal_charters (charter_id, charter_digest, charter_sequence, lifecycle_state, signed_json) VALUES (?1, ?2, ?3, 'pinned', ?4)",
298 params![
299 &charter.body().charter_id,
300 charter.digest(),
301 sqlite_i64(charter.body().sequence, "charter sequence")?,
302 &charter_json,
303 ],
304 )
305 .map_err(sqlite_error)?;
306 transaction
307 .execute(
308 "INSERT INTO fiscal_runtime_readiness (readiness_id, readiness_digest, readiness_sequence, registry_json, signed_json) VALUES (?1, ?2, ?3, ?4, ?5)",
309 params![
310 &readiness.body().readiness_id,
311 readiness.digest(),
312 sqlite_i64(readiness.body().readiness_sequence, "readiness sequence")?,
313 ®istry_json,
314 &readiness_json,
315 ],
316 )
317 .map_err(sqlite_error)?;
318 transaction
319 .execute(
320 "INSERT INTO fiscal_continuity_checkpoints (checkpoint_digest, continuity_sequence, status, signed_json) VALUES (?1, 0, 'finalized', ?2)",
321 params![checkpoint.digest(), &checkpoint_json],
322 )
323 .map_err(sqlite_error)?;
324 transaction
325 .execute(
326 "INSERT INTO fiscal_authority_state (singleton, state_json, finalized_checkpoint_digest, state_version) VALUES (1, ?1, ?2, 1)",
327 params![&authority_json, checkpoint.digest()],
328 )
329 .map_err(sqlite_error)?;
330 append_projection_commit(
331 &transaction,
332 &self.serving_owner,
333 FISCAL_GENESIS_PROJECTION_KEY,
334 1,
335 "initialize_genesis",
336 &snapshot_digest,
337 )?;
338 self.commit_write(transaction)?;
339 self.sync_after_write(&connection)
340 }
341
342 pub fn persist_proposal(
343 &self,
344 proposal: &VerifiedFiscalProposal,
345 fence: &StoreMutationFence,
346 ) -> Result<(), FiscalStoreError> {
347 let json = canonical_json_bytes(proposal.signed()).map_err(canonical_error)?;
348 self.persist_immutable_artifact(
349 "fiscal_proposals",
350 "proposal_id",
351 &proposal.body().proposal_id,
352 "proposal_digest",
353 proposal.digest(),
354 &json,
355 "persist_proposal",
356 fence,
357 )
358 }
359
360 #[allow(clippy::too_many_arguments)]
361 pub fn admit_proposal(
362 &self,
363 proposal: &VerifiedFiscalProposal,
364 current_charter: &VerifiedFiscalCharter,
365 admission_authority_id: &str,
366 signer_key_epoch: u64,
367 signer: &Keypair,
368 admitted_at: u64,
369 fence: &StoreMutationFence,
370 ) -> Result<VerifiedFiscalProposalAdmission, FiscalStoreError> {
371 let mut connection = self.connection()?;
372 let transaction = self.begin_write(&mut connection, fence)?;
373 let proposal_present = transaction
374 .query_row(
375 "SELECT proposal_digest = ?1 AND signed_json = ?2 FROM fiscal_proposals WHERE proposal_id = ?3",
376 params![
377 proposal.digest(),
378 canonical_json_bytes(proposal.signed()).map_err(canonical_error)?,
379 &proposal.body().proposal_id,
380 ],
381 |row| row.get::<_, bool>(0),
382 )
383 .optional()
384 .map_err(sqlite_error)?
385 .unwrap_or(false);
386 let charter_present = transaction
387 .query_row(
388 "SELECT charter_digest = ?1 AND signed_json = ?2 AND lifecycle_state IN ('pinned', 'active') FROM fiscal_charters WHERE charter_id = ?3",
389 params![
390 current_charter.digest(),
391 current_charter.canonical_bytes()?,
392 ¤t_charter.body().charter_id,
393 ],
394 |row| row.get::<_, bool>(0),
395 )
396 .optional()
397 .map_err(sqlite_error)?
398 .unwrap_or(false);
399 let proposal_already_admitted = transaction
400 .query_row(
401 "SELECT EXISTS(SELECT 1 FROM fiscal_proposal_admissions WHERE json_extract(state_json, '$.signedAdmission.body.proposalId') = ?1)",
402 [&proposal.body().proposal_id],
403 |row| row.get::<_, bool>(0),
404 )
405 .map_err(sqlite_error)?;
406 if !proposal_present || !charter_present || proposal_already_admitted {
407 return Err(FiscalStoreError::Conflict);
408 }
409 let current_sequence = transaction
410 .query_row(
411 "SELECT current_sequence FROM fiscal_admission_sequence WHERE singleton = 1",
412 [],
413 |row| row.get::<_, i64>(0),
414 )
415 .map_err(sqlite_error)?;
416 let admission_sequence = u64::try_from(current_sequence)
417 .map_err(|_| invariant("fiscal admission sequence is negative"))?
418 .checked_add(1)
419 .ok_or_else(|| invariant("fiscal admission sequence overflow"))?;
420 let signed = FiscalProposalAdmissionBuilder {
421 admission_sequence,
422 admitted_at,
423 admission_authority_id: admission_authority_id.to_owned(),
424 signer_key_epoch,
425 }
426 .sign(proposal, current_charter, signer)?;
427 let trust = FiscalAdmissionTrustRegistry::new(vec![FiscalAdmissionAuthority::new(
428 current_charter.body().governing_operator_id.clone(),
429 admission_authority_id.to_owned(),
430 signer_key_epoch,
431 signer.public_key(),
432 )?])?;
433 let admission = VerifiedFiscalProposalAdmission::verify(
434 signed,
435 proposal,
436 current_charter,
437 &trust,
438 admitted_at,
439 )?;
440 let state = FiscalProposalAdmissionState::admitted(&admission);
441 let state_json = canonical_json_bytes(&state).map_err(canonical_error)?;
442 let sequence_updated = transaction
443 .execute(
444 "UPDATE fiscal_admission_sequence SET current_sequence = ?1 WHERE singleton = 1 AND current_sequence = ?2",
445 params![
446 sqlite_i64(admission_sequence, "admission sequence")?,
447 current_sequence,
448 ],
449 )
450 .map_err(sqlite_error)?;
451 let admission_inserted = transaction
452 .execute(
453 "INSERT INTO fiscal_proposal_admissions (admission_id, admission_digest, admitted_at, status, state_version, state_json) VALUES (?1, ?2, ?3, 'admitted', 1, ?4)",
454 params![
455 &admission.body().admission_id,
456 admission.digest(),
457 sqlite_i64(admitted_at, "admitted_at")?,
458 &state_json,
459 ],
460 )
461 .map_err(sqlite_error)?;
462 if sequence_updated != 1 || admission_inserted != 1 {
463 return Err(FiscalStoreError::Conflict);
464 }
465 append_projection_commit(
466 &transaction,
467 &self.serving_owner,
468 &format!("admission:{}", admission.body().admission_id),
469 1,
470 "admit_fiscal_proposal",
471 &sha256_hex(&state_json),
472 )?;
473 self.commit_write(transaction)?;
474 self.sync_after_write(&connection)?;
475 Ok(admission)
476 }
477
478 pub fn persist_charter(
479 &self,
480 charter: &VerifiedFiscalCharter,
481 fence: &StoreMutationFence,
482 ) -> Result<(), FiscalStoreError> {
483 let json = charter.canonical_bytes()?;
484 let mut connection = self.connection()?;
485 let transaction = self.begin_write(&mut connection, fence)?;
486 if let Some((digest, stored)) = transaction
487 .query_row(
488 "SELECT charter_digest, signed_json FROM fiscal_charters WHERE charter_id = ?1",
489 [&charter.body().charter_id],
490 |row| Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?)),
491 )
492 .optional()
493 .map_err(sqlite_error)?
494 {
495 if digest != charter.digest() || stored != json {
496 return Err(FiscalStoreError::Conflict);
497 }
498 transaction.commit().map_err(sqlite_error)?;
499 return Ok(());
500 }
501 transaction
502 .execute(
503 "INSERT INTO fiscal_charters (charter_id, charter_digest, charter_sequence, lifecycle_state, signed_json) VALUES (?1, ?2, ?3, 'proposed', ?4)",
504 params![
505 &charter.body().charter_id,
506 charter.digest(),
507 sqlite_i64(charter.body().sequence, "charter sequence")?,
508 &json,
509 ],
510 )
511 .map_err(sqlite_error)?;
512 let projection_key = format!("charter:{}", charter.body().charter_id);
513 append_projection_commit(
514 &transaction,
515 &self.serving_owner,
516 &projection_key,
517 1,
518 "persist_fiscal_charter",
519 charter.digest(),
520 )?;
521 self.commit_write(transaction)?;
522 self.sync_after_write(&connection)
523 }
524
525 pub fn persist_schedule(
526 &self,
527 schedule: &VerifiedFiscalSchedule,
528 fence: &StoreMutationFence,
529 ) -> Result<(), FiscalStoreError> {
530 let json = schedule.canonical_bytes()?;
531 let digest = sha256_hex(&json);
532 let domain_json = String::from_utf8(
533 canonical_json_bytes(&schedule.body().domain).map_err(canonical_error)?,
534 )
535 .map_err(|error| invariant(format!("fiscal domain encoding is not UTF-8: {error}")))?;
536 let mut connection = self.connection()?;
537 let transaction = self.begin_write(&mut connection, fence)?;
538 if let Some((stored_digest, stored)) = transaction
539 .query_row(
540 "SELECT schedule_digest, signed_json FROM fiscal_schedules WHERE schedule_id = ?1",
541 [&schedule.body().schedule_id],
542 |row| Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?)),
543 )
544 .optional()
545 .map_err(sqlite_error)?
546 {
547 if stored_digest != digest || stored != json {
548 return Err(FiscalStoreError::Conflict);
549 }
550 transaction.commit().map_err(sqlite_error)?;
551 return Ok(());
552 }
553 transaction
554 .execute(
555 "INSERT INTO fiscal_schedules (schedule_id, schedule_digest, domain_json, schedule_sequence, lifecycle_state, signed_json) VALUES (?1, ?2, ?3, ?4, 'staged', ?5)",
556 params![
557 &schedule.body().schedule_id,
558 &digest,
559 &domain_json,
560 sqlite_i64(schedule.body().sequence, "schedule sequence")?,
561 &json,
562 ],
563 )
564 .map_err(sqlite_error)?;
565 let projection_key = format!("schedule:{}", schedule.body().schedule_id);
566 append_projection_commit(
567 &transaction,
568 &self.serving_owner,
569 &projection_key,
570 1,
571 "persist_fiscal_schedule",
572 &digest,
573 )?;
574 self.commit_write(transaction)?;
575 self.sync_after_write(&connection)
576 }
577
578 pub fn persist_runtime_readiness(
579 &self,
580 readiness: &VerifiedFiscalRuntimeReadiness,
581 fence: &StoreMutationFence,
582 ) -> Result<(), FiscalStoreError> {
583 let signed_json = readiness.canonical_bytes()?;
584 let registry_json = readiness.runtime_registry().canonical_bytes()?;
585 let mut connection = self.connection()?;
586 let transaction = self.begin_write(&mut connection, fence)?;
587 if let Some((digest, stored_registry, stored_signed)) = transaction
588 .query_row(
589 "SELECT readiness_digest, registry_json, signed_json FROM fiscal_runtime_readiness WHERE readiness_id = ?1",
590 [&readiness.body().readiness_id],
591 |row| {
592 Ok((
593 row.get::<_, String>(0)?,
594 row.get::<_, Vec<u8>>(1)?,
595 row.get::<_, Vec<u8>>(2)?,
596 ))
597 },
598 )
599 .optional()
600 .map_err(sqlite_error)?
601 {
602 if digest != readiness.digest()
603 || stored_registry != registry_json
604 || stored_signed != signed_json
605 {
606 return Err(FiscalStoreError::Conflict);
607 }
608 transaction.commit().map_err(sqlite_error)?;
609 return Ok(());
610 }
611 transaction
612 .execute(
613 "INSERT INTO fiscal_runtime_readiness (readiness_id, readiness_digest, readiness_sequence, registry_json, signed_json) VALUES (?1, ?2, ?3, ?4, ?5)",
614 params![
615 &readiness.body().readiness_id,
616 readiness.digest(),
617 sqlite_i64(readiness.body().readiness_sequence, "readiness sequence")?,
618 ®istry_json,
619 &signed_json,
620 ],
621 )
622 .map_err(sqlite_error)?;
623 let projection_key = format!("readiness:{}", readiness.body().readiness_id);
624 append_projection_commit(
625 &transaction,
626 &self.serving_owner,
627 &projection_key,
628 1,
629 "persist_fiscal_readiness",
630 readiness.digest(),
631 )?;
632 self.commit_write(transaction)?;
633 self.sync_after_write(&connection)
634 }
635
636 pub fn load_genesis_policy(&self) -> Result<FiscalGenesisPolicy, FiscalStoreError> {
637 let mut connection = self.connection()?;
638 let transaction = self.begin_read(&mut connection)?;
639 let (policy_id, policy_digest, policy_json) = transaction
640 .query_row(
641 "SELECT policy_id, policy_digest, policy_json FROM fiscal_genesis_policies",
642 [],
643 |row| {
644 Ok((
645 row.get::<_, String>(0)?,
646 row.get::<_, String>(1)?,
647 row.get::<_, Vec<u8>>(2)?,
648 ))
649 },
650 )
651 .optional()
652 .map_err(sqlite_error)?
653 .ok_or(FiscalStoreError::NotFound)?;
654 let policy: FiscalGenesisPolicy =
655 serde_json::from_slice(&policy_json).map_err(|error| {
656 invariant(format!("stored fiscal genesis policy is invalid: {error}"))
657 })?;
658 if canonical_json_bytes(&policy).map_err(canonical_error)? != policy_json
659 || policy.policy_id != policy_id
660 || policy.digest()? != policy_digest
661 {
662 return Err(invariant("stored fiscal genesis policy binding is invalid"));
663 }
664 transaction.commit().map_err(sqlite_error)?;
665 Ok(policy)
666 }
667
668 pub fn load_charter_registry(&self) -> Result<FiscalCharterRegistry, FiscalStoreError> {
669 let mut connection = self.connection()?;
670 let transaction = self.begin_read(&mut connection)?;
671 let mut statement = transaction
672 .prepare("SELECT signed_json FROM fiscal_charters ORDER BY charter_sequence")
673 .map_err(sqlite_error)?;
674 let signed = statement
675 .query_map([], |row| row.get::<_, Vec<u8>>(0))
676 .map_err(sqlite_error)?
677 .map(|row| {
678 let bytes = row.map_err(sqlite_error)?;
679 Ok(VerifiedFiscalCharter::from_canonical_bytes(&bytes)?
680 .signed()
681 .clone())
682 })
683 .collect::<Result<Vec<_>, FiscalStoreError>>()?;
684 drop(statement);
685 transaction.commit().map_err(sqlite_error)?;
686 Ok(FiscalCharterRegistry::new(signed)?)
687 }
688
689 pub fn load_runtime_readiness(
690 &self,
691 readiness_digest: &str,
692 policy: &FiscalGenesisPolicy,
693 ) -> Result<VerifiedFiscalRuntimeReadiness, FiscalStoreError> {
694 let mut connection = self.connection()?;
695 let transaction = self.begin_read(&mut connection)?;
696 let (registry_json, signed_json) = transaction
697 .query_row(
698 "SELECT registry_json, signed_json FROM fiscal_runtime_readiness WHERE readiness_digest = ?1",
699 [readiness_digest],
700 |row| Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?)),
701 )
702 .optional()
703 .map_err(sqlite_error)?
704 .ok_or(FiscalStoreError::NotFound)?;
705 let registry = FiscalRuntimeAdapterRegistry::from_canonical_bytes(®istry_json)?;
706 let readiness =
707 VerifiedFiscalRuntimeReadiness::from_canonical_bytes(&signed_json, policy, registry)?;
708 if readiness.digest() != readiness_digest {
709 return Err(invariant(
710 "stored fiscal runtime readiness digest is invalid",
711 ));
712 }
713 transaction.commit().map_err(sqlite_error)?;
714 Ok(readiness)
715 }
716
717 pub fn load_verified_schedule(
718 &self,
719 schedule_id: &str,
720 charters: &FiscalCharterRegistry,
721 ) -> Result<VerifiedFiscalSchedule, FiscalStoreError> {
722 let mut connection = self.connection()?;
723 let transaction = self.begin_read(&mut connection)?;
724 let mut lineage = Vec::new();
725 let mut next_id = Some(schedule_id.to_owned());
726 let mut seen = HashSet::new();
727 while let Some(id) = next_id {
728 if lineage.len() >= 4096 || !seen.insert(id.clone()) {
729 return Err(invariant(
730 "stored fiscal schedule lineage is cyclic or too deep",
731 ));
732 }
733 let bytes = transaction
734 .query_row(
735 "SELECT signed_json FROM fiscal_schedules WHERE schedule_id = ?1",
736 [&id],
737 |row| row.get::<_, Vec<u8>>(0),
738 )
739 .optional()
740 .map_err(sqlite_error)?
741 .ok_or(FiscalStoreError::NotFound)?;
742 let signed: chio_fiscal::SignedFiscalSchedule = serde_json::from_slice(&bytes)
743 .map_err(|error| {
744 invariant(format!("stored fiscal schedule is invalid: {error}"))
745 })?;
746 if canonical_json_bytes(&signed).map_err(canonical_error)? != bytes
747 || signed.body.schedule_id != id
748 {
749 return Err(invariant(
750 "stored fiscal schedule canonical binding is invalid",
751 ));
752 }
753 next_id = signed.body.supersedes_schedule_id.clone();
754 lineage.push(signed);
755 }
756 lineage.reverse();
757 let mut verified: Option<VerifiedFiscalSchedule> = None;
758 for signed in lineage {
759 let charter = charters.resolve(&signed.body.charter_id, &signed.body.charter_digest)?;
760 let next = match verified.as_ref() {
761 Some(predecessor)
762 if predecessor.body().charter_digest != signed.body.charter_digest =>
763 {
764 VerifiedFiscalSchedule::verify_rotation_replacement(
765 signed,
766 &charter,
767 predecessor,
768 )?
769 }
770 predecessor => VerifiedFiscalSchedule::verify(signed, &charter, predecessor)?,
771 };
772 verified = Some(next);
773 }
774 let verified = verified.ok_or(FiscalStoreError::NotFound)?;
775 transaction.commit().map_err(sqlite_error)?;
776 Ok(verified)
777 }
778
779 pub fn load_signed_schedules(&self) -> Result<Vec<SignedFiscalSchedule>, FiscalStoreError> {
780 self.load_signed_artifacts(
781 "SELECT signed_json FROM fiscal_schedules ORDER BY domain_json, schedule_sequence",
782 "stored fiscal schedule",
783 )
784 }
785
786 pub fn load_signed_proposals(&self) -> Result<Vec<SignedFiscalProposal>, FiscalStoreError> {
787 self.load_signed_artifacts(
788 "SELECT signed_json FROM fiscal_proposals ORDER BY proposal_id",
789 "stored fiscal proposal",
790 )
791 }
792
793 pub fn load_signed_activations(&self) -> Result<Vec<SignedFiscalActivation>, FiscalStoreError> {
794 self.load_signed_artifacts(
795 "SELECT activation.signed_json FROM fiscal_activations AS activation WHERE EXISTS(SELECT 1 FROM fiscal_staged_activation_mutations AS mutation JOIN fiscal_staged_transitions AS stage ON stage.transition_id = mutation.transition_id WHERE mutation.activation_id = activation.activation_id AND mutation.activation_digest = activation.activation_digest AND stage.status = 'db_finalized') OR EXISTS(SELECT 1 FROM fiscal_staged_rotation_mutations AS mutation JOIN fiscal_staged_transitions AS stage ON stage.transition_id = mutation.transition_id WHERE mutation.activation_id = activation.activation_id AND mutation.activation_digest = activation.activation_digest AND stage.status = 'db_finalized') ORDER BY activation.activation_id",
796 "stored fiscal activation",
797 )
798 }
799
800 fn load_signed_artifacts<T: serde::de::DeserializeOwned + Serialize>(
801 &self,
802 sql: &str,
803 label: &str,
804 ) -> Result<Vec<T>, FiscalStoreError> {
805 let mut connection = self.connection()?;
806 let transaction = self.begin_read(&mut connection)?;
807 let mut statement = transaction.prepare(sql).map_err(sqlite_error)?;
808 let artifacts = statement
809 .query_map([], |row| row.get::<_, Vec<u8>>(0))
810 .map_err(sqlite_error)?
811 .map(|row| {
812 let bytes = row.map_err(sqlite_error)?;
813 let artifact: T = serde_json::from_slice(&bytes)
814 .map_err(|error| invariant(format!("{label} is invalid: {error}")))?;
815 if canonical_json_bytes(&artifact).map_err(canonical_error)? != bytes {
816 return Err(invariant(format!("{label} is not canonical")));
817 }
818 Ok(artifact)
819 })
820 .collect::<Result<Vec<_>, FiscalStoreError>>()?;
821 drop(statement);
822 transaction.commit().map_err(sqlite_error)?;
823 Ok(artifacts)
824 }
825
826 pub fn persist_approval(
827 &self,
828 approval: &VerifiedFiscalApproval,
829 fence: &StoreMutationFence,
830 ) -> Result<(), FiscalStoreError> {
831 let json = canonical_json_bytes(approval.signed()).map_err(canonical_error)?;
832 self.persist_immutable_artifact(
833 "fiscal_approvals",
834 "approval_id",
835 &approval.body().approval_id,
836 "approval_digest",
837 approval.digest(),
838 &json,
839 "persist_approval",
840 fence,
841 )
842 }
843
844 pub fn require_approval(
845 &self,
846 approval: &VerifiedFiscalApproval,
847 ) -> Result<(), FiscalStoreError> {
848 let mut connection = self.connection()?;
849 let transaction = self.begin_read(&mut connection)?;
850 let exact = transaction
851 .query_row(
852 "SELECT approval_digest = ?1 AND signed_json = ?2 FROM fiscal_approvals WHERE approval_id = ?3",
853 params![
854 approval.digest(),
855 canonical_json_bytes(approval.signed()).map_err(canonical_error)?,
856 &approval.body().approval_id,
857 ],
858 |row| row.get::<_, bool>(0),
859 )
860 .optional()
861 .map_err(sqlite_error)?
862 .unwrap_or(false);
863 if !exact {
864 return Err(FiscalStoreError::Conflict);
865 }
866 transaction.commit().map_err(sqlite_error)
867 }
868
869 pub fn load_admission_state(
870 &self,
871 admission_id: &str,
872 ) -> Result<FiscalProposalAdmissionState, FiscalStoreError> {
873 let mut connection = self.connection()?;
874 let transaction = self.begin_read(&mut connection)?;
875 let (digest, status, version, json) = transaction
876 .query_row(
877 "SELECT admission_digest, status, state_version, state_json FROM fiscal_proposal_admissions WHERE admission_id = ?1",
878 [admission_id],
879 |row| {
880 Ok((
881 row.get::<_, String>(0)?,
882 row.get::<_, String>(1)?,
883 row.get::<_, i64>(2)?,
884 row.get::<_, Vec<u8>>(3)?,
885 ))
886 },
887 )
888 .optional()
889 .map_err(sqlite_error)?
890 .ok_or(FiscalStoreError::NotFound)?;
891 let state: FiscalProposalAdmissionState =
892 serde_json::from_slice(&json).map_err(|error| {
893 invariant(format!("stored fiscal admission state is invalid: {error}"))
894 })?;
895 let expected_status = match state.status {
896 FiscalProposalAdmissionStatus::Admitted => "admitted",
897 FiscalProposalAdmissionStatus::Activated => "activated",
898 };
899 if canonical_json_bytes(&state).map_err(canonical_error)? != json
900 || state.signed_admission.body.admission_id != admission_id
901 || state.admission_digest != digest
902 || expected_status != status
903 || sqlite_i64(state.version, "admission state version")? != version
904 {
905 return Err(invariant("stored fiscal admission binding is invalid"));
906 }
907 transaction.commit().map_err(sqlite_error)?;
908 Ok(state)
909 }
910
911 pub fn persist_activation(
912 &self,
913 activation: &VerifiedFiscalActivation,
914 fence: &StoreMutationFence,
915 ) -> Result<(), FiscalStoreError> {
916 let json = canonical_json_bytes(activation.signed()).map_err(canonical_error)?;
917 self.persist_immutable_artifact(
918 "fiscal_activations",
919 "activation_id",
920 &activation.body().activation_id,
921 "activation_digest",
922 activation.digest(),
923 &json,
924 "persist_activation",
925 fence,
926 )
927 }
928
929 pub fn bind_legacy_fee_schedule(
930 &self,
931 legacy_schedule: &SignedOpenMarketFeeSchedule,
932 schedule: &VerifiedFiscalSchedule,
933 fence: &StoreMutationFence,
934 ) -> Result<(), FiscalStoreError> {
935 legacy_schedule
936 .body
937 .validate()
938 .map_err(|error| invariant(format!("legacy fee schedule is invalid: {error}")))?;
939 if !legacy_schedule
940 .verify_signature()
941 .map_err(|error| invariant(format!("legacy fee schedule signature failed: {error}")))?
942 {
943 return Err(invariant("legacy fee schedule signature is invalid"));
944 }
945 let FiscalParams::OpenMarketFeeAndBondSchedule { legacy_body } = &schedule.body().params
946 else {
947 return Err(FiscalStoreError::Conflict);
948 };
949 if legacy_body.as_ref() != &legacy_schedule.body {
950 return Err(FiscalStoreError::Conflict);
951 }
952 let legacy_schedule_id = &legacy_schedule.body.fee_schedule_id;
953 let legacy_envelope_json =
954 canonical_json_bytes(legacy_schedule).map_err(canonical_error)?;
955 let legacy_envelope_digest = sha256_hex(&legacy_envelope_json);
956 let schedule_json = schedule.canonical_bytes()?;
957 let schedule_digest = sha256_hex(&schedule_json);
958 let mut connection = self.connection()?;
959 let transaction = self.begin_write(&mut connection, fence)?;
960 let retained = transaction
961 .query_row(
962 "SELECT fiscal_schedule_id, fiscal_schedule_digest, legacy_envelope_digest FROM fiscal_legacy_fee_schedule_bindings WHERE legacy_schedule_id = ?1",
963 [legacy_schedule_id],
964 |row| {
965 Ok((
966 row.get::<_, String>(0)?,
967 row.get::<_, String>(1)?,
968 row.get::<_, String>(2)?,
969 ))
970 },
971 )
972 .optional()
973 .map_err(sqlite_error)?;
974 if let Some((schedule_id, digest, envelope_digest)) = retained {
975 if schedule_id != schedule.body().schedule_id
976 || digest != schedule_digest
977 || envelope_digest != legacy_envelope_digest
978 {
979 return Err(FiscalStoreError::Conflict);
980 }
981 transaction.commit().map_err(sqlite_error)?;
982 return Ok(());
983 }
984 let exact_schedule = transaction
985 .query_row(
986 "SELECT schedule_digest = ?1 AND signed_json = ?2 FROM fiscal_schedules WHERE schedule_id = ?3",
987 params![&schedule_digest, &schedule_json, &schedule.body().schedule_id],
988 |row| row.get::<_, bool>(0),
989 )
990 .optional()
991 .map_err(sqlite_error)?
992 .unwrap_or(false);
993 if !exact_schedule {
994 return Err(FiscalStoreError::Conflict);
995 }
996 transaction
997 .execute(
998 "INSERT INTO fiscal_legacy_fee_schedule_bindings (legacy_schedule_id, fiscal_schedule_id, fiscal_schedule_digest, legacy_envelope_digest) VALUES (?1, ?2, ?3, ?4)",
999 params![
1000 legacy_schedule_id,
1001 &schedule.body().schedule_id,
1002 &schedule_digest,
1003 &legacy_envelope_digest,
1004 ],
1005 )
1006 .map_err(sqlite_error)?;
1007 let projection_key = format!("legacy-fee:{legacy_schedule_id}");
1008 append_projection_commit(
1009 &transaction,
1010 &self.serving_owner,
1011 &projection_key,
1012 1,
1013 "bind_legacy_fee_schedule",
1014 &schedule_digest,
1015 )?;
1016 self.commit_write(transaction)?;
1017 self.sync_after_write(&connection)
1018 }
1019
1020 pub fn load_legacy_fee_schedule_binding(
1021 &self,
1022 fiscal_schedule_id: &str,
1023 ) -> Result<FiscalLegacyFeeScheduleBindingRecord, FiscalStoreError> {
1024 let mut connection = self.connection()?;
1025 let transaction = self.begin_read(&mut connection)?;
1026 let record = transaction
1027 .query_row(
1028 "SELECT legacy_schedule_id, fiscal_schedule_id, fiscal_schedule_digest, legacy_envelope_digest FROM fiscal_legacy_fee_schedule_bindings WHERE fiscal_schedule_id = ?1",
1029 [fiscal_schedule_id],
1030 |row| {
1031 Ok(FiscalLegacyFeeScheduleBindingRecord {
1032 legacy_schedule_id: row.get(0)?,
1033 fiscal_schedule_id: row.get(1)?,
1034 fiscal_schedule_digest: row.get(2)?,
1035 legacy_envelope_digest: row.get(3)?,
1036 })
1037 },
1038 )
1039 .optional()
1040 .map_err(sqlite_error)?
1041 .ok_or(FiscalStoreError::NotFound)?;
1042 transaction.commit().map_err(sqlite_error)?;
1043 Ok(record)
1044 }
1045
1046 #[allow(clippy::too_many_arguments)]
1047 fn persist_immutable_artifact(
1048 &self,
1049 table: &str,
1050 id_column: &str,
1051 id: &str,
1052 digest_column: &str,
1053 digest: &str,
1054 json: &[u8],
1055 mutation_kind: &str,
1056 fence: &StoreMutationFence,
1057 ) -> Result<(), FiscalStoreError> {
1058 let mut connection = self.connection()?;
1059 let transaction = self.begin_write(&mut connection, fence)?;
1060 let sql =
1061 format!("SELECT {digest_column}, signed_json FROM {table} WHERE {id_column} = ?1");
1062 if let Some((stored_digest, stored_json)) = transaction
1063 .query_row(&sql, [id], |row| {
1064 Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?))
1065 })
1066 .optional()
1067 .map_err(sqlite_error)?
1068 {
1069 if stored_digest != digest || stored_json != json {
1070 return Err(FiscalStoreError::Conflict);
1071 }
1072 transaction.commit().map_err(sqlite_error)?;
1073 return Ok(());
1074 }
1075 let insert = format!(
1076 "INSERT INTO {table} ({id_column}, {digest_column}, signed_json) VALUES (?1, ?2, ?3)"
1077 );
1078 transaction
1079 .execute(&insert, params![id, digest, json])
1080 .map_err(sqlite_error)?;
1081 let projection_key = format!("artifact:{id}");
1082 append_projection_commit(
1083 &transaction,
1084 &self.serving_owner,
1085 &projection_key,
1086 1,
1087 mutation_kind,
1088 digest,
1089 )?;
1090 self.commit_write(transaction)?;
1091 self.sync_after_write(&connection)
1092 }
1093
1094 pub fn persist_admission_state(
1095 &self,
1096 state: &FiscalProposalAdmissionState,
1097 expected_version: Option<u64>,
1098 fence: &StoreMutationFence,
1099 ) -> Result<(), FiscalStoreError> {
1100 let state_json = canonical_json_bytes(state).map_err(canonical_error)?;
1101 let id = &state.signed_admission.body.admission_id;
1102 let status = match state.status {
1103 FiscalProposalAdmissionStatus::Admitted => "admitted",
1104 FiscalProposalAdmissionStatus::Activated => "activated",
1105 };
1106 let mut connection = self.connection()?;
1107 let transaction = self.begin_write(&mut connection, fence)?;
1108 if expected_version.is_none() {
1109 let current_sequence = transaction
1110 .query_row(
1111 "SELECT current_sequence FROM fiscal_admission_sequence WHERE singleton = 1",
1112 [],
1113 |row| row.get::<_, i64>(0),
1114 )
1115 .map_err(sqlite_error)?;
1116 let expected_sequence = u64::try_from(current_sequence)
1117 .map_err(|_| invariant("fiscal admission sequence is negative"))?
1118 .checked_add(1)
1119 .ok_or_else(|| invariant("fiscal admission sequence overflow"))?;
1120 if state.signed_admission.body.admission_sequence != expected_sequence
1121 || transaction
1122 .execute(
1123 "UPDATE fiscal_admission_sequence SET current_sequence = ?1 WHERE singleton = 1 AND current_sequence = ?2",
1124 params![
1125 sqlite_i64(expected_sequence, "admission sequence")?,
1126 current_sequence,
1127 ],
1128 )
1129 .map_err(sqlite_error)?
1130 != 1
1131 {
1132 return Err(FiscalStoreError::Conflict);
1133 }
1134 }
1135 let changed = match expected_version {
1136 None => transaction.execute(
1137 "INSERT INTO fiscal_proposal_admissions (admission_id, admission_digest, admitted_at, status, state_version, state_json) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
1138 params![
1139 id,
1140 &state.admission_digest,
1141 sqlite_i64(state.signed_admission.body.admitted_at, "admitted_at")?,
1142 status,
1143 sqlite_i64(state.version, "admission state version")?,
1144 &state_json,
1145 ],
1146 ),
1147 Some(version) => transaction.execute(
1148 "UPDATE fiscal_proposal_admissions SET status = ?1, state_version = ?2, state_json = ?3 WHERE admission_id = ?4 AND admission_digest = ?5 AND state_version = ?6",
1149 params![
1150 status,
1151 sqlite_i64(state.version, "admission state version")?,
1152 &state_json,
1153 id,
1154 &state.admission_digest,
1155 sqlite_i64(version, "expected admission state version")?,
1156 ],
1157 ),
1158 }
1159 .map_err(sqlite_error)?;
1160 if changed != 1 {
1161 return Err(FiscalStoreError::Conflict);
1162 }
1163 let projection_key = format!("admission:{id}");
1164 append_projection_commit(
1165 &transaction,
1166 &self.serving_owner,
1167 &projection_key,
1168 state.version,
1169 "persist_admission_state",
1170 &sha256_hex(&state_json),
1171 )?;
1172 self.commit_write(transaction)?;
1173 self.sync_after_write(&connection)
1174 }
1175
1176 pub fn stage_advance(
1177 &self,
1178 advance: &VerifiedFiscalContinuityAdvance,
1179 next_authority: &FiscalAuthorityState,
1180 fence: &StoreMutationFence,
1181 ) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
1182 self.stage_advance_inner(advance, next_authority, None, None, fence)
1183 }
1184
1185 #[allow(clippy::too_many_arguments)]
1186 pub fn stage_activation_advance(
1187 &self,
1188 advance: &VerifiedFiscalContinuityAdvance,
1189 next_authority: &FiscalAuthorityState,
1190 activation: &VerifiedFiscalActivation,
1191 activated_admission: &FiscalProposalAdmissionState,
1192 candidate: &VerifiedFiscalSchedule,
1193 predecessor: Option<&VerifiedFiscalSchedule>,
1194 fence: &StoreMutationFence,
1195 ) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
1196 let mutation = prepare_activation_mutation(
1197 advance,
1198 activation,
1199 activated_admission,
1200 candidate,
1201 predecessor,
1202 )?;
1203 self.stage_advance_inner(advance, next_authority, Some(&mutation), None, fence)
1204 }
1205
1206 #[allow(clippy::too_many_arguments)]
1207 pub fn stage_charter_rotation_advance(
1208 &self,
1209 advance: &VerifiedFiscalContinuityAdvance,
1210 next_authority: &FiscalAuthorityState,
1211 activation: &VerifiedFiscalActivation,
1212 activated_admission: &FiscalProposalAdmissionState,
1213 successor_charter: &VerifiedFiscalCharter,
1214 successor_schedules: &[VerifiedFiscalSchedule],
1215 predecessor_charter: &VerifiedFiscalCharter,
1216 predecessor_schedules: &[VerifiedFiscalSchedule],
1217 fence: &StoreMutationFence,
1218 ) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
1219 let mutation = prepare_rotation_mutation(
1220 advance,
1221 activation,
1222 activated_admission,
1223 successor_charter,
1224 successor_schedules,
1225 predecessor_charter,
1226 predecessor_schedules,
1227 )?;
1228 self.stage_advance_inner(advance, next_authority, None, Some(&mutation), fence)
1229 }
1230
1231 fn stage_advance_inner(
1232 &self,
1233 advance: &VerifiedFiscalContinuityAdvance,
1234 next_authority: &FiscalAuthorityState,
1235 activation_mutation: Option<&PreparedFiscalActivationMutation>,
1236 rotation_mutation: Option<&PreparedFiscalRotationMutation>,
1237 fence: &StoreMutationFence,
1238 ) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
1239 if activation_mutation.is_some() && rotation_mutation.is_some() {
1240 return Err(FiscalStoreError::Conflict);
1241 }
1242 next_authority.validate()?;
1243 if next_authority.finalized_checkpoint_digest != advance.next().digest() {
1244 return Err(invariant("staged authority does not match next checkpoint"));
1245 }
1246 let transition_id = advance
1247 .next()
1248 .body()
1249 .staged_transition
1250 .as_ref()
1251 .map(|transition| transition.transition_id.clone())
1252 .unwrap_or_else(|| sha256_hex(advance.canonical_proof_bytes()));
1253 let record = FiscalStagedTransitionRecord {
1254 transition_id,
1255 current_checkpoint_digest: advance.current().digest().to_owned(),
1256 next_checkpoint_digest: advance.next().digest().to_owned(),
1257 proof_json: advance.canonical_proof_bytes().to_vec(),
1258 status: FiscalStageStatus::DbStaged,
1259 stage_version: 1,
1260 };
1261 let checkpoint_json = advance.next().canonical_bytes()?;
1262 let mut connection = self.connection()?;
1263 let transaction = self.begin_write(&mut connection, fence)?;
1264 match load_transition(&transaction, &record.transition_id) {
1265 Ok(existing) => {
1266 let stored_checkpoint = transaction
1267 .query_row(
1268 "SELECT signed_json FROM fiscal_continuity_checkpoints WHERE checkpoint_digest = ?1",
1269 [&record.next_checkpoint_digest],
1270 |row| row.get::<_, Vec<u8>>(0),
1271 )
1272 .optional()
1273 .map_err(sqlite_error)?;
1274 if existing.current_checkpoint_digest != record.current_checkpoint_digest
1275 || existing.next_checkpoint_digest != record.next_checkpoint_digest
1276 || existing.proof_json != record.proof_json
1277 || stored_checkpoint.as_deref() != Some(checkpoint_json.as_slice())
1278 {
1279 return Err(FiscalStoreError::Conflict);
1280 }
1281 verify_exact_activation_mutation(
1282 &transaction,
1283 &record.transition_id,
1284 activation_mutation,
1285 )?;
1286 verify_exact_rotation_mutation(
1287 &transaction,
1288 &record.transition_id,
1289 rotation_mutation,
1290 )?;
1291 transaction.commit().map_err(sqlite_error)?;
1292 return Ok(existing);
1293 }
1294 Err(FiscalStoreError::NotFound) => {}
1295 Err(error) => return Err(error),
1296 }
1297 let current: String = transaction
1298 .query_row(
1299 "SELECT finalized_checkpoint_digest FROM fiscal_authority_state WHERE singleton = 1",
1300 [],
1301 |row| row.get(0),
1302 )
1303 .optional()
1304 .map_err(sqlite_error)?
1305 .ok_or(FiscalStoreError::NotFound)?;
1306 if current != record.current_checkpoint_digest {
1307 return Err(FiscalStoreError::Conflict);
1308 }
1309 let readiness_present = transaction
1310 .query_row(
1311 "SELECT EXISTS(SELECT 1 FROM fiscal_runtime_readiness WHERE readiness_digest = ?1)",
1312 [&advance.next().body().runtime_readiness_digest],
1313 |row| row.get::<_, bool>(0),
1314 )
1315 .map_err(sqlite_error)?;
1316 let charter_present = transaction
1317 .query_row(
1318 "SELECT EXISTS(SELECT 1 FROM fiscal_charters WHERE charter_id = ?1 AND charter_digest = ?2)",
1319 params![
1320 &advance.next().body().pinned_charter_id,
1321 &advance.next().body().pinned_charter_digest,
1322 ],
1323 |row| row.get::<_, bool>(0),
1324 )
1325 .map_err(sqlite_error)?;
1326 if !readiness_present || !charter_present {
1327 return Err(invariant(
1328 "staged fiscal checkpoint references an unavailable charter or readiness record",
1329 ));
1330 }
1331 for head in advance
1332 .next()
1333 .body()
1334 .domains
1335 .iter()
1336 .flat_map(|state| [state.active.as_ref(), state.last_known_good.as_ref()])
1337 .flatten()
1338 {
1339 let present = transaction
1340 .query_row(
1341 "SELECT EXISTS(SELECT 1 FROM fiscal_schedules WHERE schedule_id = ?1 AND schedule_digest = ?2 AND schedule_sequence = ?3)",
1342 params![
1343 &head.schedule_id,
1344 &head.schedule_digest,
1345 sqlite_i64(head.sequence, "schedule head sequence")?,
1346 ],
1347 |row| row.get::<_, bool>(0),
1348 )
1349 .map_err(sqlite_error)?;
1350 if !present {
1351 return Err(invariant(
1352 "staged fiscal checkpoint references an unavailable schedule",
1353 ));
1354 }
1355 }
1356 transaction
1357 .execute(
1358 "INSERT INTO fiscal_continuity_checkpoints (checkpoint_digest, continuity_sequence, status, signed_json) VALUES (?1, ?2, 'staged', ?3)",
1359 params![
1360 &record.next_checkpoint_digest,
1361 sqlite_i64(advance.next().body().continuity_sequence, "continuity sequence")?,
1362 &checkpoint_json,
1363 ],
1364 )
1365 .map_err(sqlite_error)?;
1366 transaction
1367 .execute(
1368 "INSERT INTO fiscal_staged_transitions (transition_id, current_checkpoint_digest, next_checkpoint_digest, proof_json, status, stage_version) VALUES (?1, ?2, ?3, ?4, ?5, 1)",
1369 params![
1370 &record.transition_id,
1371 &record.current_checkpoint_digest,
1372 &record.next_checkpoint_digest,
1373 &record.proof_json,
1374 record.status.as_str(),
1375 ],
1376 )
1377 .map_err(sqlite_error)?;
1378 if let Some(mutation) = activation_mutation {
1379 insert_activation_mutation(&transaction, mutation)?;
1380 }
1381 if let Some(mutation) = rotation_mutation {
1382 insert_rotation_mutation(&transaction, mutation)?;
1383 }
1384 let projection_key = format!("transition:{}", record.transition_id);
1385 append_projection_commit(
1386 &transaction,
1387 &self.serving_owner,
1388 &projection_key,
1389 1,
1390 "stage_fiscal_advance",
1391 &sha256_hex(&record.proof_json),
1392 )?;
1393 self.commit_write(transaction)?;
1394 self.sync_after_write(&connection)?;
1395 Ok(record)
1396 }
1397
1398 pub fn mark_anchor_advanced(
1399 &self,
1400 transition_id: &str,
1401 acknowledged: &VerifiedFiscalContinuityCheckpoint,
1402 fence: &StoreMutationFence,
1403 ) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
1404 self.transition_status_update(
1405 transition_id,
1406 FiscalStageStatus::DbStaged,
1407 FiscalStageStatus::FiscalAnchorAdvanced,
1408 Some(acknowledged),
1409 fence,
1410 )
1411 }
1412
1413 pub fn discard_unanchored_stage(
1414 &self,
1415 transition_id: &str,
1416 fence: &StoreMutationFence,
1417 ) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
1418 self.transition_status_update(
1419 transition_id,
1420 FiscalStageStatus::DbStaged,
1421 FiscalStageStatus::Discarded,
1422 None,
1423 fence,
1424 )
1425 }
1426
1427 pub fn finalize_advance(
1428 &self,
1429 transition_id: &str,
1430 acknowledged: &VerifiedFiscalContinuityCheckpoint,
1431 authority: &FiscalAuthorityState,
1432 fence: &StoreMutationFence,
1433 ) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
1434 authority.validate()?;
1435 if authority.finalized_checkpoint_digest != acknowledged.digest() {
1436 return Err(invariant(
1437 "finalized authority does not match anchor acknowledgement",
1438 ));
1439 }
1440 let acknowledged_json = acknowledged.canonical_bytes()?;
1441 let authority_json = canonical_json_bytes(authority).map_err(canonical_error)?;
1442 let mut connection = self.connection()?;
1443 let transaction = self.begin_write(&mut connection, fence)?;
1444 let mut record = load_transition(&transaction, transition_id)?;
1445 verify_acknowledgement(&transaction, &record, acknowledged, &acknowledged_json)?;
1446 if record.status == FiscalStageStatus::DbFinalized {
1447 let stored: Vec<u8> = transaction
1448 .query_row(
1449 "SELECT state_json FROM fiscal_authority_state WHERE singleton = 1",
1450 [],
1451 |row| row.get(0),
1452 )
1453 .map_err(sqlite_error)?;
1454 if stored != authority_json {
1455 return Err(FiscalStoreError::Conflict);
1456 }
1457 transaction.commit().map_err(sqlite_error)?;
1458 return Ok(record);
1459 }
1460 if record.status != FiscalStageStatus::FiscalAnchorAdvanced {
1461 return Err(FiscalStoreError::Conflict);
1462 }
1463 let next_version = record
1464 .stage_version
1465 .checked_add(1)
1466 .ok_or_else(|| invariant("fiscal stage version overflowed"))?;
1467 let updated = transaction
1468 .execute(
1469 "UPDATE fiscal_staged_transitions SET status = 'db_finalized', stage_version = ?1 WHERE transition_id = ?2 AND status = 'fiscal_anchor_advanced' AND stage_version = ?3",
1470 params![
1471 sqlite_i64(next_version, "fiscal stage version")?,
1472 transition_id,
1473 sqlite_i64(record.stage_version, "fiscal stage version")?,
1474 ],
1475 )
1476 .map_err(sqlite_error)?;
1477 let checkpoint_updated = transaction
1478 .execute(
1479 "UPDATE fiscal_continuity_checkpoints SET status = 'finalized' WHERE checkpoint_digest = ?1 AND status = 'staged' AND signed_json = ?2",
1480 params![acknowledged.digest(), &acknowledged_json],
1481 )
1482 .map_err(sqlite_error)?;
1483 let authority_updated = transaction
1484 .execute(
1485 "UPDATE fiscal_authority_state SET state_json = ?1, finalized_checkpoint_digest = ?2, state_version = state_version + 1 WHERE singleton = 1 AND finalized_checkpoint_digest = ?3 AND state_version < 9223372036854775807",
1486 params![
1487 &authority_json,
1488 acknowledged.digest(),
1489 &record.current_checkpoint_digest,
1490 ],
1491 )
1492 .map_err(sqlite_error)?;
1493 if updated != 1 || checkpoint_updated != 1 || authority_updated != 1 {
1494 return Err(FiscalStoreError::Conflict);
1495 }
1496 finalize_activation_mutation(&transaction, &self.serving_owner, transition_id)?;
1497 finalize_rotation_mutation(&transaction, &self.serving_owner, transition_id)?;
1498 record.status = FiscalStageStatus::DbFinalized;
1499 record.stage_version = next_version;
1500 let projection_key = format!("transition:{transition_id}");
1501 append_projection_commit(
1502 &transaction,
1503 &self.serving_owner,
1504 &projection_key,
1505 next_version,
1506 "finalize_fiscal_advance",
1507 acknowledged.digest(),
1508 )?;
1509 self.commit_write(transaction)?;
1510 self.sync_after_write(&connection)?;
1511 Ok(record)
1512 }
1513
1514 pub fn load_authority_state(&self) -> Result<FiscalAuthorityState, FiscalStoreError> {
1515 let mut connection = self.connection()?;
1516 let transaction = self.begin_read(&mut connection)?;
1517 let json = transaction
1518 .query_row(
1519 "SELECT state_json FROM fiscal_authority_state WHERE singleton = 1",
1520 [],
1521 |row| row.get::<_, Vec<u8>>(0),
1522 )
1523 .optional()
1524 .map_err(sqlite_error)?
1525 .ok_or(FiscalStoreError::NotFound)?;
1526 let state: FiscalAuthorityState = serde_json::from_slice(&json)
1527 .map_err(|error| invariant(format!("stored fiscal authority is invalid: {error}")))?;
1528 state.validate()?;
1529 if canonical_json_bytes(&state).map_err(canonical_error)? != json {
1530 return Err(invariant("stored fiscal authority is not canonical"));
1531 }
1532 transaction.commit().map_err(sqlite_error)?;
1533 Ok(state)
1534 }
1535
1536 pub fn load_transition(
1537 &self,
1538 transition_id: &str,
1539 ) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
1540 let mut connection = self.connection()?;
1541 let transaction = self.begin_read(&mut connection)?;
1542 let record = load_transition(&transaction, transition_id)?;
1543 transaction.commit().map_err(sqlite_error)?;
1544 Ok(record)
1545 }
1546
1547 pub fn load_open_transition(
1548 &self,
1549 ) -> Result<Option<FiscalStagedTransitionRecord>, FiscalStoreError> {
1550 let mut connection = self.connection()?;
1551 let transaction = self.begin_read(&mut connection)?;
1552 let transition_id = transaction
1553 .query_row(
1554 "SELECT transition_id FROM fiscal_staged_transitions WHERE status IN ('db_staged', 'fiscal_anchor_advanced')",
1555 [],
1556 |row| row.get::<_, String>(0),
1557 )
1558 .optional()
1559 .map_err(sqlite_error)?;
1560 let record = transition_id
1561 .as_deref()
1562 .map(|id| load_transition(&transaction, id))
1563 .transpose()?;
1564 transaction.commit().map_err(sqlite_error)?;
1565 Ok(record)
1566 }
1567
1568 pub fn load_checkpoint(
1569 &self,
1570 checkpoint_digest: &str,
1571 policy: &FiscalGenesisPolicy,
1572 charters: &FiscalCharterRegistry,
1573 ) -> Result<VerifiedFiscalContinuityCheckpoint, FiscalStoreError> {
1574 let mut connection = self.connection()?;
1575 let transaction = self.begin_read(&mut connection)?;
1576 let json = transaction
1577 .query_row(
1578 "SELECT signed_json FROM fiscal_continuity_checkpoints WHERE checkpoint_digest = ?1",
1579 [checkpoint_digest],
1580 |row| row.get::<_, Vec<u8>>(0),
1581 )
1582 .optional()
1583 .map_err(sqlite_error)?
1584 .ok_or(FiscalStoreError::NotFound)?;
1585 let checkpoint =
1586 VerifiedFiscalContinuityCheckpoint::from_canonical_bytes(&json, policy, charters)?;
1587 if checkpoint.digest() != checkpoint_digest {
1588 return Err(invariant("stored fiscal checkpoint digest is inconsistent"));
1589 }
1590 transaction.commit().map_err(sqlite_error)?;
1591 Ok(checkpoint)
1592 }
1593
1594 pub fn load_finalized_checkpoints(
1595 &self,
1596 policy: &FiscalGenesisPolicy,
1597 charters: &FiscalCharterRegistry,
1598 ) -> Result<Vec<VerifiedFiscalContinuityCheckpoint>, FiscalStoreError> {
1599 let mut connection = self.connection()?;
1600 let transaction = self.begin_read(&mut connection)?;
1601 let mut statement = transaction
1602 .prepare(
1603 "SELECT signed_json FROM fiscal_continuity_checkpoints WHERE status = 'finalized' ORDER BY continuity_sequence",
1604 )
1605 .map_err(sqlite_error)?;
1606 let checkpoints = statement
1607 .query_map([], |row| row.get::<_, Vec<u8>>(0))
1608 .map_err(sqlite_error)?
1609 .map(|row| {
1610 let bytes = row.map_err(sqlite_error)?;
1611 VerifiedFiscalContinuityCheckpoint::from_canonical_bytes(&bytes, policy, charters)
1612 .map_err(FiscalStoreError::from)
1613 })
1614 .collect::<Result<Vec<_>, FiscalStoreError>>()?;
1615 drop(statement);
1616 transaction.commit().map_err(sqlite_error)?;
1617 Ok(checkpoints)
1618 }
1619
1620 fn transition_status_update(
1621 &self,
1622 transition_id: &str,
1623 expected: FiscalStageStatus,
1624 next: FiscalStageStatus,
1625 acknowledged: Option<&VerifiedFiscalContinuityCheckpoint>,
1626 fence: &StoreMutationFence,
1627 ) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
1628 let acknowledged_json = acknowledged
1629 .map(VerifiedFiscalContinuityCheckpoint::canonical_bytes)
1630 .transpose()?;
1631 let mut connection = self.connection()?;
1632 let transaction = self.begin_write(&mut connection, fence)?;
1633 let mut record = load_transition(&transaction, transition_id)?;
1634 if let (Some(checkpoint), Some(json)) = (acknowledged, acknowledged_json.as_deref()) {
1635 verify_acknowledgement(&transaction, &record, checkpoint, json)?;
1636 }
1637 if record.status == next {
1638 transaction.commit().map_err(sqlite_error)?;
1639 return Ok(record);
1640 }
1641 if record.status != expected {
1642 return Err(FiscalStoreError::Conflict);
1643 }
1644 let next_version = record
1645 .stage_version
1646 .checked_add(1)
1647 .ok_or_else(|| invariant("fiscal stage version overflowed"))?;
1648 let updated = transaction
1649 .execute(
1650 "UPDATE fiscal_staged_transitions SET status = ?1, stage_version = ?2 WHERE transition_id = ?3 AND status = ?4 AND stage_version = ?5",
1651 params![
1652 next.as_str(),
1653 sqlite_i64(next_version, "fiscal stage version")?,
1654 transition_id,
1655 expected.as_str(),
1656 sqlite_i64(record.stage_version, "fiscal stage version")?,
1657 ],
1658 )
1659 .map_err(sqlite_error)?;
1660 if updated != 1 {
1661 return Err(FiscalStoreError::Conflict);
1662 }
1663 record.status = next;
1664 record.stage_version = next_version;
1665 let projection_key = format!("transition:{transition_id}");
1666 append_projection_commit(
1667 &transaction,
1668 &self.serving_owner,
1669 &projection_key,
1670 next_version,
1671 "advance_fiscal_stage",
1672 &record.next_checkpoint_digest,
1673 )?;
1674 self.commit_write(transaction)?;
1675 self.sync_after_write(&connection)?;
1676 Ok(record)
1677 }
1678}
1679
1680fn prepare_activation_mutation(
1681 advance: &VerifiedFiscalContinuityAdvance,
1682 activation: &VerifiedFiscalActivation,
1683 activated_admission: &FiscalProposalAdmissionState,
1684 candidate: &VerifiedFiscalSchedule,
1685 predecessor: Option<&VerifiedFiscalSchedule>,
1686) -> Result<PreparedFiscalActivationMutation, FiscalStoreError> {
1687 let FiscalActivationTarget::Schedule {
1688 schedule_id,
1689 supersedes_schedule_id,
1690 } = &activation.body().target
1691 else {
1692 return Err(FiscalStoreError::Conflict);
1693 };
1694 let transition = FiscalStagedTransition::new(
1695 activation.body().activation_id.clone(),
1696 activation.digest().to_owned(),
1697 )?;
1698 let candidate_head = FiscalScheduleHead::from_signed(candidate.signed())?;
1699 let predecessor_head = predecessor
1700 .map(|schedule| FiscalScheduleHead::from_signed(schedule.signed()))
1701 .transpose()?;
1702 let next_domain = advance
1703 .next()
1704 .body()
1705 .domains
1706 .iter()
1707 .find(|state| state.domain == candidate.body().domain)
1708 .ok_or(FiscalStoreError::Conflict)?;
1709 let expected_admission_version = activated_admission
1710 .version
1711 .checked_sub(1)
1712 .ok_or(FiscalStoreError::Conflict)?;
1713 if advance.next().body().staged_transition.as_ref() != Some(&transition)
1714 || schedule_id != &candidate.body().schedule_id
1715 || next_domain.active.as_ref() != Some(&candidate_head)
1716 || next_domain.last_known_good.as_ref() != Some(&candidate_head)
1717 || activated_admission.status != FiscalProposalAdmissionStatus::Activated
1718 || activated_admission.signed_admission.body.admission_id != activation.body().admission_id
1719 || activated_admission.admission_digest != activation.body().admission_digest
1720 || activated_admission.activation_digest.as_deref() != Some(activation.digest())
1721 || activated_admission.activated_sequence != Some(candidate.body().sequence)
1722 || expected_admission_version == 0
1723 || supersedes_schedule_id.as_deref()
1724 != predecessor.map(|schedule| schedule.body().schedule_id.as_str())
1725 || candidate.body().supersedes_schedule_id.as_deref()
1726 != predecessor.map(|schedule| schedule.body().schedule_id.as_str())
1727 {
1728 return Err(FiscalStoreError::Conflict);
1729 }
1730 let expected_admission = FiscalProposalAdmissionState {
1731 signed_admission: activated_admission.signed_admission.clone(),
1732 admission_digest: activated_admission.admission_digest.clone(),
1733 version: expected_admission_version,
1734 status: FiscalProposalAdmissionStatus::Admitted,
1735 activation_digest: None,
1736 activated_sequence: None,
1737 };
1738 Ok(PreparedFiscalActivationMutation {
1739 transition_id: transition.transition_id,
1740 activation_id: activation.body().activation_id.clone(),
1741 activation_digest: activation.digest().to_owned(),
1742 admission_id: activated_admission
1743 .signed_admission
1744 .body
1745 .admission_id
1746 .clone(),
1747 admission_digest: activated_admission.admission_digest.clone(),
1748 expected_admission_version,
1749 expected_admission_json: canonical_json_bytes(&expected_admission)
1750 .map_err(canonical_error)?,
1751 activated_admission_version: activated_admission.version,
1752 activated_admission_json: canonical_json_bytes(activated_admission)
1753 .map_err(canonical_error)?,
1754 candidate_schedule_id: candidate.body().schedule_id.clone(),
1755 candidate_schedule_digest: candidate_head.schedule_digest,
1756 predecessor_schedule_id: predecessor.map(|schedule| schedule.body().schedule_id.clone()),
1757 predecessor_schedule_digest: predecessor_head.map(|head| head.schedule_digest),
1758 })
1759}
1760
1761#[allow(clippy::too_many_arguments)]
1762fn prepare_rotation_mutation(
1763 advance: &VerifiedFiscalContinuityAdvance,
1764 activation: &VerifiedFiscalActivation,
1765 activated_admission: &FiscalProposalAdmissionState,
1766 successor_charter: &VerifiedFiscalCharter,
1767 successor_schedules: &[VerifiedFiscalSchedule],
1768 predecessor_charter: &VerifiedFiscalCharter,
1769 predecessor_schedules: &[VerifiedFiscalSchedule],
1770) -> Result<PreparedFiscalRotationMutation, FiscalStoreError> {
1771 let FiscalActivationTarget::CharterRotation {
1772 successor_charter_digest,
1773 predecessor_charter_digest,
1774 successor_schedules: signed_successors,
1775 } = &activation.body().target
1776 else {
1777 return Err(FiscalStoreError::Conflict);
1778 };
1779 let transition = FiscalStagedTransition::new(
1780 activation.body().activation_id.clone(),
1781 activation.digest().to_owned(),
1782 )?;
1783 let expected_admission_version = activated_admission
1784 .version
1785 .checked_sub(1)
1786 .ok_or(FiscalStoreError::Conflict)?;
1787 let activated_sequence = successor_schedules
1788 .iter()
1789 .map(|schedule| schedule.body().sequence)
1790 .max()
1791 .ok_or(FiscalStoreError::Conflict)?;
1792 if advance.next().body().staged_transition.as_ref() != Some(&transition)
1793 || advance.next().body().pinned_charter_id != successor_charter.body().charter_id
1794 || advance.next().body().pinned_charter_digest != successor_charter.digest()
1795 || advance.next().body().pinned_charter_sequence != successor_charter.body().sequence
1796 || successor_charter_digest != successor_charter.digest()
1797 || predecessor_charter_digest != predecessor_charter.digest()
1798 || successor_charter
1799 .body()
1800 .predecessor_charter_digest
1801 .as_deref()
1802 != Some(predecessor_charter.digest())
1803 || successor_schedules.len() != predecessor_schedules.len()
1804 || signed_successors.len() != successor_schedules.len()
1805 || signed_successors
1806 .iter()
1807 .zip(successor_schedules)
1808 .any(|(signed, verified)| signed != verified.signed())
1809 || activated_admission.status != FiscalProposalAdmissionStatus::Activated
1810 || activated_admission.signed_admission.body.admission_id != activation.body().admission_id
1811 || activated_admission.admission_digest != activation.body().admission_digest
1812 || activated_admission.activation_digest.as_deref() != Some(activation.digest())
1813 || activated_admission.activated_sequence != Some(activated_sequence)
1814 || expected_admission_version == 0
1815 {
1816 return Err(FiscalStoreError::Conflict);
1817 }
1818 let mut schedules = Vec::with_capacity(successor_schedules.len());
1819 for (candidate, predecessor) in successor_schedules.iter().zip(predecessor_schedules) {
1820 let candidate_head = FiscalScheduleHead::from_signed(candidate.signed())?;
1821 let predecessor_head = FiscalScheduleHead::from_signed(predecessor.signed())?;
1822 let next_domain = advance
1823 .next()
1824 .body()
1825 .domains
1826 .iter()
1827 .find(|state| state.domain == candidate.body().domain)
1828 .ok_or(FiscalStoreError::Conflict)?;
1829 if candidate.body().domain != predecessor.body().domain
1830 || candidate.body().supersedes_schedule_id.as_deref()
1831 != Some(predecessor.body().schedule_id.as_str())
1832 || next_domain.active.as_ref() != Some(&candidate_head)
1833 || next_domain.last_known_good.as_ref() != Some(&candidate_head)
1834 {
1835 return Err(FiscalStoreError::Conflict);
1836 }
1837 schedules.push(PreparedFiscalRotationScheduleMutation {
1838 domain: candidate.body().domain,
1839 candidate_schedule_id: candidate.body().schedule_id.clone(),
1840 candidate_schedule_digest: candidate_head.schedule_digest,
1841 predecessor_schedule_id: predecessor.body().schedule_id.clone(),
1842 predecessor_schedule_digest: predecessor_head.schedule_digest,
1843 });
1844 }
1845 if schedules
1846 .windows(2)
1847 .any(|pair| pair[0].domain >= pair[1].domain)
1848 {
1849 return Err(FiscalStoreError::Conflict);
1850 }
1851 let expected_admission = FiscalProposalAdmissionState {
1852 signed_admission: activated_admission.signed_admission.clone(),
1853 admission_digest: activated_admission.admission_digest.clone(),
1854 version: expected_admission_version,
1855 status: FiscalProposalAdmissionStatus::Admitted,
1856 activation_digest: None,
1857 activated_sequence: None,
1858 };
1859 Ok(PreparedFiscalRotationMutation {
1860 transition_id: transition.transition_id,
1861 activation_id: activation.body().activation_id.clone(),
1862 activation_digest: activation.digest().to_owned(),
1863 admission_id: activated_admission
1864 .signed_admission
1865 .body
1866 .admission_id
1867 .clone(),
1868 admission_digest: activated_admission.admission_digest.clone(),
1869 expected_admission_version,
1870 expected_admission_json: canonical_json_bytes(&expected_admission)
1871 .map_err(canonical_error)?,
1872 activated_admission_version: activated_admission.version,
1873 activated_admission_json: canonical_json_bytes(activated_admission)
1874 .map_err(canonical_error)?,
1875 successor_charter_id: successor_charter.body().charter_id.clone(),
1876 successor_charter_digest: successor_charter.digest().to_owned(),
1877 predecessor_charter_id: predecessor_charter.body().charter_id.clone(),
1878 predecessor_charter_digest: predecessor_charter.digest().to_owned(),
1879 schedules,
1880 })
1881}
1882
1883fn insert_activation_mutation(
1884 transaction: &Transaction<'_>,
1885 mutation: &PreparedFiscalActivationMutation,
1886) -> Result<(), FiscalStoreError> {
1887 let activation_exact = transaction
1888 .query_row(
1889 "SELECT activation_digest = ?1 FROM fiscal_activations WHERE activation_id = ?2",
1890 params![&mutation.activation_digest, &mutation.activation_id],
1891 |row| row.get::<_, bool>(0),
1892 )
1893 .optional()
1894 .map_err(sqlite_error)?
1895 .unwrap_or(false);
1896 let admission_exact = transaction
1897 .query_row(
1898 "SELECT admission_digest = ?1 AND state_version = ?2 AND status = 'admitted' AND state_json = ?3 FROM fiscal_proposal_admissions WHERE admission_id = ?4",
1899 params![
1900 &mutation.admission_digest,
1901 sqlite_i64(mutation.expected_admission_version, "expected admission version")?,
1902 &mutation.expected_admission_json,
1903 &mutation.admission_id,
1904 ],
1905 |row| row.get::<_, bool>(0),
1906 )
1907 .optional()
1908 .map_err(sqlite_error)?
1909 .unwrap_or(false);
1910 let candidate_exact = schedule_has_state(
1911 transaction,
1912 &mutation.candidate_schedule_id,
1913 &mutation.candidate_schedule_digest,
1914 "staged",
1915 )?;
1916 let predecessor_exact = match (
1917 &mutation.predecessor_schedule_id,
1918 &mutation.predecessor_schedule_digest,
1919 ) {
1920 (Some(id), Some(digest)) => schedule_has_state(transaction, id, digest, "active")?,
1921 (None, None) => true,
1922 _ => false,
1923 };
1924 if !activation_exact || !admission_exact || !candidate_exact || !predecessor_exact {
1925 return Err(FiscalStoreError::Conflict);
1926 }
1927 transaction
1928 .execute(
1929 "INSERT INTO fiscal_staged_activation_mutations (transition_id, activation_id, activation_digest, admission_id, admission_digest, expected_admission_version, activated_admission_version, activated_admission_json, candidate_schedule_id, candidate_schedule_digest, predecessor_schedule_id, predecessor_schedule_digest) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
1930 params![
1931 &mutation.transition_id,
1932 &mutation.activation_id,
1933 &mutation.activation_digest,
1934 &mutation.admission_id,
1935 &mutation.admission_digest,
1936 sqlite_i64(mutation.expected_admission_version, "expected admission version")?,
1937 sqlite_i64(mutation.activated_admission_version, "activated admission version")?,
1938 &mutation.activated_admission_json,
1939 &mutation.candidate_schedule_id,
1940 &mutation.candidate_schedule_digest,
1941 &mutation.predecessor_schedule_id,
1942 &mutation.predecessor_schedule_digest,
1943 ],
1944 )
1945 .map_err(|error| invariant(format!("staged activation mutation insert failed: {error}")))?;
1946 Ok(())
1947}
1948
1949fn verify_exact_activation_mutation(
1950 transaction: &Transaction<'_>,
1951 transition_id: &str,
1952 expected: Option<&PreparedFiscalActivationMutation>,
1953) -> Result<(), FiscalStoreError> {
1954 let retained = transaction
1955 .query_row(
1956 "SELECT activation_id, activation_digest, admission_id, admission_digest, expected_admission_version, activated_admission_version, activated_admission_json, candidate_schedule_id, candidate_schedule_digest, predecessor_schedule_id, predecessor_schedule_digest FROM fiscal_staged_activation_mutations WHERE transition_id = ?1",
1957 [transition_id],
1958 |row| {
1959 Ok((
1960 row.get::<_, String>(0)?,
1961 row.get::<_, String>(1)?,
1962 row.get::<_, String>(2)?,
1963 row.get::<_, String>(3)?,
1964 row.get::<_, i64>(4)?,
1965 row.get::<_, i64>(5)?,
1966 row.get::<_, Vec<u8>>(6)?,
1967 row.get::<_, String>(7)?,
1968 row.get::<_, String>(8)?,
1969 row.get::<_, Option<String>>(9)?,
1970 row.get::<_, Option<String>>(10)?,
1971 ))
1972 },
1973 )
1974 .optional()
1975 .map_err(sqlite_error)?;
1976 match (retained, expected) {
1977 (None, None) => Ok(()),
1978 (Some(row), Some(expected))
1979 if row.0 == expected.activation_id
1980 && row.1 == expected.activation_digest
1981 && row.2 == expected.admission_id
1982 && row.3 == expected.admission_digest
1983 && read_u64(row.4, "expected admission version")?
1984 == expected.expected_admission_version
1985 && read_u64(row.5, "activated admission version")?
1986 == expected.activated_admission_version
1987 && row.6 == expected.activated_admission_json
1988 && row.7 == expected.candidate_schedule_id
1989 && row.8 == expected.candidate_schedule_digest
1990 && row.9 == expected.predecessor_schedule_id
1991 && row.10 == expected.predecessor_schedule_digest =>
1992 {
1993 Ok(())
1994 }
1995 _ => Err(FiscalStoreError::Conflict),
1996 }
1997}
1998
1999fn insert_rotation_mutation(
2000 transaction: &Transaction<'_>,
2001 mutation: &PreparedFiscalRotationMutation,
2002) -> Result<(), FiscalStoreError> {
2003 let activation_exact = transaction
2004 .query_row(
2005 "SELECT activation_digest = ?1 FROM fiscal_activations WHERE activation_id = ?2",
2006 params![&mutation.activation_digest, &mutation.activation_id],
2007 |row| row.get::<_, bool>(0),
2008 )
2009 .optional()
2010 .map_err(sqlite_error)?
2011 .unwrap_or(false);
2012 let admission_exact = transaction
2013 .query_row(
2014 "SELECT admission_digest = ?1 AND state_version = ?2 AND status = 'admitted' AND state_json = ?3 FROM fiscal_proposal_admissions WHERE admission_id = ?4",
2015 params![
2016 &mutation.admission_digest,
2017 sqlite_i64(mutation.expected_admission_version, "expected admission version")?,
2018 &mutation.expected_admission_json,
2019 &mutation.admission_id,
2020 ],
2021 |row| row.get::<_, bool>(0),
2022 )
2023 .optional()
2024 .map_err(sqlite_error)?
2025 .unwrap_or(false);
2026 let successor_exact = charter_has_state(
2027 transaction,
2028 &mutation.successor_charter_id,
2029 &mutation.successor_charter_digest,
2030 "proposed",
2031 )?;
2032 let predecessor_exact = charter_has_current_state(
2033 transaction,
2034 &mutation.predecessor_charter_id,
2035 &mutation.predecessor_charter_digest,
2036 )?;
2037 if !activation_exact || !admission_exact || !successor_exact || !predecessor_exact {
2038 return Err(FiscalStoreError::Conflict);
2039 }
2040 for schedule in &mutation.schedules {
2041 if !schedule_has_state(
2042 transaction,
2043 &schedule.candidate_schedule_id,
2044 &schedule.candidate_schedule_digest,
2045 "staged",
2046 )? || !schedule_has_state(
2047 transaction,
2048 &schedule.predecessor_schedule_id,
2049 &schedule.predecessor_schedule_digest,
2050 "active",
2051 )? {
2052 return Err(FiscalStoreError::Conflict);
2053 }
2054 }
2055 transaction
2056 .execute(
2057 "INSERT INTO fiscal_staged_rotation_mutations (transition_id, activation_id, activation_digest, admission_id, admission_digest, expected_admission_version, activated_admission_version, activated_admission_json, successor_charter_id, successor_charter_digest, predecessor_charter_id, predecessor_charter_digest) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
2058 params![
2059 &mutation.transition_id,
2060 &mutation.activation_id,
2061 &mutation.activation_digest,
2062 &mutation.admission_id,
2063 &mutation.admission_digest,
2064 sqlite_i64(mutation.expected_admission_version, "expected admission version")?,
2065 sqlite_i64(mutation.activated_admission_version, "activated admission version")?,
2066 &mutation.activated_admission_json,
2067 &mutation.successor_charter_id,
2068 &mutation.successor_charter_digest,
2069 &mutation.predecessor_charter_id,
2070 &mutation.predecessor_charter_digest,
2071 ],
2072 )
2073 .map_err(sqlite_error)?;
2074 for schedule in &mutation.schedules {
2075 let domain_json =
2076 String::from_utf8(canonical_json_bytes(&schedule.domain).map_err(canonical_error)?)
2077 .map_err(|error| {
2078 invariant(format!("fiscal domain encoding is not UTF-8: {error}"))
2079 })?;
2080 transaction
2081 .execute(
2082 "INSERT INTO fiscal_staged_rotation_schedules (transition_id, domain_json, candidate_schedule_id, candidate_schedule_digest, predecessor_schedule_id, predecessor_schedule_digest) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
2083 params![
2084 &mutation.transition_id,
2085 domain_json,
2086 &schedule.candidate_schedule_id,
2087 &schedule.candidate_schedule_digest,
2088 &schedule.predecessor_schedule_id,
2089 &schedule.predecessor_schedule_digest,
2090 ],
2091 )
2092 .map_err(sqlite_error)?;
2093 }
2094 Ok(())
2095}
2096
2097fn verify_exact_rotation_mutation(
2098 transaction: &Transaction<'_>,
2099 transition_id: &str,
2100 expected: Option<&PreparedFiscalRotationMutation>,
2101) -> Result<(), FiscalStoreError> {
2102 let retained = transaction
2103 .query_row(
2104 "SELECT activation_id, activation_digest, admission_id, admission_digest, expected_admission_version, activated_admission_version, activated_admission_json, successor_charter_id, successor_charter_digest, predecessor_charter_id, predecessor_charter_digest FROM fiscal_staged_rotation_mutations WHERE transition_id = ?1",
2105 [transition_id],
2106 |row| {
2107 Ok((
2108 row.get::<_, String>(0)?,
2109 row.get::<_, String>(1)?,
2110 row.get::<_, String>(2)?,
2111 row.get::<_, String>(3)?,
2112 row.get::<_, i64>(4)?,
2113 row.get::<_, i64>(5)?,
2114 row.get::<_, Vec<u8>>(6)?,
2115 row.get::<_, String>(7)?,
2116 row.get::<_, String>(8)?,
2117 row.get::<_, String>(9)?,
2118 row.get::<_, String>(10)?,
2119 ))
2120 },
2121 )
2122 .optional()
2123 .map_err(sqlite_error)?;
2124 match (retained, expected) {
2125 (None, None) => Ok(()),
2126 (Some(row), Some(expected))
2127 if row.0 == expected.activation_id
2128 && row.1 == expected.activation_digest
2129 && row.2 == expected.admission_id
2130 && row.3 == expected.admission_digest
2131 && read_u64(row.4, "expected admission version")?
2132 == expected.expected_admission_version
2133 && read_u64(row.5, "activated admission version")?
2134 == expected.activated_admission_version
2135 && row.6 == expected.activated_admission_json
2136 && row.7 == expected.successor_charter_id
2137 && row.8 == expected.successor_charter_digest
2138 && row.9 == expected.predecessor_charter_id
2139 && row.10 == expected.predecessor_charter_digest =>
2140 {
2141 let mut statement = transaction
2142 .prepare(
2143 "SELECT domain_json, candidate_schedule_id, candidate_schedule_digest, predecessor_schedule_id, predecessor_schedule_digest FROM fiscal_staged_rotation_schedules WHERE transition_id = ?1 ORDER BY rowid",
2144 )
2145 .map_err(sqlite_error)?;
2146 let retained_schedules = statement
2147 .query_map([transition_id], |row| {
2148 Ok((
2149 row.get::<_, String>(0)?,
2150 row.get::<_, String>(1)?,
2151 row.get::<_, String>(2)?,
2152 row.get::<_, String>(3)?,
2153 row.get::<_, String>(4)?,
2154 ))
2155 })
2156 .map_err(sqlite_error)?
2157 .collect::<Result<Vec<_>, _>>()
2158 .map_err(sqlite_error)?;
2159 let expected_schedules = expected
2160 .schedules
2161 .iter()
2162 .map(|schedule| {
2163 let domain = canonical_json_bytes(&schedule.domain).map_err(canonical_error)?;
2164 let domain = String::from_utf8(domain).map_err(|error| {
2165 invariant(format!("fiscal domain encoding is not UTF-8: {error}"))
2166 })?;
2167 Ok((
2168 domain,
2169 schedule.candidate_schedule_id.clone(),
2170 schedule.candidate_schedule_digest.clone(),
2171 schedule.predecessor_schedule_id.clone(),
2172 schedule.predecessor_schedule_digest.clone(),
2173 ))
2174 })
2175 .collect::<Result<Vec<_>, FiscalStoreError>>()?;
2176 if retained_schedules == expected_schedules {
2177 Ok(())
2178 } else {
2179 Err(FiscalStoreError::Conflict)
2180 }
2181 }
2182 _ => Err(FiscalStoreError::Conflict),
2183 }
2184}
2185
2186fn charter_has_state(
2187 transaction: &Transaction<'_>,
2188 charter_id: &str,
2189 charter_digest: &str,
2190 state: &str,
2191) -> Result<bool, FiscalStoreError> {
2192 transaction
2193 .query_row(
2194 "SELECT EXISTS(SELECT 1 FROM fiscal_charters WHERE charter_id = ?1 AND charter_digest = ?2 AND lifecycle_state = ?3)",
2195 params![charter_id, charter_digest, state],
2196 |row| row.get(0),
2197 )
2198 .map_err(sqlite_error)
2199}
2200
2201fn charter_has_current_state(
2202 transaction: &Transaction<'_>,
2203 charter_id: &str,
2204 charter_digest: &str,
2205) -> Result<bool, FiscalStoreError> {
2206 transaction
2207 .query_row(
2208 "SELECT EXISTS(SELECT 1 FROM fiscal_charters WHERE charter_id = ?1 AND charter_digest = ?2 AND lifecycle_state IN ('pinned', 'active'))",
2209 params![charter_id, charter_digest],
2210 |row| row.get(0),
2211 )
2212 .map_err(sqlite_error)
2213}
2214
2215fn schedule_has_state(
2216 transaction: &Transaction<'_>,
2217 schedule_id: &str,
2218 schedule_digest: &str,
2219 state: &str,
2220) -> Result<bool, FiscalStoreError> {
2221 transaction
2222 .query_row(
2223 "SELECT EXISTS(SELECT 1 FROM fiscal_schedules WHERE schedule_id = ?1 AND schedule_digest = ?2 AND lifecycle_state = ?3)",
2224 params![schedule_id, schedule_digest, state],
2225 |row| row.get(0),
2226 )
2227 .map_err(sqlite_error)
2228}
2229
2230fn finalize_activation_mutation(
2231 transaction: &Transaction<'_>,
2232 owner: &SqliteServingOwner,
2233 transition_id: &str,
2234) -> Result<(), FiscalStoreError> {
2235 let mutation = transaction
2236 .query_row(
2237 "SELECT admission_id, admission_digest, expected_admission_version, activated_admission_version, activated_admission_json, candidate_schedule_id, candidate_schedule_digest, predecessor_schedule_id, predecessor_schedule_digest FROM fiscal_staged_activation_mutations WHERE transition_id = ?1",
2238 [transition_id],
2239 |row| {
2240 Ok((
2241 row.get::<_, String>(0)?,
2242 row.get::<_, String>(1)?,
2243 row.get::<_, i64>(2)?,
2244 row.get::<_, i64>(3)?,
2245 row.get::<_, Vec<u8>>(4)?,
2246 row.get::<_, String>(5)?,
2247 row.get::<_, String>(6)?,
2248 row.get::<_, Option<String>>(7)?,
2249 row.get::<_, Option<String>>(8)?,
2250 ))
2251 },
2252 )
2253 .optional()
2254 .map_err(sqlite_error)?;
2255 let Some(mutation) = mutation else {
2256 return Ok(());
2257 };
2258 let expected_version = read_u64(mutation.2, "expected admission version")?;
2259 let activated_version = read_u64(mutation.3, "activated admission version")?;
2260 let admission_updated = transaction
2261 .execute(
2262 "UPDATE fiscal_proposal_admissions SET status = 'activated', state_version = ?1, state_json = ?2 WHERE admission_id = ?3 AND admission_digest = ?4 AND state_version = ?5 AND status = 'admitted'",
2263 params![
2264 sqlite_i64(activated_version, "activated admission version")?,
2265 &mutation.4,
2266 &mutation.0,
2267 &mutation.1,
2268 sqlite_i64(expected_version, "expected admission version")?,
2269 ],
2270 )
2271 .map_err(sqlite_error)?;
2272 let candidate_updated = transaction
2273 .execute(
2274 "UPDATE fiscal_schedules SET lifecycle_state = 'active' WHERE schedule_id = ?1 AND schedule_digest = ?2 AND lifecycle_state = 'staged'",
2275 params![&mutation.5, &mutation.6],
2276 )
2277 .map_err(sqlite_error)?;
2278 let predecessor_updated = match (&mutation.7, &mutation.8) {
2279 (Some(id), Some(digest)) => transaction
2280 .execute(
2281 "UPDATE fiscal_schedules SET lifecycle_state = 'superseded' WHERE schedule_id = ?1 AND schedule_digest = ?2 AND lifecycle_state = 'active'",
2282 params![id, digest],
2283 )
2284 .map_err(sqlite_error)?,
2285 (None, None) => 1,
2286 _ => 0,
2287 };
2288 if admission_updated != 1 || candidate_updated != 1 || predecessor_updated != 1 {
2289 return Err(FiscalStoreError::Conflict);
2290 }
2291 append_projection_commit(
2292 transaction,
2293 owner,
2294 &format!("admission:{}", mutation.0),
2295 activated_version,
2296 "activate_fiscal_admission",
2297 &sha256_hex(&mutation.4),
2298 )
2299}
2300
2301fn finalize_rotation_mutation(
2302 transaction: &Transaction<'_>,
2303 owner: &SqliteServingOwner,
2304 transition_id: &str,
2305) -> Result<(), FiscalStoreError> {
2306 let mutation = transaction
2307 .query_row(
2308 "SELECT admission_id, admission_digest, expected_admission_version, activated_admission_version, activated_admission_json, successor_charter_id, successor_charter_digest, predecessor_charter_id, predecessor_charter_digest FROM fiscal_staged_rotation_mutations WHERE transition_id = ?1",
2309 [transition_id],
2310 |row| {
2311 Ok((
2312 row.get::<_, String>(0)?,
2313 row.get::<_, String>(1)?,
2314 row.get::<_, i64>(2)?,
2315 row.get::<_, i64>(3)?,
2316 row.get::<_, Vec<u8>>(4)?,
2317 row.get::<_, String>(5)?,
2318 row.get::<_, String>(6)?,
2319 row.get::<_, String>(7)?,
2320 row.get::<_, String>(8)?,
2321 ))
2322 },
2323 )
2324 .optional()
2325 .map_err(sqlite_error)?;
2326 let Some(mutation) = mutation else {
2327 return Ok(());
2328 };
2329 let expected_version = read_u64(mutation.2, "expected admission version")?;
2330 let activated_version = read_u64(mutation.3, "activated admission version")?;
2331 let admission_updated = transaction
2332 .execute(
2333 "UPDATE fiscal_proposal_admissions SET status = 'activated', state_version = ?1, state_json = ?2 WHERE admission_id = ?3 AND admission_digest = ?4 AND state_version = ?5 AND status = 'admitted'",
2334 params![
2335 sqlite_i64(activated_version, "activated admission version")?,
2336 &mutation.4,
2337 &mutation.0,
2338 &mutation.1,
2339 sqlite_i64(expected_version, "expected admission version")?,
2340 ],
2341 )
2342 .map_err(sqlite_error)?;
2343 let successor_updated = transaction
2344 .execute(
2345 "UPDATE fiscal_charters SET lifecycle_state = 'active' WHERE charter_id = ?1 AND charter_digest = ?2 AND lifecycle_state = 'proposed'",
2346 params![&mutation.5, &mutation.6],
2347 )
2348 .map_err(sqlite_error)?;
2349 let predecessor_updated = transaction
2350 .execute(
2351 "UPDATE fiscal_charters SET lifecycle_state = 'superseded' WHERE charter_id = ?1 AND charter_digest = ?2 AND lifecycle_state IN ('pinned', 'active')",
2352 params![&mutation.7, &mutation.8],
2353 )
2354 .map_err(sqlite_error)?;
2355 let schedule_count: i64 = transaction
2356 .query_row(
2357 "SELECT COUNT(*) FROM fiscal_staged_rotation_schedules WHERE transition_id = ?1",
2358 [transition_id],
2359 |row| row.get(0),
2360 )
2361 .map_err(sqlite_error)?;
2362 let candidates_updated = transaction
2363 .execute(
2364 "UPDATE fiscal_schedules SET lifecycle_state = 'active' WHERE lifecycle_state = 'staged' AND schedule_id IN (SELECT candidate_schedule_id FROM fiscal_staged_rotation_schedules WHERE transition_id = ?1)",
2365 [transition_id],
2366 )
2367 .map_err(sqlite_error)?;
2368 let predecessors_updated = transaction
2369 .execute(
2370 "UPDATE fiscal_schedules SET lifecycle_state = 'superseded' WHERE lifecycle_state = 'active' AND schedule_id IN (SELECT predecessor_schedule_id FROM fiscal_staged_rotation_schedules WHERE transition_id = ?1)",
2371 [transition_id],
2372 )
2373 .map_err(sqlite_error)?;
2374 if admission_updated != 1
2375 || successor_updated != 1
2376 || predecessor_updated != 1
2377 || i64::try_from(candidates_updated).map_err(|_| invariant("candidate update overflow"))?
2378 != schedule_count
2379 || i64::try_from(predecessors_updated)
2380 .map_err(|_| invariant("predecessor update overflow"))?
2381 != schedule_count
2382 {
2383 return Err(FiscalStoreError::Conflict);
2384 }
2385 append_projection_commit(
2386 transaction,
2387 owner,
2388 &format!("admission:{}", mutation.0),
2389 activated_version,
2390 "activate_fiscal_admission",
2391 &sha256_hex(&mutation.4),
2392 )
2393}
2394
2395pub(crate) fn initialize_fiscal_schema(
2396 connection: &mut Connection,
2397) -> Result<(), FiscalStoreError> {
2398 let on_disk = crate::check_schema_version(
2399 connection,
2400 FISCAL_STORE_SCHEMA_KEY,
2401 FISCAL_STORE_SUPPORTED_SCHEMA_VERSION,
2402 &["fiscal_authority_state", "capability_grant_budgets"],
2403 )
2404 .map_err(|error| invariant(error.to_string()))?;
2405 if on_disk == FISCAL_STORE_SUPPORTED_SCHEMA_VERSION {
2406 return verify_fiscal_sql_invariants(connection);
2407 }
2408 let transaction = connection
2409 .transaction_with_behavior(TransactionBehavior::Immediate)
2410 .map_err(sqlite_error)?;
2411 ensure_legacy_envelope_digest_column(&transaction)?;
2412 transaction
2413 .execute_batch(FISCAL_STORE_SCHEMA)
2414 .map_err(sqlite_error)?;
2415 crate::stamp_schema_version(
2416 &transaction,
2417 FISCAL_STORE_SCHEMA_KEY,
2418 FISCAL_STORE_SUPPORTED_SCHEMA_VERSION,
2419 )
2420 .map_err(|error| invariant(error.to_string()))?;
2421 verify_fiscal_sql_invariants(&transaction)?;
2422 transaction.commit().map_err(sqlite_error)
2423}
2424
2425pub(crate) fn verify_fiscal_sql_invariants(
2426 connection: &Connection,
2427) -> Result<(), FiscalStoreError> {
2428 let invalid = connection
2429 .query_row(
2430 r#"
2431 SELECT
2432 EXISTS(
2433 SELECT 1 FROM fiscal_authority_state AS authority
2434 WHERE NOT EXISTS(
2435 SELECT 1 FROM fiscal_continuity_checkpoints AS checkpoint
2436 WHERE checkpoint.checkpoint_digest = authority.finalized_checkpoint_digest
2437 AND checkpoint.status = 'finalized'
2438 )
2439 )
2440 OR EXISTS(
2441 SELECT 1 FROM fiscal_staged_transitions AS stage
2442 WHERE NOT EXISTS(
2443 SELECT 1 FROM fiscal_projection_commits AS commit_record
2444 WHERE commit_record.projection_key = 'transition:' || stage.transition_id
2445 AND commit_record.projection_sequence = stage.stage_version
2446 )
2447 OR NOT EXISTS(
2448 SELECT 1 FROM fiscal_continuity_checkpoints AS checkpoint
2449 WHERE checkpoint.checkpoint_digest = stage.next_checkpoint_digest
2450 AND (
2451 (stage.status = 'db_finalized' AND checkpoint.status = 'finalized')
2452 OR (stage.status <> 'db_finalized' AND checkpoint.status = 'staged')
2453 )
2454 )
2455 OR (
2456 stage.status <> 'db_finalized'
2457 AND NOT EXISTS(
2458 SELECT 1 FROM fiscal_authority_state AS authority
2459 WHERE authority.finalized_checkpoint_digest = stage.current_checkpoint_digest
2460 )
2461 )
2462 )
2463 OR EXISTS(
2464 SELECT 1 FROM fiscal_proposal_admissions AS admission
2465 WHERE NOT EXISTS(
2466 SELECT 1 FROM fiscal_projection_commits AS commit_record
2467 WHERE commit_record.projection_key = 'admission:' || admission.admission_id
2468 AND commit_record.projection_sequence = admission.state_version
2469 )
2470 )
2471 OR EXISTS(
2472 SELECT 1 FROM fiscal_legacy_fee_schedule_bindings AS binding
2473 WHERE binding.legacy_envelope_digest IS NULL
2474 OR length(binding.legacy_envelope_digest) <> 64
2475 OR binding.legacy_envelope_digest GLOB '*[^0-9a-f]*'
2476 OR NOT EXISTS(
2477 SELECT 1 FROM fiscal_schedules AS schedule
2478 WHERE schedule.schedule_id = binding.fiscal_schedule_id
2479 AND schedule.schedule_digest = binding.fiscal_schedule_digest
2480 )
2481 )
2482 OR EXISTS(
2483 SELECT 1
2484 FROM fiscal_staged_activation_mutations AS mutation
2485 JOIN fiscal_staged_transitions AS stage
2486 ON stage.transition_id = mutation.transition_id
2487 WHERE NOT EXISTS(
2488 SELECT 1 FROM fiscal_activations AS activation
2489 WHERE activation.activation_id = mutation.activation_id
2490 AND activation.activation_digest = mutation.activation_digest
2491 )
2492 OR (
2493 stage.status = 'db_finalized'
2494 AND (
2495 NOT EXISTS(
2496 SELECT 1 FROM fiscal_proposal_admissions AS admission
2497 WHERE admission.admission_id = mutation.admission_id
2498 AND admission.admission_digest = mutation.admission_digest
2499 AND admission.status = 'activated'
2500 AND admission.state_version = mutation.activated_admission_version
2501 AND admission.state_json = mutation.activated_admission_json
2502 )
2503 OR NOT EXISTS(
2504 SELECT 1 FROM fiscal_schedules AS candidate
2505 WHERE candidate.schedule_id = mutation.candidate_schedule_id
2506 AND candidate.schedule_digest = mutation.candidate_schedule_digest
2507 AND candidate.lifecycle_state IN ('active', 'superseded')
2508 )
2509 OR (
2510 mutation.predecessor_schedule_id IS NOT NULL
2511 AND NOT EXISTS(
2512 SELECT 1 FROM fiscal_schedules AS predecessor
2513 WHERE predecessor.schedule_id = mutation.predecessor_schedule_id
2514 AND predecessor.schedule_digest = mutation.predecessor_schedule_digest
2515 AND predecessor.lifecycle_state = 'superseded'
2516 )
2517 )
2518 OR NOT EXISTS(
2519 SELECT 1 FROM fiscal_projection_commits AS commit_record
2520 WHERE commit_record.projection_key = 'admission:' || mutation.admission_id
2521 AND commit_record.projection_sequence = mutation.activated_admission_version
2522 )
2523 )
2524 )
2525 OR (
2526 stage.status <> 'db_finalized'
2527 AND (
2528 NOT EXISTS(
2529 SELECT 1 FROM fiscal_proposal_admissions AS admission
2530 WHERE admission.admission_id = mutation.admission_id
2531 AND admission.admission_digest = mutation.admission_digest
2532 AND admission.status = 'admitted'
2533 AND admission.state_version = mutation.expected_admission_version
2534 )
2535 OR NOT EXISTS(
2536 SELECT 1 FROM fiscal_schedules AS candidate
2537 WHERE candidate.schedule_id = mutation.candidate_schedule_id
2538 AND candidate.schedule_digest = mutation.candidate_schedule_digest
2539 AND candidate.lifecycle_state = 'staged'
2540 )
2541 OR (
2542 mutation.predecessor_schedule_id IS NOT NULL
2543 AND NOT EXISTS(
2544 SELECT 1 FROM fiscal_schedules AS predecessor
2545 WHERE predecessor.schedule_id = mutation.predecessor_schedule_id
2546 AND predecessor.schedule_digest = mutation.predecessor_schedule_digest
2547 AND predecessor.lifecycle_state = 'active'
2548 )
2549 )
2550 )
2551 )
2552 )
2553 OR EXISTS(
2554 SELECT 1
2555 FROM fiscal_staged_rotation_mutations AS mutation
2556 JOIN fiscal_staged_transitions AS stage
2557 ON stage.transition_id = mutation.transition_id
2558 WHERE NOT EXISTS(
2559 SELECT 1 FROM fiscal_activations AS activation
2560 WHERE activation.activation_id = mutation.activation_id
2561 AND activation.activation_digest = mutation.activation_digest
2562 )
2563 OR NOT EXISTS(
2564 SELECT 1 FROM fiscal_staged_rotation_schedules AS replacement
2565 WHERE replacement.transition_id = mutation.transition_id
2566 )
2567 OR (
2568 stage.status = 'db_finalized'
2569 AND (
2570 NOT EXISTS(
2571 SELECT 1 FROM fiscal_proposal_admissions AS admission
2572 WHERE admission.admission_id = mutation.admission_id
2573 AND admission.admission_digest = mutation.admission_digest
2574 AND admission.status = 'activated'
2575 AND admission.state_version = mutation.activated_admission_version
2576 AND admission.state_json = mutation.activated_admission_json
2577 )
2578 OR NOT EXISTS(
2579 SELECT 1 FROM fiscal_charters AS successor
2580 WHERE successor.charter_id = mutation.successor_charter_id
2581 AND successor.charter_digest = mutation.successor_charter_digest
2582 AND successor.lifecycle_state IN ('active', 'superseded')
2583 )
2584 OR NOT EXISTS(
2585 SELECT 1 FROM fiscal_charters AS predecessor
2586 WHERE predecessor.charter_id = mutation.predecessor_charter_id
2587 AND predecessor.charter_digest = mutation.predecessor_charter_digest
2588 AND predecessor.lifecycle_state = 'superseded'
2589 )
2590 OR EXISTS(
2591 SELECT 1 FROM fiscal_staged_rotation_schedules AS replacement
2592 WHERE replacement.transition_id = mutation.transition_id
2593 AND (
2594 NOT EXISTS(
2595 SELECT 1 FROM fiscal_schedules AS candidate
2596 WHERE candidate.schedule_id = replacement.candidate_schedule_id
2597 AND candidate.schedule_digest = replacement.candidate_schedule_digest
2598 AND candidate.lifecycle_state IN ('active', 'superseded')
2599 )
2600 OR NOT EXISTS(
2601 SELECT 1 FROM fiscal_schedules AS predecessor
2602 WHERE predecessor.schedule_id = replacement.predecessor_schedule_id
2603 AND predecessor.schedule_digest = replacement.predecessor_schedule_digest
2604 AND predecessor.lifecycle_state = 'superseded'
2605 )
2606 )
2607 )
2608 OR NOT EXISTS(
2609 SELECT 1 FROM fiscal_projection_commits AS commit_record
2610 WHERE commit_record.projection_key = 'admission:' || mutation.admission_id
2611 AND commit_record.projection_sequence = mutation.activated_admission_version
2612 )
2613 )
2614 )
2615 OR (
2616 stage.status <> 'db_finalized'
2617 AND (
2618 NOT EXISTS(
2619 SELECT 1 FROM fiscal_proposal_admissions AS admission
2620 WHERE admission.admission_id = mutation.admission_id
2621 AND admission.admission_digest = mutation.admission_digest
2622 AND admission.status = 'admitted'
2623 AND admission.state_version = mutation.expected_admission_version
2624 )
2625 OR NOT EXISTS(
2626 SELECT 1 FROM fiscal_charters AS successor
2627 WHERE successor.charter_id = mutation.successor_charter_id
2628 AND successor.charter_digest = mutation.successor_charter_digest
2629 AND successor.lifecycle_state = 'proposed'
2630 )
2631 OR NOT EXISTS(
2632 SELECT 1 FROM fiscal_charters AS predecessor
2633 WHERE predecessor.charter_id = mutation.predecessor_charter_id
2634 AND predecessor.charter_digest = mutation.predecessor_charter_digest
2635 AND predecessor.lifecycle_state IN ('pinned', 'active')
2636 )
2637 OR EXISTS(
2638 SELECT 1 FROM fiscal_staged_rotation_schedules AS replacement
2639 WHERE replacement.transition_id = mutation.transition_id
2640 AND (
2641 NOT EXISTS(
2642 SELECT 1 FROM fiscal_schedules AS candidate
2643 WHERE candidate.schedule_id = replacement.candidate_schedule_id
2644 AND candidate.schedule_digest = replacement.candidate_schedule_digest
2645 AND candidate.lifecycle_state = 'staged'
2646 )
2647 OR NOT EXISTS(
2648 SELECT 1 FROM fiscal_schedules AS predecessor
2649 WHERE predecessor.schedule_id = replacement.predecessor_schedule_id
2650 AND predecessor.schedule_digest = replacement.predecessor_schedule_digest
2651 AND predecessor.lifecycle_state = 'active'
2652 )
2653 )
2654 )
2655 )
2656 )
2657 )
2658 OR EXISTS(
2659 SELECT 1 FROM fiscal_projection_commits AS commit_record
2660 WHERE commit_record.projection_sequence > 1
2661 AND NOT EXISTS(
2662 SELECT 1 FROM fiscal_projection_commits AS previous
2663 WHERE previous.projection_key = commit_record.projection_key
2664 AND previous.projection_sequence = commit_record.projection_sequence - 1
2665 AND previous.commit_digest = commit_record.previous_commit_digest
2666 )
2667 )
2668 OR EXISTS(
2669 SELECT 1 FROM fiscal_projection_commits
2670 WHERE projection_sequence = 1 AND previous_commit_digest <> ?1
2671 )
2672 "#,
2673 [ZERO_DIGEST],
2674 |row| row.get::<_, bool>(0),
2675 )
2676 .map_err(sqlite_error)?;
2677 if invalid {
2678 Err(invariant("fiscal store projection is inconsistent"))
2679 } else {
2680 Ok(())
2681 }
2682}
2683
2684fn ensure_legacy_envelope_digest_column(
2685 transaction: &Transaction<'_>,
2686) -> Result<(), FiscalStoreError> {
2687 let table_present = transaction
2688 .query_row(
2689 "SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE type = 'table' AND name = 'fiscal_legacy_fee_schedule_bindings')",
2690 [],
2691 |row| row.get::<_, bool>(0),
2692 )
2693 .map_err(sqlite_error)?;
2694 if !table_present {
2695 return Ok(());
2696 }
2697 let mut statement = transaction
2698 .prepare("PRAGMA table_info(fiscal_legacy_fee_schedule_bindings)")
2699 .map_err(sqlite_error)?;
2700 let columns = statement
2701 .query_map([], |row| row.get::<_, String>(1))
2702 .map_err(sqlite_error)?
2703 .collect::<Result<Vec<_>, _>>()
2704 .map_err(sqlite_error)?;
2705 if !columns
2706 .iter()
2707 .any(|column| column == "legacy_envelope_digest")
2708 {
2709 transaction
2710 .execute(
2711 "ALTER TABLE fiscal_legacy_fee_schedule_bindings ADD COLUMN legacy_envelope_digest TEXT",
2712 [],
2713 )
2714 .map_err(sqlite_error)?;
2715 }
2716 Ok(())
2717}
2718
2719#[allow(clippy::too_many_arguments)]
2720fn verify_exact_genesis(
2721 transaction: &Transaction<'_>,
2722 policy_json: &[u8],
2723 authority_json: &[u8],
2724 charter_json: &[u8],
2725 readiness_json: &[u8],
2726 registry_json: &[u8],
2727 checkpoint_json: &[u8],
2728 policy: &FiscalGenesisPolicy,
2729 charter: &VerifiedFiscalCharter,
2730 readiness: &VerifiedFiscalRuntimeReadiness,
2731 checkpoint: &VerifiedFiscalContinuityCheckpoint,
2732) -> Result<(), FiscalStoreError> {
2733 let exact = transaction
2734 .query_row(
2735 r#"
2736 SELECT
2737 (SELECT policy_json FROM fiscal_genesis_policies WHERE policy_id = ?1) = ?2
2738 AND (SELECT state_json FROM fiscal_authority_state WHERE singleton = 1) = ?3
2739 AND (SELECT signed_json FROM fiscal_charters WHERE charter_id = ?4) = ?5
2740 AND (SELECT signed_json FROM fiscal_runtime_readiness WHERE readiness_id = ?6) = ?7
2741 AND (SELECT registry_json FROM fiscal_runtime_readiness WHERE readiness_id = ?6) = ?8
2742 AND (SELECT signed_json FROM fiscal_continuity_checkpoints WHERE checkpoint_digest = ?9) = ?10
2743 "#,
2744 params![
2745 &policy.policy_id,
2746 policy_json,
2747 authority_json,
2748 &charter.body().charter_id,
2749 charter_json,
2750 &readiness.body().readiness_id,
2751 readiness_json,
2752 registry_json,
2753 checkpoint.digest(),
2754 checkpoint_json,
2755 ],
2756 |row| row.get::<_, bool>(0),
2757 )
2758 .map_err(sqlite_error)?;
2759 if exact {
2760 Ok(())
2761 } else {
2762 Err(FiscalStoreError::Conflict)
2763 }
2764}
2765
2766fn verify_acknowledgement(
2767 transaction: &Transaction<'_>,
2768 record: &FiscalStagedTransitionRecord,
2769 checkpoint: &VerifiedFiscalContinuityCheckpoint,
2770 canonical: &[u8],
2771) -> Result<(), FiscalStoreError> {
2772 if checkpoint.digest() != record.next_checkpoint_digest {
2773 return Err(FiscalStoreError::Conflict);
2774 }
2775 let stored = transaction
2776 .query_row(
2777 "SELECT signed_json FROM fiscal_continuity_checkpoints WHERE checkpoint_digest = ?1",
2778 [&record.next_checkpoint_digest],
2779 |row| row.get::<_, Vec<u8>>(0),
2780 )
2781 .optional()
2782 .map_err(sqlite_error)?
2783 .ok_or(FiscalStoreError::NotFound)?;
2784 if stored != canonical {
2785 return Err(FiscalStoreError::Conflict);
2786 }
2787 Ok(())
2788}
2789
2790fn load_transition(
2791 transaction: &Transaction<'_>,
2792 transition_id: &str,
2793) -> Result<FiscalStagedTransitionRecord, FiscalStoreError> {
2794 transaction
2795 .query_row(
2796 "SELECT transition_id, current_checkpoint_digest, next_checkpoint_digest, proof_json, status, stage_version FROM fiscal_staged_transitions WHERE transition_id = ?1",
2797 [transition_id],
2798 |row| {
2799 let status = row.get::<_, String>(4)?;
2800 Ok((
2801 row.get::<_, String>(0)?,
2802 row.get::<_, String>(1)?,
2803 row.get::<_, String>(2)?,
2804 row.get::<_, Vec<u8>>(3)?,
2805 status,
2806 row.get::<_, i64>(5)?,
2807 ))
2808 },
2809 )
2810 .optional()
2811 .map_err(sqlite_error)?
2812 .ok_or(FiscalStoreError::NotFound)
2813 .and_then(|row| {
2814 Ok(FiscalStagedTransitionRecord {
2815 transition_id: row.0,
2816 current_checkpoint_digest: row.1,
2817 next_checkpoint_digest: row.2,
2818 proof_json: row.3,
2819 status: FiscalStageStatus::parse(&row.4)?,
2820 stage_version: read_u64(row.5, "fiscal stage version")?,
2821 })
2822 })
2823}
2824
2825fn append_projection_commit(
2826 transaction: &Transaction<'_>,
2827 owner: &SqliteServingOwner,
2828 projection_key: &str,
2829 projection_sequence: u64,
2830 mutation_kind: &str,
2831 snapshot_digest: &str,
2832) -> Result<(), FiscalStoreError> {
2833 let previous = if projection_sequence == 1 {
2834 ZERO_DIGEST.to_owned()
2835 } else {
2836 let previous_sequence = projection_sequence
2837 .checked_sub(1)
2838 .ok_or_else(|| invariant("fiscal projection sequence underflowed"))?;
2839 transaction
2840 .query_row(
2841 "SELECT commit_digest FROM fiscal_projection_commits WHERE projection_key = ?1 AND projection_sequence = ?2",
2842 params![projection_key, sqlite_i64(previous_sequence, "previous fiscal projection sequence")?],
2843 |row| row.get::<_, String>(0),
2844 )
2845 .optional()
2846 .map_err(sqlite_error)?
2847 .ok_or_else(|| invariant("previous fiscal projection commit is absent"))?
2848 };
2849 let commit_digest = canonical_digest(&FiscalProjectionCommit {
2850 format: "chio.sqlite-fiscal-projection-commit.v1",
2851 projection_key,
2852 projection_sequence,
2853 mutation_kind,
2854 snapshot_digest,
2855 previous_commit_digest: &previous,
2856 store_uuid: &owner.fence.store_uuid,
2857 store_lease_id: &owner.fence.lease_id,
2858 store_owner_epoch: owner.fence.owner_epoch,
2859 })?;
2860 transaction
2861 .execute(
2862 "INSERT INTO fiscal_projection_commits (projection_key, projection_sequence, mutation_kind, snapshot_digest, previous_commit_digest, commit_digest, store_uuid, store_lease_id, store_owner_epoch) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
2863 params![
2864 projection_key,
2865 sqlite_i64(projection_sequence, "fiscal projection sequence")?,
2866 mutation_kind,
2867 snapshot_digest,
2868 &previous,
2869 &commit_digest,
2870 &owner.fence.store_uuid,
2871 &owner.fence.lease_id,
2872 sqlite_i64(owner.fence.owner_epoch, "store owner epoch")?,
2873 ],
2874 )
2875 .map_err(sqlite_error)?;
2876 owner
2877 .append_global_commit(
2878 transaction,
2879 mutation_kind,
2880 "fiscal",
2881 projection_key,
2882 projection_sequence,
2883 )
2884 .map_err(map_owner_error)
2885}
2886
2887fn verify_owner(
2888 transaction: &Transaction<'_>,
2889 owner: &SqliteServingOwner,
2890 fence: Option<&StoreMutationFence>,
2891) -> Result<(), FiscalStoreError> {
2892 crate::admission_operation_store::verify_active_owner(transaction, owner, fence).map_err(
2893 |error| match error {
2894 chio_kernel::admission_operation::AdmissionOperationStoreError::Fenced => {
2895 FiscalStoreError::Fenced
2896 }
2897 other => FiscalStoreError::Unavailable(other.to_string()),
2898 },
2899 )
2900}
2901
2902fn canonical_digest(value: &impl Serialize) -> Result<String, FiscalStoreError> {
2903 canonical_json_bytes(value)
2904 .map(|bytes| sha256_hex(&bytes))
2905 .map_err(canonical_error)
2906}
2907
2908fn canonical_error(error: impl std::fmt::Display) -> FiscalStoreError {
2909 invariant(format!("canonical fiscal encoding failed: {error}"))
2910}
2911
2912fn sqlite_error(error: rusqlite::Error) -> FiscalStoreError {
2913 FiscalStoreError::Unavailable(error.to_string())
2914}
2915
2916fn map_owner_error(error: SqliteServingOwnerError) -> FiscalStoreError {
2917 match error {
2918 SqliteServingOwnerError::OutcomeUnknown(detail) => FiscalStoreError::OutcomeUnknown(detail),
2919 other => FiscalStoreError::Unavailable(other.to_string()),
2920 }
2921}
2922
2923fn invariant(detail: impl Into<String>) -> FiscalStoreError {
2924 FiscalStoreError::Invariant(detail.into())
2925}
2926
2927fn sqlite_i64(value: u64, field: &str) -> Result<i64, FiscalStoreError> {
2928 i64::try_from(value).map_err(|_| invariant(format!("{field} exceeds SQLite INTEGER")))
2929}
2930
2931fn read_u64(value: i64, field: &str) -> Result<u64, FiscalStoreError> {
2932 u64::try_from(value).map_err(|_| invariant(format!("{field} is negative")))
2933}
2934
2935#[cfg(test)]
2936#[path = "fiscal_store_tests.rs"]
2937mod tests;