open.rs (40312B)
1 //! SQLite storage lifecycle modes and owned paths. 2 3 use std::error::Error as StdError; 4 use std::fmt; 5 use std::path::{Component, Path, PathBuf}; 6 use std::time::Duration; 7 8 use crate::{OpenOptions, event::SqliteStorage, lock::WriterLock, migration}; 9 use radroots_storage::event::SourceGeneration; 10 use sqlx::{ 11 ConnectOptions, Connection, Row, SqliteConnection, SqlitePool, 12 sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions, SqliteSynchronous}, 13 }; 14 15 const RUNTIME_DATABASE_NAME: &str = "runtime.sqlite"; 16 const PRIVATE_DATABASE_NAME: &str = "private.sqlite"; 17 const MAX_CONNECTIONS_PER_DATABASE: u32 = 4; 18 19 /// Explicit behavior for opening owned SQLite files. 20 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 21 #[non_exhaustive] 22 pub enum OpenMode { 23 /// Open both existing files without permitting mutation or migration. 24 ReadOnly, 25 /// Open both existing files with write and migration authority. 26 ReadWriteExisting, 27 /// Open or create the two owned files with write and migration authority. 28 Create, 29 } 30 31 impl OpenMode { 32 /// Returns whether the mode may mutate owned database files. 33 pub fn is_writable(self) -> bool { 34 matches!(self, Self::ReadWriteExisting | Self::Create) 35 } 36 37 /// Returns whether missing owned files may be created. 38 pub fn may_create(self) -> bool { 39 matches!(self, Self::Create) 40 } 41 } 42 43 /// The complete SQLite file set owned by one Radroots storage backend. 44 #[derive(Clone, Debug, Eq, PartialEq)] 45 pub struct Paths { 46 runtime: PathBuf, 47 private: PathBuf, 48 } 49 50 impl Paths { 51 /// Derives the governed file names from an absolute owner directory. 52 pub fn from_directory(directory: impl AsRef<Path>) -> Result<Self, Error> { 53 let directory = directory.as_ref(); 54 validate_absolute_normal_path(directory)?; 55 Self::from_files( 56 directory.join(RUNTIME_DATABASE_NAME), 57 directory.join(PRIVATE_DATABASE_NAME), 58 ) 59 } 60 61 /// Validates explicit runtime and private file paths. 62 pub fn from_files( 63 runtime: impl Into<PathBuf>, 64 private: impl Into<PathBuf>, 65 ) -> Result<Self, Error> { 66 let runtime = runtime.into(); 67 let private = private.into(); 68 validate_owned_file_path(&runtime, RUNTIME_DATABASE_NAME)?; 69 validate_owned_file_path(&private, PRIVATE_DATABASE_NAME)?; 70 if runtime == private { 71 return Err(Error::PathsOverlap(runtime)); 72 } 73 Ok(Self { runtime, private }) 74 } 75 76 /// Returns the canonical event, journal, outbox, and projection database. 77 pub fn runtime(&self) -> &Path { 78 &self.runtime 79 } 80 81 /// Returns the encrypted private-artifact database. 82 pub fn private(&self) -> &Path { 83 &self.private 84 } 85 86 #[cfg_attr(coverage_nightly, coverage(off))] 87 pub(crate) fn validate_filesystem(&self, mode: OpenMode) -> Result<(), Error> { 88 for path in [&self.runtime, &self.private] { 89 validate_parent(path)?; 90 match std::fs::symlink_metadata(path) { 91 Ok(metadata) if metadata.file_type().is_symlink() => { 92 return Err(Error::SymlinkPath(path.clone())); 93 } 94 Ok(metadata) if !metadata.is_file() => { 95 return Err(Error::NotAFile(path.clone())); 96 } 97 Ok(_) => {} 98 Err(source) if source.kind() == std::io::ErrorKind::NotFound => { 99 if !mode.may_create() { 100 return Err(Error::MissingFile(path.clone())); 101 } 102 } 103 Err(source) => { 104 return Err(Error::Inspect { 105 path: path.clone(), 106 source, 107 }); 108 } 109 } 110 } 111 Ok(()) 112 } 113 } 114 115 fn validate_owned_file_path(path: &Path, expected_name: &'static str) -> Result<(), Error> { 116 validate_absolute_normal_path(path)?; 117 if path.file_name().and_then(|name| name.to_str()) != Some(expected_name) { 118 return Err(Error::UnexpectedFileName { 119 path: path.to_path_buf(), 120 expected: expected_name, 121 }); 122 } 123 Ok(()) 124 } 125 126 fn validate_absolute_normal_path(path: &Path) -> Result<(), Error> { 127 if !path.is_absolute() 128 || path 129 .components() 130 .any(|part| matches!(part, Component::CurDir | Component::ParentDir)) 131 { 132 return Err(Error::InvalidPath(path.to_path_buf())); 133 } 134 Ok(()) 135 } 136 137 #[cfg_attr(coverage_nightly, coverage(off))] 138 fn validate_parent(path: &Path) -> Result<(), Error> { 139 let parent = path 140 .parent() 141 .ok_or_else(|| Error::InvalidPath(path.to_path_buf()))?; 142 match std::fs::metadata(parent) { 143 Ok(metadata) if metadata.is_dir() => Ok(()), 144 Ok(_) => Err(Error::ParentNotDirectory(parent.to_path_buf())), 145 Err(source) if source.kind() == std::io::ErrorKind::NotFound => { 146 Err(Error::MissingParent(parent.to_path_buf())) 147 } 148 Err(source) => Err(Error::Inspect { 149 path: parent.to_path_buf(), 150 source, 151 }), 152 } 153 } 154 155 /// Stable, secret-safe SQLite lifecycle or filesystem failure. 156 #[derive(Debug)] 157 #[non_exhaustive] 158 pub enum Error { 159 /// Local capacity failure; inspect existing state before retrying startup. 160 SpaceInsufficient, 161 InvalidPath(PathBuf), 162 UnexpectedFileName { 163 path: PathBuf, 164 expected: &'static str, 165 }, 166 PathsOverlap(PathBuf), 167 MissingParent(PathBuf), 168 ParentNotDirectory(PathBuf), 169 MissingFile(PathBuf), 170 SymlinkPath(PathBuf), 171 NotAFile(PathBuf), 172 Inspect { 173 path: PathBuf, 174 source: std::io::Error, 175 }, 176 InvalidBusyTimeout { 177 minimum: Duration, 178 maximum: Duration, 179 actual: Duration, 180 }, 181 InvalidSourceGenerationTimestamp { 182 actual: u64, 183 }, 184 WriterLockOpen { 185 path: PathBuf, 186 source: std::io::Error, 187 }, 188 WriterAlreadyActive { 189 path: PathBuf, 190 }, 191 WriterLockFailed { 192 path: PathBuf, 193 source: std::io::Error, 194 }, 195 WriterUnlockFailed { 196 path: PathBuf, 197 source: std::io::Error, 198 }, 199 SchemaMetadataUnavailable { 200 database: &'static str, 201 }, 202 SchemaIdentityMismatch { 203 database: &'static str, 204 expected: u32, 205 actual: u32, 206 }, 207 SchemaTooOld { 208 database: &'static str, 209 minimum: u32, 210 actual: u32, 211 }, 212 SchemaTooNew { 213 database: &'static str, 214 supported: u32, 215 actual: u32, 216 }, 217 SchemaMigrationRequired { 218 database: &'static str, 219 current: u32, 220 actual: u32, 221 }, 222 UnrecognizedSchema { 223 database: &'static str, 224 }, 225 SchemaCatalogMismatch { 226 database: &'static str, 227 version: u32, 228 }, 229 SchemaMigrationFailed { 230 database: &'static str, 231 target_version: u32, 232 }, 233 AuthoredMigrationBlocked { 234 prepared_or_recoverable: u64, 235 signed_without_complete_event: u64, 236 invalid_or_unsupported: u64, 237 }, 238 DatabaseOpenFailed { 239 database: &'static str, 240 }, 241 DatabaseCorrupt { 242 database: &'static str, 243 }, 244 DatabaseCloseFailed { 245 database: &'static str, 246 }, 247 ConnectionPolicyMismatch { 248 database: &'static str, 249 }, 250 SourceGenerationRequired, 251 SourceGenerationUnavailable, 252 SourceGenerationMismatch, 253 CorruptSourceGeneration, 254 InvalidBackupRoot(PathBuf), 255 BackupRootRequired, 256 BackupBundleAlreadyExists(PathBuf), 257 BackupBackendUnavailable, 258 UnsupportedBackupVersion, 259 BackupCaptureFailed { 260 member: &'static str, 261 }, 262 BackupBundleMissing(PathBuf), 263 BackupVerificationFailed { 264 member: &'static str, 265 }, 266 BackupUnexpectedEntry(PathBuf), 267 RestoreRequiresWritableStorage, 268 RestoreStagingAlreadyExists(PathBuf), 269 RestoreStagingFailed { 270 member: &'static str, 271 }, 272 RestoreMarkerCorrupt(PathBuf), 273 RestoreRecoveryConflict(PathBuf), 274 RestoreReplacementFailed { 275 member: &'static str, 276 }, 277 RestoreFilesystem { 278 operation: &'static str, 279 source: std::io::Error, 280 }, 281 InvalidLegacyImportPlan, 282 InvalidLegacySource(PathBuf), 283 LegacyImportBackupAlreadyExists(PathBuf), 284 LegacyImportBackupFailed { 285 source_kind: &'static str, 286 }, 287 LegacyImportSourceInvalid { 288 source_kind: &'static str, 289 }, 290 LegacyImportEvidenceInvalid, 291 LegacyImportTargetMismatch, 292 LegacyImportMigrationHistoryInvalid, 293 InvalidLegacyImportJournal, 294 LegacyImportConflict, 295 LegacyImportJournalFailed, 296 InvalidLegacyImportStageRequest, 297 LegacyImportStagingFailed, 298 LegacyImportRowInvalid { 299 source_kind: &'static str, 300 legacy_sequence: i64, 301 }, 302 UnsupportedLegacySchema { 303 source_kind: &'static str, 304 user_version: i64, 305 catalog_sha256: String, 306 }, 307 LegacyImportFilesystem { 308 operation: &'static str, 309 source: std::io::Error, 310 }, 311 BackupFilesystem { 312 operation: &'static str, 313 source: std::io::Error, 314 }, 315 } 316 317 #[cfg_attr(coverage_nightly, coverage(off))] 318 impl fmt::Display for Error { 319 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 320 match self { 321 Self::SpaceInsufficient => formatter.write_str("storage space is insufficient"), 322 Self::InvalidPath(path) => { 323 write!(formatter, "invalid owned SQLite path: {}", path.display()) 324 } 325 Self::UnexpectedFileName { path, expected } => write!( 326 formatter, 327 "owned SQLite path {} must use file name {expected}", 328 path.display() 329 ), 330 Self::PathsOverlap(path) => { 331 write!(formatter, "owned SQLite paths overlap: {}", path.display()) 332 } 333 Self::MissingParent(path) => write!( 334 formatter, 335 "owned SQLite parent is missing: {}", 336 path.display() 337 ), 338 Self::ParentNotDirectory(path) => write!( 339 formatter, 340 "owned SQLite parent is not a directory: {}", 341 path.display() 342 ), 343 Self::MissingFile(path) => write!( 344 formatter, 345 "owned SQLite file is missing: {}", 346 path.display() 347 ), 348 Self::SymlinkPath(path) => write!( 349 formatter, 350 "owned SQLite path cannot be a symlink: {}", 351 path.display() 352 ), 353 Self::NotAFile(path) => write!( 354 formatter, 355 "owned SQLite path is not a file: {}", 356 path.display() 357 ), 358 Self::Inspect { path, .. } => write!( 359 formatter, 360 "failed to inspect owned SQLite path: {}", 361 path.display() 362 ), 363 Self::InvalidBusyTimeout { 364 minimum, 365 maximum, 366 actual, 367 } => write!( 368 formatter, 369 "SQLite busy timeout {actual:?} must be within {minimum:?}..={maximum:?}" 370 ), 371 Self::InvalidSourceGenerationTimestamp { actual } => write!( 372 formatter, 373 "source generation creation time {actual} must fit a positive SQLite integer" 374 ), 375 Self::WriterLockOpen { path, .. } => write!( 376 formatter, 377 "failed to open SQLite writer lock: {}", 378 path.display() 379 ), 380 Self::WriterAlreadyActive { path } => write!( 381 formatter, 382 "another writable SQLite storage client holds: {}", 383 path.display() 384 ), 385 Self::WriterLockFailed { path, .. } => write!( 386 formatter, 387 "failed to acquire SQLite writer lock: {}", 388 path.display() 389 ), 390 Self::WriterUnlockFailed { path, .. } => write!( 391 formatter, 392 "failed to release SQLite writer lock: {}", 393 path.display() 394 ), 395 Self::SchemaMetadataUnavailable { database } => { 396 write!(formatter, "failed to inspect {database} schema metadata") 397 } 398 Self::SchemaIdentityMismatch { 399 database, 400 expected, 401 actual, 402 } => write!( 403 formatter, 404 "{database} application id {actual} does not match required id {expected}" 405 ), 406 Self::SchemaTooOld { 407 database, 408 minimum, 409 actual, 410 } => write!( 411 formatter, 412 "{database} schema version {actual} is older than supported version {minimum}" 413 ), 414 Self::SchemaTooNew { 415 database, 416 supported, 417 actual, 418 } => write!( 419 formatter, 420 "{database} schema version {actual} is newer than supported version {supported}" 421 ), 422 Self::SchemaMigrationRequired { 423 database, 424 current, 425 actual, 426 } => write!( 427 formatter, 428 "{database} schema version {actual} requires writable migration to {current}" 429 ), 430 Self::UnrecognizedSchema { database } => { 431 write!( 432 formatter, 433 "{database} has an unrecognized unversioned schema" 434 ) 435 } 436 Self::SchemaCatalogMismatch { database, version } => write!( 437 formatter, 438 "{database} object catalog does not match schema version {version}" 439 ), 440 Self::SchemaMigrationFailed { 441 database, 442 target_version, 443 } => write!( 444 formatter, 445 "{database} migration to schema version {target_version} failed" 446 ), 447 Self::AuthoredMigrationBlocked { 448 prepared_or_recoverable, 449 signed_without_complete_event, 450 invalid_or_unsupported, 451 } => write!( 452 formatter, 453 "SQLite authored-operation migration is blocked: {prepared_or_recoverable} incomplete, {signed_without_complete_event} missing exact signed output, {invalid_or_unsupported} invalid or unsupported" 454 ), 455 Self::DatabaseOpenFailed { database } => { 456 write!( 457 formatter, 458 "failed to open governed SQLite database {database}" 459 ) 460 } 461 Self::DatabaseCorrupt { database } => { 462 write!(formatter, "governed SQLite database {database} is corrupt") 463 } 464 Self::DatabaseCloseFailed { database } => write!( 465 formatter, 466 "failed to close migration connection for {database}" 467 ), 468 Self::ConnectionPolicyMismatch { database } => write!( 469 formatter, 470 "{database} does not satisfy the governed SQLite connection policy" 471 ), 472 Self::SourceGenerationRequired => formatter.write_str( 473 "a fresh writable SQLite store requires a host-supplied source generation", 474 ), 475 Self::SourceGenerationUnavailable => { 476 formatter.write_str("SQLite storage has no active source generation") 477 } 478 Self::SourceGenerationMismatch => formatter 479 .write_str("host-supplied source generation does not match durable storage"), 480 Self::CorruptSourceGeneration => { 481 formatter.write_str("SQLite storage source generation is corrupt") 482 } 483 Self::InvalidBackupRoot(path) => { 484 write!(formatter, "invalid SQLite backup root: {}", path.display()) 485 } 486 Self::BackupRootRequired => { 487 formatter.write_str("SQLite backup requires a configured host-owned root") 488 } 489 Self::BackupBundleAlreadyExists(path) => write!( 490 formatter, 491 "SQLite backup bundle path already exists: {}", 492 path.display() 493 ), 494 Self::BackupBackendUnavailable => { 495 formatter.write_str("SQLite backup backend is unavailable") 496 } 497 Self::UnsupportedBackupVersion => { 498 formatter.write_str("SQLite backup format version is unsupported") 499 } 500 Self::BackupCaptureFailed { member } => { 501 write!(formatter, "failed to capture SQLite backup member {member}") 502 } 503 Self::BackupBundleMissing(path) => { 504 write!( 505 formatter, 506 "SQLite backup bundle is missing: {}", 507 path.display() 508 ) 509 } 510 Self::BackupVerificationFailed { member } => { 511 write!( 512 formatter, 513 "SQLite backup member failed verification: {member}" 514 ) 515 } 516 Self::BackupUnexpectedEntry(path) => write!( 517 formatter, 518 "SQLite backup bundle contains an unexpected entry: {}", 519 path.display() 520 ), 521 Self::RestoreRequiresWritableStorage => { 522 formatter.write_str("SQLite restore requires writable storage authority") 523 } 524 Self::RestoreStagingAlreadyExists(path) => write!( 525 formatter, 526 "SQLite restore staging path already exists: {}", 527 path.display() 528 ), 529 Self::RestoreStagingFailed { member } => { 530 write!(formatter, "failed to stage SQLite restore member {member}") 531 } 532 Self::RestoreMarkerCorrupt(path) => write!( 533 formatter, 534 "SQLite restore interruption marker is corrupt: {}", 535 path.display() 536 ), 537 Self::RestoreRecoveryConflict(path) => write!( 538 formatter, 539 "SQLite restore recovery state conflicts at: {}", 540 path.display() 541 ), 542 Self::RestoreReplacementFailed { member } => { 543 write!( 544 formatter, 545 "failed to replace SQLite restore member {member}" 546 ) 547 } 548 Self::RestoreFilesystem { operation, .. } => { 549 write!( 550 formatter, 551 "SQLite restore filesystem operation failed: {operation}" 552 ) 553 } 554 Self::InvalidLegacyImportPlan => { 555 formatter.write_str("SQLite legacy import plan is invalid") 556 } 557 Self::InvalidLegacySource(path) => write!( 558 formatter, 559 "SQLite legacy import source is invalid: {}", 560 path.display() 561 ), 562 Self::LegacyImportBackupAlreadyExists(path) => write!( 563 formatter, 564 "SQLite legacy import backup already exists: {}", 565 path.display() 566 ), 567 Self::LegacyImportBackupFailed { source_kind } => write!( 568 formatter, 569 "failed to back up SQLite legacy {source_kind} source" 570 ), 571 Self::LegacyImportSourceInvalid { source_kind } => write!( 572 formatter, 573 "SQLite legacy {source_kind} source failed integrity validation" 574 ), 575 Self::LegacyImportEvidenceInvalid => { 576 formatter.write_str("SQLite legacy import evidence is invalid") 577 } 578 Self::LegacyImportTargetMismatch => { 579 formatter.write_str("SQLite legacy import target generation does not match") 580 } 581 Self::LegacyImportMigrationHistoryInvalid => { 582 formatter.write_str("SQLite legacy event-store migration history is invalid") 583 } 584 Self::InvalidLegacyImportJournal => { 585 formatter.write_str("SQLite legacy import journal is invalid") 586 } 587 Self::LegacyImportConflict => { 588 formatter.write_str("SQLite legacy import conflicts with durable state") 589 } 590 Self::LegacyImportJournalFailed => { 591 formatter.write_str("SQLite legacy import journal operation failed") 592 } 593 Self::InvalidLegacyImportStageRequest => { 594 formatter.write_str("SQLite legacy import staging request is invalid") 595 } 596 Self::LegacyImportStagingFailed => { 597 formatter.write_str("SQLite legacy import staging operation failed") 598 } 599 Self::LegacyImportRowInvalid { 600 source_kind, 601 legacy_sequence, 602 } => write!( 603 formatter, 604 "SQLite legacy {source_kind} row {legacy_sequence} is invalid" 605 ), 606 Self::UnsupportedLegacySchema { 607 source_kind, 608 user_version, 609 catalog_sha256, 610 } => write!( 611 formatter, 612 "unsupported SQLite legacy {source_kind} schema at user_version {user_version} with catalog SHA-256 {catalog_sha256}" 613 ), 614 Self::LegacyImportFilesystem { operation, .. } => write!( 615 formatter, 616 "SQLite legacy import filesystem operation failed: {operation}" 617 ), 618 Self::BackupFilesystem { operation, .. } => { 619 write!( 620 formatter, 621 "SQLite backup filesystem operation failed: {operation}" 622 ) 623 } 624 } 625 } 626 } 627 628 impl SqliteStorage { 629 /// Opens both governed databases, applying only authorized forward 630 /// migrations and retaining the writer guard for the backend lifetime. 631 #[cfg_attr(coverage_nightly, coverage(off))] 632 pub async fn open(options: OpenOptions) -> Result<Self, Error> { 633 let writer_lock = WriterLock::acquire(options.paths(), options.mode())?; 634 crate::backup::recover_interrupted_restore(options.paths(), options.mode()).await?; 635 options.validate_filesystem()?; 636 let runtime_exists = 637 options 638 .paths() 639 .runtime() 640 .try_exists() 641 .map_err(|source| Error::Inspect { 642 path: options.paths().runtime().to_path_buf(), 643 source, 644 })?; 645 if options.mode().may_create() && !runtime_exists && options.source_generation().is_none() { 646 return Err(Error::SourceGenerationRequired); 647 } 648 options.validate_filesystem()?; 649 650 migration::preflight_existing(&options).await?; 651 652 let runtime_options = connect_options( 653 options.paths().runtime(), 654 options.mode(), 655 options.busy_timeout(), 656 ); 657 let private_options = connect_options( 658 options.paths().private(), 659 options.mode(), 660 options.busy_timeout(), 661 ); 662 let mut runtime_connection = 663 connect(runtime_options.clone(), RUNTIME_DATABASE_NAME).await?; 664 let mut private_connection = 665 connect(private_options.clone(), PRIVATE_DATABASE_NAME).await?; 666 667 migration::migrate_runtime(&mut runtime_connection, options.mode()).await?; 668 migration::migrate_private(&mut private_connection, options.mode()).await?; 669 let generation = active_source_generation( 670 &mut runtime_connection, 671 options.mode(), 672 options.source_generation_bootstrap(), 673 ) 674 .await?; 675 verify_connection( 676 &mut runtime_connection, 677 RUNTIME_DATABASE_NAME, 678 options.busy_timeout(), 679 ) 680 .await?; 681 verify_connection( 682 &mut private_connection, 683 PRIVATE_DATABASE_NAME, 684 options.busy_timeout(), 685 ) 686 .await?; 687 runtime_connection.close().await.map_err(|source| { 688 crate::backend::startup_error( 689 &source, 690 Error::DatabaseCloseFailed { 691 database: RUNTIME_DATABASE_NAME, 692 }, 693 ) 694 })?; 695 private_connection.close().await.map_err(|source| { 696 crate::backend::startup_error( 697 &source, 698 Error::DatabaseCloseFailed { 699 database: PRIVATE_DATABASE_NAME, 700 }, 701 ) 702 })?; 703 704 let runtime_pool = pool(runtime_options, RUNTIME_DATABASE_NAME).await?; 705 let private_pool = pool(private_options, PRIVATE_DATABASE_NAME).await?; 706 verify_pool(&runtime_pool, RUNTIME_DATABASE_NAME, options.busy_timeout()).await?; 707 verify_pool(&private_pool, PRIVATE_DATABASE_NAME, options.busy_timeout()).await?; 708 709 Ok(Self::from_opened( 710 runtime_pool, 711 private_pool, 712 generation, 713 &options, 714 writer_lock, 715 )) 716 } 717 } 718 719 fn connect_options(path: &Path, mode: OpenMode, busy_timeout: Duration) -> SqliteConnectOptions { 720 let mut options = SqliteConnectOptions::new() 721 .filename(path) 722 .read_only(!mode.is_writable()) 723 .create_if_missing(mode.may_create()) 724 .foreign_keys(true) 725 .busy_timeout(busy_timeout) 726 .synchronous(SqliteSynchronous::Full) 727 .pragma("fullfsync", "ON") 728 .disable_statement_logging(); 729 if mode.is_writable() { 730 options = options.journal_mode(SqliteJournalMode::Wal); 731 } 732 options 733 } 734 735 #[cfg_attr(coverage_nightly, coverage(off))] 736 async fn connect( 737 options: SqliteConnectOptions, 738 database: &'static str, 739 ) -> Result<SqliteConnection, Error> { 740 SqliteConnection::connect_with(&options) 741 .await 742 .map_err(|source| map_database_open_error(&source, database)) 743 } 744 745 #[cfg_attr(coverage_nightly, coverage(off))] 746 async fn pool(options: SqliteConnectOptions, database: &'static str) -> Result<SqlitePool, Error> { 747 SqlitePoolOptions::new() 748 .max_connections(MAX_CONNECTIONS_PER_DATABASE) 749 .min_connections(1) 750 .connect_with(options) 751 .await 752 .map_err(|source| map_database_open_error(&source, database)) 753 } 754 755 pub(crate) fn map_database_open_error(source: &sqlx::Error, database: &'static str) -> Error { 756 let is_corrupt = source 757 .as_database_error() 758 .and_then(|error| error.code()) 759 .and_then(|code| code.parse::<i32>().ok()) 760 .is_some_and(|code| matches!(code & 0xff, 11 | 26)); 761 if crate::backend::is_capacity(source) { 762 Error::SpaceInsufficient 763 } else if is_corrupt { 764 Error::DatabaseCorrupt { database } 765 } else { 766 Error::DatabaseOpenFailed { database } 767 } 768 } 769 770 // Inspect the existing generation before migrations without bootstrapping it. 771 // The actual writer transaction still validates or installs it after migration. 772 pub(crate) async fn preflight_source_generation( 773 connection: &mut SqliteConnection, 774 mode: OpenMode, 775 expected: Option<(SourceGeneration, u64)>, 776 ) -> Result<(), Error> { 777 let rows = active_generation_rows(connection).await?; 778 if rows.is_empty() && mode.is_writable() { 779 expected.ok_or(Error::SourceGenerationRequired)?; 780 Ok(()) 781 } else { 782 existing_source_generation(rows.as_slice(), expected).map(|_| ()) 783 } 784 } 785 786 #[cfg_attr(coverage_nightly, coverage(off))] 787 async fn active_source_generation( 788 connection: &mut SqliteConnection, 789 mode: OpenMode, 790 expected: Option<(SourceGeneration, u64)>, 791 ) -> Result<SourceGeneration, Error> { 792 if !mode.is_writable() { 793 let rows = active_generation_rows(connection).await?; 794 return existing_source_generation(rows.as_slice(), expected); 795 } 796 let mut transaction = connection 797 .begin_with("BEGIN IMMEDIATE") 798 .await 799 .map_err(|source| { 800 crate::backend::startup_error(&source, Error::SourceGenerationUnavailable) 801 })?; 802 let rows = active_generation_rows(&mut transaction).await?; 803 let generation = match rows.as_slice() { 804 [] => { 805 let (generation, created_at) = expected.ok_or(Error::SourceGenerationRequired)?; 806 sqlx::query( 807 "INSERT INTO radroots_runtime_source_generations ( 808 generation, sequence_head, state, created_at_unix_ms 809 ) VALUES (?, 0, 'active', ?)", 810 ) 811 .bind(generation.as_bytes().as_slice()) 812 .bind(i64::try_from(created_at).map_err(|_| Error::CorruptSourceGeneration)?) 813 .execute(&mut *transaction) 814 .await 815 .map_err(|source| { 816 crate::backend::startup_error(&source, Error::SourceGenerationUnavailable) 817 })?; 818 generation 819 } 820 [_] => existing_source_generation(rows.as_slice(), expected)?, 821 _ => return Err(Error::CorruptSourceGeneration), 822 }; 823 transaction.commit().await.map_err(|source| { 824 crate::backend::startup_error(&source, Error::SourceGenerationUnavailable) 825 })?; 826 Ok(generation) 827 } 828 829 #[cfg_attr(coverage_nightly, coverage(off))] 830 async fn active_generation_rows( 831 connection: &mut SqliteConnection, 832 ) -> Result<Vec<sqlx::sqlite::SqliteRow>, Error> { 833 sqlx::query( 834 "SELECT generation, created_at_unix_ms 835 FROM radroots_runtime_source_generations 836 WHERE state = 'active' ORDER BY generation", 837 ) 838 .fetch_all(connection) 839 .await 840 .map_err(|source| crate::backend::startup_error(&source, Error::SourceGenerationUnavailable)) 841 } 842 843 fn existing_source_generation( 844 rows: &[sqlx::sqlite::SqliteRow], 845 expected: Option<(SourceGeneration, u64)>, 846 ) -> Result<SourceGeneration, Error> { 847 let [row] = rows else { 848 return if rows.is_empty() { 849 Err(Error::SourceGenerationUnavailable) 850 } else { 851 Err(Error::CorruptSourceGeneration) 852 }; 853 }; 854 let durable = decode_source_generation(row)?; 855 let created_at = u64::try_from( 856 row.try_get::<i64, _>("created_at_unix_ms") 857 .map_err(|_| Error::CorruptSourceGeneration)?, 858 ) 859 .map_err(|_| Error::CorruptSourceGeneration)?; 860 if expected.is_some_and(|candidate| candidate != (durable, created_at)) { 861 Err(Error::SourceGenerationMismatch) 862 } else { 863 Ok(durable) 864 } 865 } 866 867 fn decode_source_generation(row: &sqlx::sqlite::SqliteRow) -> Result<SourceGeneration, Error> { 868 SourceGeneration::new( 869 row.try_get::<Vec<u8>, _>("generation") 870 .map_err(|_| Error::CorruptSourceGeneration)? 871 .try_into() 872 .map_err(|_| Error::CorruptSourceGeneration)?, 873 ) 874 .map_err(|_| Error::CorruptSourceGeneration) 875 } 876 877 #[cfg_attr(coverage_nightly, coverage(off))] 878 async fn verify_pool( 879 pool: &SqlitePool, 880 database: &'static str, 881 busy_timeout: Duration, 882 ) -> Result<(), Error> { 883 let mut connection = pool 884 .acquire() 885 .await 886 .map_err(|source| map_database_open_error(&source, database))?; 887 verify_connection(&mut connection, database, busy_timeout).await 888 } 889 890 #[cfg_attr(coverage_nightly, coverage(off))] 891 async fn verify_connection( 892 connection: &mut SqliteConnection, 893 database: &'static str, 894 busy_timeout: Duration, 895 ) -> Result<(), Error> { 896 let foreign_keys = sqlx::query_scalar::<_, i64>("PRAGMA foreign_keys") 897 .fetch_one(&mut *connection) 898 .await 899 .map_err(|source| { 900 crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) 901 })?; 902 let journal_mode = sqlx::query_scalar::<_, String>("PRAGMA journal_mode") 903 .fetch_one(&mut *connection) 904 .await 905 .map_err(|source| { 906 crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) 907 })?; 908 let configured_busy_timeout = sqlx::query_scalar::<_, i64>("PRAGMA busy_timeout") 909 .fetch_one(&mut *connection) 910 .await 911 .map_err(|source| { 912 crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) 913 })?; 914 let synchronous = sqlx::query_scalar::<_, i64>("PRAGMA synchronous") 915 .fetch_one(&mut *connection) 916 .await 917 .map_err(|source| { 918 crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) 919 })?; 920 let expected_busy_timeout = i64::try_from(busy_timeout.as_millis()) 921 .map_err(|_| Error::ConnectionPolicyMismatch { database })?; 922 let fullfsync = sqlx::query_scalar::<_, i64>("PRAGMA fullfsync") 923 .fetch_one(&mut *connection) 924 .await 925 .map_err(|source| { 926 crate::backend::startup_error(&source, Error::ConnectionPolicyMismatch { database }) 927 })?; 928 if foreign_keys == 1 929 && journal_mode.eq_ignore_ascii_case("wal") 930 && configured_busy_timeout == expected_busy_timeout 931 && synchronous == 2 932 && fullfsync == 1 933 { 934 Ok(()) 935 } else { 936 Err(Error::ConnectionPolicyMismatch { database }) 937 } 938 } 939 940 impl StdError for Error { 941 fn source(&self) -> Option<&(dyn StdError + 'static)> { 942 match self { 943 Self::Inspect { source, .. } 944 | Self::WriterLockOpen { source, .. } 945 | Self::WriterLockFailed { source, .. } 946 | Self::WriterUnlockFailed { source, .. } 947 | Self::RestoreFilesystem { source, .. } 948 | Self::LegacyImportFilesystem { source, .. } 949 | Self::BackupFilesystem { source, .. } => Some(source), 950 _ => None, 951 } 952 } 953 } 954 955 #[cfg(test)] 956 mod error_mapping_tests { 957 use super::*; 958 959 #[test] 960 fn startup_capacity_is_typed_redacted_and_preserves_generic_fallbacks() { 961 for source in [ 962 coded_error(Some("13")), 963 coded_error(Some("269")), 964 sqlx::Error::Io(std::io::Error::new( 965 std::io::ErrorKind::StorageFull, 966 "private path", 967 )), 968 sqlx::Error::Io(std::io::Error::new( 969 std::io::ErrorKind::QuotaExceeded, 970 "private path", 971 )), 972 ] { 973 let mapped = map_database_open_error(&source, RUNTIME_DATABASE_NAME); 974 assert!(matches!(mapped, Error::SpaceInsufficient)); 975 assert_eq!(mapped.to_string(), "storage space is insufficient"); 976 assert_eq!(format!("{mapped:?}"), "SpaceInsufficient"); 977 assert!(matches!( 978 crate::backend::startup_error(&source, Error::SourceGenerationUnavailable), 979 Error::SpaceInsufficient 980 )); 981 } 982 assert!(matches!( 983 crate::backend::startup_error( 984 &sqlx::Error::PoolClosed, 985 Error::SourceGenerationUnavailable 986 ), 987 Error::SourceGenerationUnavailable 988 )); 989 } 990 use sqlx::error::{DatabaseError, ErrorKind}; 991 use std::borrow::Cow; 992 993 #[derive(Debug)] 994 struct CodedDatabaseError(Option<&'static str>); 995 996 impl std::fmt::Display for CodedDatabaseError { 997 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { 998 formatter.write_str("synthetic database error") 999 } 1000 } 1001 1002 impl StdError for CodedDatabaseError {} 1003 1004 impl DatabaseError for CodedDatabaseError { 1005 fn message(&self) -> &str { 1006 "synthetic database error" 1007 } 1008 1009 fn code(&self) -> Option<Cow<'_, str>> { 1010 self.0.map(Cow::Borrowed) 1011 } 1012 1013 fn as_error(&self) -> &(dyn StdError + Send + Sync + 'static) { 1014 self 1015 } 1016 1017 fn as_error_mut(&mut self) -> &mut (dyn StdError + Send + Sync + 'static) { 1018 self 1019 } 1020 1021 fn into_error(self: Box<Self>) -> Box<dyn StdError + Send + Sync + 'static> { 1022 self 1023 } 1024 1025 fn kind(&self) -> ErrorKind { 1026 ErrorKind::Other 1027 } 1028 } 1029 1030 fn coded_error(code: Option<&'static str>) -> sqlx::Error { 1031 sqlx::Error::Database(Box::new(CodedDatabaseError(code))) 1032 } 1033 1034 #[test] 1035 fn database_open_errors_classify_primary_and_extended_corruption_codes() { 1036 assert!(matches!( 1037 map_database_open_error(&coded_error(Some("11")), RUNTIME_DATABASE_NAME), 1038 Error::DatabaseCorrupt { 1039 database: RUNTIME_DATABASE_NAME 1040 } 1041 )); 1042 assert!(matches!( 1043 map_database_open_error(&coded_error(Some("26")), RUNTIME_DATABASE_NAME), 1044 Error::DatabaseCorrupt { 1045 database: RUNTIME_DATABASE_NAME 1046 } 1047 )); 1048 assert!(matches!( 1049 map_database_open_error(&coded_error(Some("267")), RUNTIME_DATABASE_NAME), 1050 Error::DatabaseCorrupt { 1051 database: RUNTIME_DATABASE_NAME 1052 } 1053 )); 1054 assert!(matches!( 1055 map_database_open_error(&coded_error(Some("523")), RUNTIME_DATABASE_NAME), 1056 Error::DatabaseCorrupt { 1057 database: RUNTIME_DATABASE_NAME 1058 } 1059 )); 1060 } 1061 1062 #[test] 1063 fn database_open_errors_fail_closed_for_absent_malformed_and_other_codes() { 1064 assert!(matches!( 1065 map_database_open_error(&coded_error(None), PRIVATE_DATABASE_NAME), 1066 Error::DatabaseOpenFailed { 1067 database: PRIVATE_DATABASE_NAME 1068 } 1069 )); 1070 assert!(matches!( 1071 map_database_open_error(&coded_error(Some("not-a-number")), PRIVATE_DATABASE_NAME), 1072 Error::DatabaseOpenFailed { 1073 database: PRIVATE_DATABASE_NAME 1074 } 1075 )); 1076 assert!(matches!( 1077 map_database_open_error(&coded_error(Some("5")), PRIVATE_DATABASE_NAME), 1078 Error::DatabaseOpenFailed { 1079 database: PRIVATE_DATABASE_NAME 1080 } 1081 )); 1082 assert!(matches!( 1083 map_database_open_error(&sqlx::Error::PoolClosed, PRIVATE_DATABASE_NAME), 1084 Error::DatabaseOpenFailed { 1085 database: PRIVATE_DATABASE_NAME 1086 } 1087 )); 1088 } 1089 } 1090 1091 #[cfg(test)] 1092 #[cfg_attr(coverage_nightly, coverage(off))] 1093 mod policy_tests { 1094 use super::*; 1095 use serde::Deserialize; 1096 1097 const POLICY: &str = include_str!("../../../contracts/storage/connection_policy_v1.toml"); 1098 1099 #[derive(Deserialize)] 1100 struct Policy { 1101 schema_version: u32, 1102 databases: Vec<String>, 1103 max_connections_per_database: u32, 1104 foreign_keys: bool, 1105 journal_mode: String, 1106 synchronous: String, 1107 fullfsync: bool, 1108 busy_timeout_min_ms: u64, 1109 busy_timeout_default_ms: u64, 1110 busy_timeout_max_ms: u64, 1111 fresh_source_generation: String, 1112 read_only_migrations: bool, 1113 raw_handles_public: bool, 1114 } 1115 1116 #[test] 1117 fn implementation_matches_the_governed_connection_policy() { 1118 let policy = toml::from_str::<Policy>(POLICY).expect("connection policy"); 1119 assert_eq!(policy.schema_version, 1); 1120 assert_eq!( 1121 policy.databases, 1122 [RUNTIME_DATABASE_NAME, PRIVATE_DATABASE_NAME] 1123 ); 1124 assert_eq!( 1125 policy.max_connections_per_database, 1126 MAX_CONNECTIONS_PER_DATABASE 1127 ); 1128 assert!(policy.foreign_keys); 1129 assert_eq!(policy.journal_mode, "wal"); 1130 assert_eq!(policy.synchronous, "full"); 1131 assert!(policy.fullfsync); 1132 assert_eq!(policy.busy_timeout_min_ms, 1); 1133 assert_eq!(policy.busy_timeout_default_ms, 5_000); 1134 assert_eq!(policy.busy_timeout_max_ms, 60_000); 1135 assert_eq!( 1136 policy.fresh_source_generation, 1137 "host_supplied_entropy_and_timestamp" 1138 ); 1139 assert!(!policy.read_only_migrations); 1140 assert!(!policy.raw_handles_public); 1141 } 1142 } 1143 1144 #[cfg(test)] 1145 #[cfg_attr(coverage_nightly, coverage(off))] 1146 #[path = "authored_durability_policy_tests.rs"] 1147 mod authored_durability_policy_tests;