legacy.rs (234080B)
1 //! Explicit one-shot legacy import planning and immutable source backup. 2 3 use std::{ 4 collections::BTreeSet, 5 fs::{self, File}, 6 io::Read, 7 path::{Component, Path, PathBuf}, 8 }; 9 10 use radroots_event_codec::Codec; 11 use radroots_storage::{backup::MemberDigest, event::SourceGeneration, status::EventStoreMode}; 12 use sha2::{Digest, Sha256}; 13 use sqlx::{Connection, Row, SqliteConnection, sqlite::SqliteConnectOptions}; 14 15 use crate::{Error, SqliteStorage}; 16 17 const LEGACY_SOURCE_MAX: usize = 4; 18 /// Maximum predecessor event rows converted by one staging transaction. 19 pub const LEGACY_STAGE_PAGE_LIMIT_MAX: u16 = 256; 20 const LEGACY_MANIFEST: &str = "manifest.v1"; 21 const EVENT_STORE_LEDGER: &str = "radroots_event_store_schema_migrations"; 22 const EVENT_STORE_LEDGER_DDL: &str = "CREATE TABLE radroots_event_store_schema_migrations ( 23 version INTEGER PRIMARY KEY NOT NULL CHECK (version > 0), 24 name TEXT NOT NULL UNIQUE CHECK (length(name) > 0), 25 up_sha256 TEXT NOT NULL CHECK (length(up_sha256) = 64 AND up_sha256 NOT GLOB '*[^0-9a-f]*'), 26 down_sha256 TEXT NOT NULL CHECK (length(down_sha256) = 64 AND down_sha256 NOT GLOB '*[^0-9a-f]*'), 27 schema_sha256 TEXT NOT NULL CHECK (length(schema_sha256) = 64 AND schema_sha256 NOT GLOB '*[^0-9a-f]*') 28 ) STRICT, WITHOUT ROWID"; 29 30 const EVENT_STORE_MIGRATIONS: [LegacyEventMigration; 4] = [ 31 LegacyEventMigration { 32 version: 1, 33 name: "event_store", 34 up_sha256: "4c03906a1cffd418a48d40907aa9a1ca51bb41766cff7250c4dfc7c2fd6eddde", 35 down_sha256: "fa84d587f657f601947eaeb9cd239c962a48f6fcdce723588476e8d22f3c1f53", 36 schema_sha256: "5b1f92779640f1a2dbd75e37a96996bda6c8be58883190f69eb3eced22a48f03", 37 }, 38 LegacyEventMigration { 39 version: 2, 40 name: "nip09", 41 up_sha256: "0c1730ff36eaebd285f9c0c94b9b7346af60266afa55c24a18e30446d369581a", 42 down_sha256: "c51a099d9501f1e692c13d2226296a68ed9e6bfa5e8e46b2f12c6574dbe59e31", 43 schema_sha256: "1fee6b2bb8cdc4602d9c89fecd97c3f51312b9a4339dbf5049b04c692ba50b12", 44 }, 45 LegacyEventMigration { 46 version: 3, 47 name: "food_availability_projection", 48 up_sha256: "4e7edfb981b25f76055efc7802ec30b4034eeae9b9c0809ea4ea7c574678748a", 49 down_sha256: "29d663320109d9dd0df6a00b6a53d8d988438d01f7a66960a9d4ba3482ffffb8", 50 schema_sha256: "dd12467e04addcbddb5ea0f386c12a8ac05ef5ebaaf949f24dd2c62745f5aaac", 51 }, 52 LegacyEventMigration { 53 version: 4, 54 name: "source_maintenance", 55 up_sha256: "ab2724188f8d08c897eebea2533a635e7c74282a25e84e4c0c37e78b08837a43", 56 down_sha256: "fe44fd53c51545c08ea479b385e6781079dab70fc63da2a3c205d727a00ce860", 57 schema_sha256: "074f85b663444ac150239ecd8441ea4a96ad83a798a55e22d2e5e2f7ee943a8c", 58 }, 59 ]; 60 const OUTBOX_CATALOG_SHA256: &str = 61 "e7eeba00de78ec6d990c620e7c056018166e8a00bb703e472ef6f67a00870293"; 62 const PRIVATE_CATALOG_SHA256: &str = 63 "5aa3664e3ecb4461bde0589c3e8f73be041b715a83a67add366af975f827614e"; 64 const STUDIO_CATALOG_SHA256: &str = 65 "3e13518dba056db82090a336833618ca1bc3a44ba49067967ad9bf4c22768193"; 66 67 #[derive(Clone, Copy)] 68 struct LegacyEventMigration { 69 version: u32, 70 name: &'static str, 71 up_sha256: &'static str, 72 down_sha256: &'static str, 73 schema_sha256: &'static str, 74 } 75 76 /// Stable identity for one forward-only legacy import attempt. 77 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 78 pub struct LegacyImportId([u8; 16]); 79 80 impl LegacyImportId { 81 /// Creates a non-zero caller-supplied import identity. 82 pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> { 83 if bytes_are_zero(&bytes) { 84 Err(Error::InvalidLegacyImportPlan) 85 } else { 86 Ok(Self(bytes)) 87 } 88 } 89 90 /// Returns the stable identity bytes. 91 pub const fn as_bytes(&self) -> &[u8; 16] { 92 &self.0 93 } 94 } 95 96 /// Supported predecessor database families accepted by the one-shot planner. 97 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 98 #[non_exhaustive] 99 pub enum LegacySourceKind { 100 EventStore, 101 Outbox, 102 Private, 103 Studio, 104 } 105 106 impl LegacySourceKind { 107 /// Returns the stable policy identifier for this source family. 108 pub const fn as_str(self) -> &'static str { 109 match self { 110 Self::EventStore => "event_store", 111 Self::Outbox => "outbox", 112 Self::Private => "private", 113 Self::Studio => "studio", 114 } 115 } 116 117 const fn backup_file_name(self) -> &'static str { 118 match self { 119 Self::EventStore => "event_store.sqlite", 120 Self::Outbox => "outbox.sqlite", 121 Self::Private => "private.sqlite", 122 Self::Studio => "studio.sqlite", 123 } 124 } 125 } 126 127 /// Exact supported predecessor schema selected by fail-closed classification. 128 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 129 #[non_exhaustive] 130 pub enum LegacySchema { 131 EventStoreV1, 132 EventStoreV2, 133 EventStoreV3, 134 EventStoreV4, 135 OutboxV1, 136 PrivateV1, 137 StudioV1HostHandoff, 138 } 139 140 impl LegacySchema { 141 /// Returns the stable schema identifier recorded by the importer. 142 pub const fn as_str(self) -> &'static str { 143 match self { 144 Self::EventStoreV1 => "event_store_v1", 145 Self::EventStoreV2 => "event_store_v2", 146 Self::EventStoreV3 => "event_store_v3", 147 Self::EventStoreV4 => "event_store_v4", 148 Self::OutboxV1 => "outbox_v1", 149 Self::PrivateV1 => "private_v1", 150 Self::StudioV1HostHandoff => "studio_v1_host_handoff", 151 } 152 } 153 154 /// Returns whether this source is converted into owned storage or handed to its host. 155 pub const fn disposition(self) -> LegacyImportDisposition { 156 match self { 157 Self::StudioV1HostHandoff => LegacyImportDisposition::HostHandoff, 158 Self::EventStoreV1 159 | Self::EventStoreV2 160 | Self::EventStoreV3 161 | Self::EventStoreV4 162 | Self::OutboxV1 163 | Self::PrivateV1 => LegacyImportDisposition::Import, 164 } 165 } 166 } 167 168 /// Required destination behavior for one classified predecessor source. 169 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 170 #[non_exhaustive] 171 pub enum LegacyImportDisposition { 172 Import, 173 HostHandoff, 174 } 175 176 impl LegacyImportDisposition { 177 /// Returns the stable durable-journal value. 178 pub const fn as_str(self) -> &'static str { 179 match self { 180 Self::Import => "import", 181 Self::HostHandoff => "host_handoff", 182 } 183 } 184 } 185 186 /// Exact schema evidence for one classified predecessor snapshot. 187 #[derive(Clone, Debug, Eq, PartialEq)] 188 pub struct LegacySourceClassification { 189 kind: LegacySourceKind, 190 schema: LegacySchema, 191 user_version: u32, 192 catalog_sha256: MemberDigest, 193 } 194 195 impl LegacySourceClassification { 196 /// Returns the predecessor source family. 197 pub const fn kind(&self) -> LegacySourceKind { 198 self.kind 199 } 200 201 /// Returns the exact supported predecessor schema. 202 pub const fn schema(&self) -> LegacySchema { 203 self.schema 204 } 205 206 /// Returns the observed SQLite application user version. 207 pub const fn user_version(&self) -> u32 { 208 self.user_version 209 } 210 211 /// Returns the exact governed SQLite schema-catalog fingerprint. 212 pub const fn catalog_sha256(&self) -> MemberDigest { 213 self.catalog_sha256 214 } 215 } 216 217 /// Fully reverified classification of a prepared import evidence bundle. 218 #[derive(Clone, Debug, Eq, PartialEq)] 219 pub struct ClassifiedLegacyImport { 220 prepared: PreparedLegacyImport, 221 sources: Vec<LegacySourceClassification>, 222 } 223 224 impl ClassifiedLegacyImport { 225 /// Returns the stable import-attempt identity. 226 pub const fn import_id(&self) -> LegacyImportId { 227 self.prepared.import_id() 228 } 229 230 /// Returns the exact destination storage generation. 231 pub const fn target_generation(&self) -> SourceGeneration { 232 self.prepared.target_generation() 233 } 234 235 /// Returns the reverified finalized evidence bundle. 236 pub fn bundle_path(&self) -> &Path { 237 self.prepared.bundle_path() 238 } 239 240 /// Returns exact classifications in stable source-family order. 241 pub fn sources(&self) -> &[LegacySourceClassification] { 242 &self.sources 243 } 244 } 245 246 /// Durable whole-import lifecycle state. 247 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 248 #[non_exhaustive] 249 pub enum LegacyImportState { 250 Classified, 251 Staging, 252 Ready, 253 Committing, 254 Complete, 255 } 256 257 impl LegacyImportState { 258 /// Returns the stable SQLite journal value. 259 pub const fn as_str(self) -> &'static str { 260 match self { 261 Self::Classified => "classified", 262 Self::Staging => "staging", 263 Self::Ready => "ready", 264 Self::Committing => "committing", 265 Self::Complete => "complete", 266 } 267 } 268 } 269 270 /// Durable per-source staging lifecycle state. 271 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 272 #[non_exhaustive] 273 pub enum LegacyImportMemberState { 274 Pending, 275 Staging, 276 Ready, 277 Complete, 278 } 279 280 impl LegacyImportMemberState { 281 /// Returns the stable SQLite journal value. 282 pub const fn as_str(self) -> &'static str { 283 match self { 284 Self::Pending => "pending", 285 Self::Staging => "staging", 286 Self::Ready => "ready", 287 Self::Complete => "complete", 288 } 289 } 290 } 291 292 /// Durable recovery state for one classified predecessor source. 293 #[derive(Clone, Debug, Eq, PartialEq)] 294 pub struct LegacyImportMemberJournal { 295 classification: LegacySourceClassification, 296 state: LegacyImportMemberState, 297 resume_cursor: Option<Vec<u8>>, 298 staged_row_count: u64, 299 updated_at_unix_ms: u64, 300 } 301 302 impl LegacyImportMemberJournal { 303 /// Returns the exact source classification bound to this member. 304 pub const fn classification(&self) -> &LegacySourceClassification { 305 &self.classification 306 } 307 308 /// Returns the durable staging state. 309 pub const fn state(&self) -> LegacyImportMemberState { 310 self.state 311 } 312 313 /// Returns the opaque source-specific resume cursor. 314 pub fn resume_cursor(&self) -> Option<&[u8]> { 315 self.resume_cursor.as_deref() 316 } 317 318 /// Returns the number of rows durably staged so far. 319 pub const fn staged_row_count(&self) -> u64 { 320 self.staged_row_count 321 } 322 323 /// Returns the positive last-update timestamp supplied by the host. 324 pub const fn updated_at_unix_ms(&self) -> u64 { 325 self.updated_at_unix_ms 326 } 327 } 328 329 /// Exact durable recovery journal for one target-bound import. 330 #[derive(Clone, Debug, Eq, PartialEq)] 331 pub struct LegacyImportJournal { 332 import_id: LegacyImportId, 333 target_generation: SourceGeneration, 334 manifest_sha256: MemberDigest, 335 classification_sha256: MemberDigest, 336 state: LegacyImportState, 337 started_at_unix_ms: u64, 338 updated_at_unix_ms: u64, 339 completed_at_unix_ms: Option<u64>, 340 members: Vec<LegacyImportMemberJournal>, 341 } 342 343 /// Result of one bounded, durable legacy event-store staging transaction. 344 #[derive(Clone, Debug, Eq, PartialEq)] 345 pub struct LegacyEventStagePage { 346 staged_rows: u16, 347 staged_row_count: u64, 348 resume_cursor: Option<[u8; 8]>, 349 complete: bool, 350 } 351 352 /// Stable predecessor table order for bounded legacy outbox graph staging. 353 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 354 #[non_exhaustive] 355 pub enum LegacyOutboxTable { 356 Operations, 357 Events, 358 DeliveryPlans, 359 DeliveryTargets, 360 DeliveryAttempts, 361 } 362 363 impl LegacyOutboxTable { 364 /// Returns the stable staging table-kind value. 365 pub const fn as_str(self) -> &'static str { 366 match self { 367 Self::Operations => "operations", 368 Self::Events => "events", 369 Self::DeliveryPlans => "delivery_plans", 370 Self::DeliveryTargets => "delivery_targets", 371 Self::DeliveryAttempts => "delivery_attempts", 372 } 373 } 374 375 const fn code(self) -> u8 { 376 match self { 377 Self::Operations => 1, 378 Self::Events => 2, 379 Self::DeliveryPlans => 3, 380 Self::DeliveryTargets => 4, 381 Self::DeliveryAttempts => 5, 382 } 383 } 384 385 const fn next(self) -> Option<Self> { 386 match self { 387 Self::Operations => Some(Self::Events), 388 Self::Events => Some(Self::DeliveryPlans), 389 Self::DeliveryPlans => Some(Self::DeliveryTargets), 390 Self::DeliveryTargets => Some(Self::DeliveryAttempts), 391 Self::DeliveryAttempts => None, 392 } 393 } 394 } 395 396 /// Result of one bounded legacy outbox table staging transaction. 397 #[derive(Clone, Debug, Eq, PartialEq)] 398 pub struct LegacyOutboxStagePage { 399 table: LegacyOutboxTable, 400 staged_rows: u16, 401 staged_row_count: u64, 402 resume_cursor: [u8; 9], 403 complete: bool, 404 } 405 406 /// Stable predecessor table order for protected legacy private-store staging. 407 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 408 #[non_exhaustive] 409 pub enum LegacyPrivateTable { 410 Metadata, 411 WrappedProfileKeys, 412 SigningSecrets, 413 FarmLocations, 414 TradeArtifacts, 415 CursorKeys, 416 Nip46Sessions, 417 RotationProgress, 418 } 419 420 impl LegacyPrivateTable { 421 pub const fn as_str(self) -> &'static str { 422 match self { 423 Self::Metadata => "metadata", 424 Self::WrappedProfileKeys => "wrapped_profile_keys", 425 Self::SigningSecrets => "signing_secrets", 426 Self::FarmLocations => "farm_locations", 427 Self::TradeArtifacts => "trade_artifacts", 428 Self::CursorKeys => "cursor_keys", 429 Self::Nip46Sessions => "nip46_sessions", 430 Self::RotationProgress => "rotation_progress", 431 } 432 } 433 434 const fn code(self) -> u8 { 435 match self { 436 Self::Metadata => 1, 437 Self::WrappedProfileKeys => 2, 438 Self::SigningSecrets => 3, 439 Self::FarmLocations => 4, 440 Self::TradeArtifacts => 5, 441 Self::CursorKeys => 6, 442 Self::Nip46Sessions => 7, 443 Self::RotationProgress => 8, 444 } 445 } 446 447 const fn next(self) -> Option<Self> { 448 match self { 449 Self::Metadata => Some(Self::WrappedProfileKeys), 450 Self::WrappedProfileKeys => Some(Self::SigningSecrets), 451 Self::SigningSecrets => Some(Self::FarmLocations), 452 Self::FarmLocations => Some(Self::TradeArtifacts), 453 Self::TradeArtifacts => Some(Self::CursorKeys), 454 Self::CursorKeys => Some(Self::Nip46Sessions), 455 Self::Nip46Sessions => Some(Self::RotationProgress), 456 Self::RotationProgress => None, 457 } 458 } 459 } 460 461 /// Result of one recoverable protected legacy private-store page. 462 #[derive(Clone, Debug, Eq, PartialEq)] 463 pub struct LegacyPrivateStagePage { 464 table: LegacyPrivateTable, 465 staged_rows: u16, 466 staged_row_count: u64, 467 resume_cursor: Vec<u8>, 468 complete: bool, 469 } 470 471 /// Immutable host-owned handoff descriptor for one classified Studio snapshot. 472 #[derive(Clone, Debug, Eq, PartialEq)] 473 pub struct LegacyStudioHandoff { 474 import_id: LegacyImportId, 475 evidence_path: PathBuf, 476 byte_length: u64, 477 source_sha256: MemberDigest, 478 catalog_sha256: MemberDigest, 479 handoff_sha256: MemberDigest, 480 } 481 482 impl LegacyStudioHandoff { 483 /// Returns the import attempt bound to this handoff. 484 pub const fn import_id(&self) -> LegacyImportId { 485 self.import_id 486 } 487 488 /// Returns the immutable backed-up Studio database offered to the host. 489 pub fn evidence_path(&self) -> &Path { 490 &self.evidence_path 491 } 492 493 /// Returns the exact backed-up Studio database length. 494 pub const fn byte_length(&self) -> u64 { 495 self.byte_length 496 } 497 498 /// Returns the exact backed-up Studio database digest. 499 pub const fn source_sha256(&self) -> MemberDigest { 500 self.source_sha256 501 } 502 503 /// Returns the exact classified Studio schema-catalog digest. 504 pub const fn catalog_sha256(&self) -> MemberDigest { 505 self.catalog_sha256 506 } 507 508 /// Returns the deterministic identity the host must acknowledge. 509 pub const fn handoff_sha256(&self) -> MemberDigest { 510 self.handoff_sha256 511 } 512 } 513 514 /// Host-supplied proof that a specific Studio handoff was durably accepted. 515 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 516 pub struct LegacyStudioHandoffReceipt { 517 handoff_sha256: MemberDigest, 518 host_commitment_sha256: MemberDigest, 519 } 520 521 /// Snapshot-consistent proof that every legacy member is ready to commit. 522 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 523 pub struct LegacyImportValidation { 524 imported_row_count: u64, 525 validation_sha256: MemberDigest, 526 } 527 528 /// Durable receipt for one fully sealed, forward-only legacy import. 529 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 530 pub struct LegacyImportCommitReceipt { 531 validation_sha256: MemberDigest, 532 imported_row_count: u64, 533 completed_at_unix_ms: u64, 534 } 535 536 impl LegacyImportCommitReceipt { 537 /// Returns the exact validation identity sealed by both databases. 538 pub const fn validation_sha256(&self) -> MemberDigest { 539 self.validation_sha256 540 } 541 542 /// Returns the exact retained SDK-owned predecessor row count. 543 pub const fn imported_row_count(&self) -> u64 { 544 self.imported_row_count 545 } 546 547 /// Returns the positive host-supplied completion timestamp. 548 pub const fn completed_at_unix_ms(&self) -> u64 { 549 self.completed_at_unix_ms 550 } 551 } 552 553 impl LegacyImportValidation { 554 /// Returns the exact number of predecessor rows staged for SDK-owned storage. 555 pub const fn imported_row_count(&self) -> u64 { 556 self.imported_row_count 557 } 558 559 /// Returns the deterministic commit identity for all staged rows and receipts. 560 pub const fn validation_sha256(&self) -> MemberDigest { 561 self.validation_sha256 562 } 563 } 564 565 impl LegacyStudioHandoffReceipt { 566 /// Binds an exact handoff to a non-zero host-owned durable commitment. 567 pub const fn new( 568 handoff_sha256: MemberDigest, 569 host_commitment_sha256: MemberDigest, 570 ) -> Result<Self, Error> { 571 if bytes_are_zero(host_commitment_sha256.as_bytes()) { 572 Err(Error::InvalidLegacyImportStageRequest) 573 } else { 574 Ok(Self { 575 handoff_sha256, 576 host_commitment_sha256, 577 }) 578 } 579 } 580 581 /// Returns the acknowledged handoff identity. 582 pub const fn handoff_sha256(&self) -> MemberDigest { 583 self.handoff_sha256 584 } 585 586 /// Returns the opaque host-owned durable commitment. 587 pub const fn host_commitment_sha256(&self) -> MemberDigest { 588 self.host_commitment_sha256 589 } 590 } 591 592 impl LegacyPrivateStagePage { 593 pub const fn table(&self) -> LegacyPrivateTable { 594 self.table 595 } 596 pub const fn staged_rows(&self) -> u16 { 597 self.staged_rows 598 } 599 pub const fn staged_row_count(&self) -> u64 { 600 self.staged_row_count 601 } 602 pub fn resume_cursor(&self) -> &[u8] { 603 &self.resume_cursor 604 } 605 pub const fn is_complete(&self) -> bool { 606 self.complete 607 } 608 } 609 610 impl LegacyOutboxStagePage { 611 /// Returns the predecessor table processed by this page. 612 pub const fn table(&self) -> LegacyOutboxTable { 613 self.table 614 } 615 616 /// Returns rows newly staged by this transaction. 617 pub const fn staged_rows(&self) -> u16 { 618 self.staged_rows 619 } 620 621 /// Returns the cumulative durable outbox graph row count. 622 pub const fn staged_row_count(&self) -> u64 { 623 self.staged_row_count 624 } 625 626 /// Returns the exact table-discriminated predecessor cursor. 627 pub const fn resume_cursor(&self) -> &[u8; 9] { 628 &self.resume_cursor 629 } 630 631 /// Reports whether all five predecessor tables reached their exact end. 632 pub const fn is_complete(&self) -> bool { 633 self.complete 634 } 635 } 636 637 impl LegacyEventStagePage { 638 /// Returns rows newly converted by this transaction. 639 pub const fn staged_rows(&self) -> u16 { 640 self.staged_rows 641 } 642 643 /// Returns the total durable event staging row count for this import. 644 pub const fn staged_row_count(&self) -> u64 { 645 self.staged_row_count 646 } 647 648 /// Returns the exact big-endian predecessor `event_envelopes.seq` cursor. 649 pub const fn resume_cursor(&self) -> Option<&[u8; 8]> { 650 self.resume_cursor.as_ref() 651 } 652 653 /// Reports whether the source member has reached its exact end. 654 pub const fn is_complete(&self) -> bool { 655 self.complete 656 } 657 } 658 659 impl LegacyImportJournal { 660 /// Returns the stable import identity. 661 pub const fn import_id(&self) -> LegacyImportId { 662 self.import_id 663 } 664 665 /// Returns the exact destination storage generation. 666 pub const fn target_generation(&self) -> SourceGeneration { 667 self.target_generation 668 } 669 670 /// Returns the exact finalized evidence-manifest digest. 671 pub const fn manifest_sha256(&self) -> MemberDigest { 672 self.manifest_sha256 673 } 674 675 /// Returns the exact ordered-classification digest. 676 pub const fn classification_sha256(&self) -> MemberDigest { 677 self.classification_sha256 678 } 679 680 /// Returns the durable whole-import state. 681 pub const fn state(&self) -> LegacyImportState { 682 self.state 683 } 684 685 /// Returns the positive host-supplied start timestamp. 686 pub const fn started_at_unix_ms(&self) -> u64 { 687 self.started_at_unix_ms 688 } 689 690 /// Returns the last positive host-supplied update timestamp. 691 pub const fn updated_at_unix_ms(&self) -> u64 { 692 self.updated_at_unix_ms 693 } 694 695 /// Returns the host-supplied completion timestamp once terminal. 696 pub const fn completed_at_unix_ms(&self) -> Option<u64> { 697 self.completed_at_unix_ms 698 } 699 700 /// Returns one exact durable row per classified source. 701 pub fn members(&self) -> &[LegacyImportMemberJournal] { 702 &self.members 703 } 704 } 705 706 /// One explicitly typed existing predecessor database. 707 #[derive(Clone, Debug, Eq, PartialEq)] 708 pub struct LegacySource { 709 kind: LegacySourceKind, 710 path: PathBuf, 711 } 712 713 impl LegacySource { 714 /// Binds a source family to one absolute existing regular SQLite file. 715 pub fn new(kind: LegacySourceKind, path: impl Into<PathBuf>) -> Result<Self, Error> { 716 let path = path.into(); 717 validate_source_path(&path)?; 718 Ok(Self { kind, path }) 719 } 720 721 /// Returns the declared predecessor database family. 722 pub const fn kind(&self) -> LegacySourceKind { 723 self.kind 724 } 725 726 /// Returns the exact caller-supplied predecessor database path. 727 pub fn path(&self) -> &Path { 728 &self.path 729 } 730 } 731 732 /// Immutable authority for one pre-backed-up, forward-only import attempt. 733 #[derive(Clone, Debug, Eq, PartialEq)] 734 pub struct LegacyImportPlan { 735 import_id: LegacyImportId, 736 sources: Vec<LegacySource>, 737 backup_root: PathBuf, 738 requested_at_unix_ms: u64, 739 } 740 741 impl LegacyImportPlan { 742 /// Creates a deterministic import plan with one source per family. 743 pub fn new( 744 import_id: LegacyImportId, 745 mut sources: Vec<LegacySource>, 746 backup_root: impl Into<PathBuf>, 747 requested_at_unix_ms: u64, 748 ) -> Result<Self, Error> { 749 let backup_root = backup_root.into(); 750 crate::backup::validate_backup_root(&backup_root)?; 751 if requested_at_unix_ms == 0 || sources.is_empty() || sources.len() > LEGACY_SOURCE_MAX { 752 return Err(Error::InvalidLegacyImportPlan); 753 } 754 let mut kinds = BTreeSet::new(); 755 let mut paths = BTreeSet::new(); 756 for source in &sources { 757 validate_source_path(source.path())?; 758 if !kinds.insert(source.kind()) || !paths.insert(source.path().to_path_buf()) { 759 return Err(Error::InvalidLegacyImportPlan); 760 } 761 } 762 sources.sort_by_key(LegacySource::kind); 763 Ok(Self { 764 import_id, 765 sources, 766 backup_root, 767 requested_at_unix_ms, 768 }) 769 } 770 771 /// Returns the stable import-attempt identity. 772 pub const fn import_id(&self) -> LegacyImportId { 773 self.import_id 774 } 775 776 /// Returns sources in stable source-family order. 777 pub fn sources(&self) -> &[LegacySource] { 778 &self.sources 779 } 780 781 /// Returns the existing host-owned directory for immutable import evidence. 782 pub fn backup_root(&self) -> &Path { 783 &self.backup_root 784 } 785 786 /// Returns the positive host-supplied import request timestamp. 787 pub const fn requested_at_unix_ms(&self) -> u64 { 788 self.requested_at_unix_ms 789 } 790 } 791 792 /// Exact immutable evidence for one backed-up predecessor member. 793 #[derive(Clone, Debug, Eq, PartialEq)] 794 pub struct LegacySourceSnapshot { 795 kind: LegacySourceKind, 796 relative_path: String, 797 byte_length: u64, 798 sha256: MemberDigest, 799 } 800 801 impl LegacySourceSnapshot { 802 /// Returns the predecessor database family. 803 pub const fn kind(&self) -> LegacySourceKind { 804 self.kind 805 } 806 807 /// Returns the stable bundle-relative snapshot path. 808 pub fn relative_path(&self) -> &str { 809 self.relative_path.as_str() 810 } 811 812 /// Returns the exact snapshot length. 813 pub const fn byte_length(&self) -> u64 { 814 self.byte_length 815 } 816 817 /// Returns the exact snapshot SHA-256 digest. 818 pub const fn sha256(&self) -> MemberDigest { 819 self.sha256 820 } 821 } 822 823 /// Durable result of the mandatory pre-import source backup. 824 #[derive(Clone, Debug, Eq, PartialEq)] 825 pub struct PreparedLegacyImport { 826 import_id: LegacyImportId, 827 target_generation: SourceGeneration, 828 bundle_path: PathBuf, 829 manifest_byte_length: u64, 830 manifest_sha256: MemberDigest, 831 snapshots: Vec<LegacySourceSnapshot>, 832 } 833 834 impl PreparedLegacyImport { 835 /// Returns the stable import-attempt identity. 836 pub const fn import_id(&self) -> LegacyImportId { 837 self.import_id 838 } 839 840 /// Returns the exact destination storage generation. 841 pub const fn target_generation(&self) -> SourceGeneration { 842 self.target_generation 843 } 844 845 /// Returns the finalized immutable evidence bundle. 846 pub fn bundle_path(&self) -> &Path { 847 &self.bundle_path 848 } 849 850 /// Returns the exact manifest length. 851 pub const fn manifest_byte_length(&self) -> u64 { 852 self.manifest_byte_length 853 } 854 855 /// Returns the exact manifest SHA-256 digest. 856 pub const fn manifest_sha256(&self) -> MemberDigest { 857 self.manifest_sha256 858 } 859 860 /// Returns the exact evidence inventory in stable source-family order. 861 pub fn snapshots(&self) -> &[LegacySourceSnapshot] { 862 &self.snapshots 863 } 864 } 865 866 impl SqliteStorage { 867 /// Captures and verifies every legacy source before any import mutation. 868 #[cfg_attr(coverage_nightly, coverage(off))] 869 pub async fn prepare_legacy_import( 870 &self, 871 plan: &LegacyImportPlan, 872 ) -> Result<PreparedLegacyImport, Error> { 873 self.lifecycle 874 .require_open() 875 .map_err(|_| Error::BackupBackendUnavailable)?; 876 if self.mode != EventStoreMode::ReadWrite { 877 return Err(Error::RestoreRequiresWritableStorage); 878 } 879 for source in plan.sources() { 880 validate_source_path(source.path())?; 881 if let Some(paths) = self.paths.as_deref() { 882 for owned in [paths.runtime(), paths.private()] { 883 if paths_refer_to_same_file(source.path(), owned)? { 884 return Err(Error::InvalidLegacySource(source.path().to_path_buf())); 885 } 886 } 887 } 888 } 889 890 let layout = LegacyBackupLayout::new(plan); 891 layout.create()?; 892 let mut snapshots = Vec::with_capacity(plan.sources().len()); 893 for source in plan.sources() { 894 let destination = layout.staging.join(source.kind().backup_file_name()); 895 capture_legacy_source(source, &destination).await?; 896 snapshots.push(snapshot(source.kind(), &destination)?); 897 } 898 let manifest_path = layout.staging.join(LEGACY_MANIFEST); 899 write_manifest(plan, self.generation, &snapshots, &manifest_path)?; 900 let (manifest_byte_length, manifest_sha256) = file_digest(&manifest_path)?; 901 sync_directory(&layout.staging, "sync legacy import staging bundle")?; 902 fs::rename(&layout.staging, &layout.finalized).map_err(|source| { 903 Error::LegacyImportFilesystem { 904 operation: "finalize legacy import backup bundle", 905 source, 906 } 907 })?; 908 sync_directory(plan.backup_root(), "sync legacy import backup root")?; 909 Ok(PreparedLegacyImport { 910 import_id: plan.import_id(), 911 target_generation: self.generation, 912 bundle_path: layout.finalized, 913 manifest_byte_length, 914 manifest_sha256, 915 snapshots, 916 }) 917 } 918 919 /// Revalidates a prepared bundle and classifies every exact predecessor schema. 920 #[cfg_attr(coverage_nightly, coverage(off))] 921 pub async fn classify_legacy_import( 922 &self, 923 prepared: &PreparedLegacyImport, 924 ) -> Result<ClassifiedLegacyImport, Error> { 925 self.lifecycle 926 .require_open() 927 .map_err(|_| Error::BackupBackendUnavailable)?; 928 if self.mode != EventStoreMode::ReadWrite { 929 return Err(Error::RestoreRequiresWritableStorage); 930 } 931 if prepared.target_generation() != self.generation { 932 return Err(Error::LegacyImportTargetMismatch); 933 } 934 verify_prepared_evidence(prepared).await?; 935 let mut sources = Vec::with_capacity(prepared.snapshots().len()); 936 for snapshot in prepared.snapshots() { 937 sources.push( 938 classify_snapshot( 939 snapshot.kind(), 940 &prepared.bundle_path().join(snapshot.relative_path()), 941 ) 942 .await?, 943 ); 944 } 945 Ok(ClassifiedLegacyImport { 946 prepared: prepared.clone(), 947 sources, 948 }) 949 } 950 951 /// Atomically creates or resumes the exact durable journal for a classification. 952 #[cfg_attr(coverage_nightly, coverage(off))] 953 pub async fn begin_legacy_import( 954 &self, 955 classified: &ClassifiedLegacyImport, 956 started_at_unix_ms: u64, 957 ) -> Result<LegacyImportJournal, Error> { 958 self.require_legacy_import_writer(classified.target_generation())?; 959 if started_at_unix_ms == 0 || classified.sources().is_empty() { 960 return Err(Error::InvalidLegacyImportJournal); 961 } 962 verify_prepared_evidence(&classified.prepared).await?; 963 let started_at = 964 i64::try_from(started_at_unix_ms).map_err(|_| Error::InvalidLegacyImportJournal)?; 965 let classification_sha256 = classification_digest(classified); 966 let mut transaction = self 967 .pool 968 .begin_with("BEGIN IMMEDIATE") 969 .await 970 .map_err(|_| Error::LegacyImportJournalFailed)?; 971 let existing = sqlx::query_scalar::<_, i64>( 972 "SELECT COUNT(*) FROM radroots_runtime_legacy_imports WHERE import_id = ?", 973 ) 974 .bind(classified.import_id().as_bytes().as_slice()) 975 .fetch_one(&mut *transaction) 976 .await 977 .map_err(|_| Error::LegacyImportJournalFailed)?; 978 if existing == 0 { 979 let active = sqlx::query_scalar::<_, i64>( 980 "SELECT COUNT(*) FROM radroots_runtime_legacy_imports WHERE target_generation = ?", 981 ) 982 .bind(classified.target_generation().as_bytes().as_slice()) 983 .fetch_one(&mut *transaction) 984 .await 985 .map_err(|_| Error::LegacyImportJournalFailed)?; 986 if active != 0 { 987 transaction 988 .rollback() 989 .await 990 .map_err(|_| Error::LegacyImportJournalFailed)?; 991 return Err(Error::LegacyImportConflict); 992 } 993 sqlx::query( 994 "INSERT INTO radroots_runtime_legacy_imports( 995 import_id, target_generation, manifest_sha256, 996 classification_sha256, state, started_at_ms, 997 updated_at_ms, completed_at_ms 998 ) VALUES (?, ?, ?, ?, 'classified', ?, ?, NULL)", 999 ) 1000 .bind(classified.import_id().as_bytes().as_slice()) 1001 .bind(classified.target_generation().as_bytes().as_slice()) 1002 .bind(classified.prepared.manifest_sha256().as_bytes().as_slice()) 1003 .bind(classification_sha256.as_bytes().as_slice()) 1004 .bind(started_at) 1005 .bind(started_at) 1006 .execute(&mut *transaction) 1007 .await 1008 .map_err(|_| Error::LegacyImportJournalFailed)?; 1009 for source in classified.sources() { 1010 sqlx::query( 1011 "INSERT INTO radroots_runtime_legacy_import_members( 1012 import_id, source_kind, legacy_schema, disposition, 1013 catalog_sha256, state, resume_cursor, staged_row_count, 1014 updated_at_ms 1015 ) VALUES (?, ?, ?, ?, ?, 'pending', NULL, 0, ?)", 1016 ) 1017 .bind(classified.import_id().as_bytes().as_slice()) 1018 .bind(source.kind().as_str()) 1019 .bind(source.schema().as_str()) 1020 .bind(source.schema().disposition().as_str()) 1021 .bind(source.catalog_sha256().as_bytes().as_slice()) 1022 .bind(started_at) 1023 .execute(&mut *transaction) 1024 .await 1025 .map_err(|_| Error::LegacyImportJournalFailed)?; 1026 } 1027 } 1028 transaction 1029 .commit() 1030 .await 1031 .map_err(|_| Error::LegacyImportJournalFailed)?; 1032 let journal = self 1033 .legacy_import_journal(classified.import_id()) 1034 .await? 1035 .ok_or(Error::InvalidLegacyImportJournal)?; 1036 if journal_matches_classified(&journal, classified, classification_sha256) { 1037 Ok(journal) 1038 } else { 1039 Err(Error::LegacyImportConflict) 1040 } 1041 } 1042 1043 /// Reads exact durable recovery state without advancing the importer. 1044 #[cfg_attr(coverage_nightly, coverage(off))] 1045 pub async fn legacy_import_journal( 1046 &self, 1047 import_id: LegacyImportId, 1048 ) -> Result<Option<LegacyImportJournal>, Error> { 1049 self.lifecycle 1050 .require_open() 1051 .map_err(|_| Error::BackupBackendUnavailable)?; 1052 let mut transaction = self 1053 .pool 1054 .begin_with("BEGIN") 1055 .await 1056 .map_err(|_| Error::LegacyImportJournalFailed)?; 1057 let row = sqlx::query( 1058 "SELECT import_id, target_generation, manifest_sha256, 1059 classification_sha256, state, started_at_ms, updated_at_ms, 1060 completed_at_ms 1061 FROM radroots_runtime_legacy_imports WHERE import_id = ?", 1062 ) 1063 .bind(import_id.as_bytes().as_slice()) 1064 .fetch_optional(&mut *transaction) 1065 .await 1066 .map_err(|_| Error::LegacyImportJournalFailed)?; 1067 let Some(row) = row else { 1068 transaction 1069 .commit() 1070 .await 1071 .map_err(|_| Error::LegacyImportJournalFailed)?; 1072 return Ok(None); 1073 }; 1074 let durable_import_id = decode_import_id( 1075 row.try_get("import_id") 1076 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1077 )?; 1078 if durable_import_id != import_id { 1079 return Err(Error::InvalidLegacyImportJournal); 1080 } 1081 let target_generation = decode_generation( 1082 row.try_get("target_generation") 1083 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1084 )?; 1085 if target_generation != self.generation { 1086 return Err(Error::InvalidLegacyImportJournal); 1087 } 1088 let manifest_sha256 = decode_digest( 1089 row.try_get("manifest_sha256") 1090 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1091 )?; 1092 let classification_sha256 = decode_digest( 1093 row.try_get("classification_sha256") 1094 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1095 )?; 1096 let state = parse_import_state( 1097 row.try_get::<String, _>("state") 1098 .map_err(|_| Error::InvalidLegacyImportJournal)? 1099 .as_str(), 1100 )?; 1101 let started_at_unix_ms = decode_positive_time( 1102 row.try_get("started_at_ms") 1103 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1104 )?; 1105 let updated_at_unix_ms = decode_positive_time( 1106 row.try_get("updated_at_ms") 1107 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1108 )?; 1109 let completed_at_unix_ms = row 1110 .try_get::<Option<i64>, _>("completed_at_ms") 1111 .map_err(|_| Error::InvalidLegacyImportJournal)? 1112 .map(decode_positive_time) 1113 .transpose()?; 1114 let member_rows = sqlx::query( 1115 "SELECT source_kind, legacy_schema, disposition, catalog_sha256, 1116 state, resume_cursor, staged_row_count, updated_at_ms 1117 FROM radroots_runtime_legacy_import_members 1118 WHERE import_id = ? ORDER BY source_kind", 1119 ) 1120 .bind(import_id.as_bytes().as_slice()) 1121 .fetch_all(&mut *transaction) 1122 .await 1123 .map_err(|_| Error::LegacyImportJournalFailed)?; 1124 if member_rows.is_empty() || member_rows.len() > LEGACY_SOURCE_MAX { 1125 return Err(Error::InvalidLegacyImportJournal); 1126 } 1127 let mut members = Vec::with_capacity(member_rows.len()); 1128 for row in member_rows { 1129 let kind = parse_source_kind( 1130 row.try_get::<String, _>("source_kind") 1131 .map_err(|_| Error::InvalidLegacyImportJournal)? 1132 .as_str(), 1133 )?; 1134 let schema = parse_legacy_schema( 1135 row.try_get::<String, _>("legacy_schema") 1136 .map_err(|_| Error::InvalidLegacyImportJournal)? 1137 .as_str(), 1138 )?; 1139 let disposition = row 1140 .try_get::<String, _>("disposition") 1141 .map_err(|_| Error::InvalidLegacyImportJournal)?; 1142 if disposition != schema.disposition().as_str() || schema_source_kind(schema) != kind { 1143 return Err(Error::InvalidLegacyImportJournal); 1144 } 1145 members.push(LegacyImportMemberJournal { 1146 classification: LegacySourceClassification { 1147 kind, 1148 schema, 1149 user_version: expected_user_version(schema), 1150 catalog_sha256: decode_digest( 1151 row.try_get("catalog_sha256") 1152 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1153 )?, 1154 }, 1155 state: parse_member_state( 1156 row.try_get::<String, _>("state") 1157 .map_err(|_| Error::InvalidLegacyImportJournal)? 1158 .as_str(), 1159 )?, 1160 resume_cursor: row 1161 .try_get("resume_cursor") 1162 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1163 staged_row_count: u64::try_from( 1164 row.try_get::<i64, _>("staged_row_count") 1165 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1166 ) 1167 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1168 updated_at_unix_ms: decode_positive_time( 1169 row.try_get("updated_at_ms") 1170 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1171 )?, 1172 }); 1173 } 1174 if updated_at_unix_ms < started_at_unix_ms 1175 || members 1176 .iter() 1177 .any(|member| member.updated_at_unix_ms() < started_at_unix_ms) 1178 || !journal_member_states_are_consistent(state, &members) 1179 { 1180 return Err(Error::InvalidLegacyImportJournal); 1181 } 1182 transaction 1183 .commit() 1184 .await 1185 .map_err(|_| Error::LegacyImportJournalFailed)?; 1186 Ok(Some(LegacyImportJournal { 1187 import_id, 1188 target_generation, 1189 manifest_sha256, 1190 classification_sha256, 1191 state, 1192 started_at_unix_ms, 1193 updated_at_unix_ms, 1194 completed_at_unix_ms, 1195 members, 1196 })) 1197 } 1198 1199 /// Converts one bounded page of an exact legacy event store into isolated staging. 1200 #[cfg_attr(coverage_nightly, coverage(off))] 1201 pub async fn stage_legacy_events( 1202 &self, 1203 classified: &ClassifiedLegacyImport, 1204 limit: u16, 1205 updated_at_unix_ms: u64, 1206 ) -> Result<LegacyEventStagePage, Error> { 1207 self.require_legacy_import_writer(classified.target_generation())?; 1208 if limit == 0 || limit > LEGACY_STAGE_PAGE_LIMIT_MAX || updated_at_unix_ms == 0 { 1209 return Err(Error::InvalidLegacyImportStageRequest); 1210 } 1211 verify_prepared_evidence(&classified.prepared).await?; 1212 let classification_sha256 = classification_digest(classified); 1213 let journal = self 1214 .legacy_import_journal(classified.import_id()) 1215 .await? 1216 .ok_or(Error::InvalidLegacyImportJournal)?; 1217 if !journal_matches_classified(&journal, classified, classification_sha256) { 1218 return Err(Error::LegacyImportConflict); 1219 } 1220 let classification = classified 1221 .sources() 1222 .iter() 1223 .find(|source| source.kind() == LegacySourceKind::EventStore) 1224 .ok_or(Error::LegacyImportConflict)?; 1225 if !matches!( 1226 classification.schema(), 1227 LegacySchema::EventStoreV1 1228 | LegacySchema::EventStoreV2 1229 | LegacySchema::EventStoreV3 1230 | LegacySchema::EventStoreV4 1231 ) { 1232 return Err(Error::LegacyImportConflict); 1233 } 1234 let snapshot = classified 1235 .prepared 1236 .snapshots() 1237 .iter() 1238 .find(|snapshot| snapshot.kind() == LegacySourceKind::EventStore) 1239 .ok_or(Error::LegacyImportConflict)?; 1240 let source_path = classified.bundle_path().join(snapshot.relative_path()); 1241 let updated_at = i64::try_from(updated_at_unix_ms) 1242 .map_err(|_| Error::InvalidLegacyImportStageRequest)?; 1243 1244 let mut transaction = self 1245 .pool 1246 .begin_with("BEGIN IMMEDIATE") 1247 .await 1248 .map_err(|_| Error::LegacyImportStagingFailed)?; 1249 let import_row = sqlx::query( 1250 "SELECT state, updated_at_ms FROM radroots_runtime_legacy_imports 1251 WHERE import_id = ? AND target_generation = ? 1252 AND manifest_sha256 = ? AND classification_sha256 = ?", 1253 ) 1254 .bind(classified.import_id().as_bytes().as_slice()) 1255 .bind(classified.target_generation().as_bytes().as_slice()) 1256 .bind(classified.prepared.manifest_sha256().as_bytes().as_slice()) 1257 .bind(classification_sha256.as_bytes().as_slice()) 1258 .fetch_optional(&mut *transaction) 1259 .await 1260 .map_err(|_| Error::LegacyImportStagingFailed)? 1261 .ok_or(Error::LegacyImportConflict)?; 1262 let import_state = parse_import_state( 1263 import_row 1264 .try_get::<String, _>("state") 1265 .map_err(|_| Error::InvalidLegacyImportJournal)? 1266 .as_str(), 1267 )?; 1268 let import_updated_at = import_row 1269 .try_get::<i64, _>("updated_at_ms") 1270 .map_err(|_| Error::InvalidLegacyImportJournal)?; 1271 if updated_at < import_updated_at 1272 || !matches!( 1273 import_state, 1274 LegacyImportState::Classified 1275 | LegacyImportState::Staging 1276 | LegacyImportState::Ready 1277 ) 1278 { 1279 return Err(Error::LegacyImportConflict); 1280 } 1281 let member_row = sqlx::query( 1282 "SELECT state, resume_cursor, staged_row_count, updated_at_ms 1283 FROM radroots_runtime_legacy_import_members 1284 WHERE import_id = ? AND source_kind = 'event_store'", 1285 ) 1286 .bind(classified.import_id().as_bytes().as_slice()) 1287 .fetch_optional(&mut *transaction) 1288 .await 1289 .map_err(|_| Error::LegacyImportStagingFailed)? 1290 .ok_or(Error::LegacyImportConflict)?; 1291 let member_state = parse_member_state( 1292 member_row 1293 .try_get::<String, _>("state") 1294 .map_err(|_| Error::InvalidLegacyImportJournal)? 1295 .as_str(), 1296 )?; 1297 let durable_cursor = member_row 1298 .try_get::<Option<Vec<u8>>, _>("resume_cursor") 1299 .map_err(|_| Error::InvalidLegacyImportJournal)?; 1300 let staged_row_count = u64::try_from( 1301 member_row 1302 .try_get::<i64, _>("staged_row_count") 1303 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1304 ) 1305 .map_err(|_| Error::InvalidLegacyImportJournal)?; 1306 let member_updated_at = member_row 1307 .try_get::<i64, _>("updated_at_ms") 1308 .map_err(|_| Error::InvalidLegacyImportJournal)?; 1309 let resume_sequence = decode_event_stage_cursor(durable_cursor.as_deref())?; 1310 if updated_at < member_updated_at { 1311 return Err(Error::LegacyImportConflict); 1312 } 1313 if member_state == LegacyImportMemberState::Ready { 1314 transaction 1315 .commit() 1316 .await 1317 .map_err(|_| Error::LegacyImportStagingFailed)?; 1318 return Ok(LegacyEventStagePage { 1319 staged_rows: 0, 1320 staged_row_count, 1321 resume_cursor: durable_cursor 1322 .as_deref() 1323 .map(decode_exact_event_stage_cursor) 1324 .transpose()?, 1325 complete: true, 1326 }); 1327 } 1328 if !matches!( 1329 member_state, 1330 LegacyImportMemberState::Pending | LegacyImportMemberState::Staging 1331 ) { 1332 return Err(Error::LegacyImportConflict); 1333 } 1334 1335 if import_state == LegacyImportState::Classified { 1336 sqlx::query( 1337 "UPDATE radroots_runtime_legacy_imports 1338 SET state = 'staging', updated_at_ms = ? WHERE import_id = ?", 1339 ) 1340 .bind(updated_at) 1341 .bind(classified.import_id().as_bytes().as_slice()) 1342 .execute(&mut *transaction) 1343 .await 1344 .map_err(|_| Error::LegacyImportStagingFailed)?; 1345 } 1346 if member_state == LegacyImportMemberState::Pending { 1347 sqlx::query( 1348 "UPDATE radroots_runtime_legacy_import_members 1349 SET state = 'staging', updated_at_ms = ? 1350 WHERE import_id = ? AND source_kind = 'event_store'", 1351 ) 1352 .bind(updated_at) 1353 .bind(classified.import_id().as_bytes().as_slice()) 1354 .execute(&mut *transaction) 1355 .await 1356 .map_err(|_| Error::LegacyImportStagingFailed)?; 1357 } 1358 1359 let mut source = SqliteConnection::connect_with( 1360 &SqliteConnectOptions::new() 1361 .filename(&source_path) 1362 .read_only(true), 1363 ) 1364 .await 1365 .map_err(|_| Error::LegacyImportStagingFailed)?; 1366 let rows = sqlx::query( 1367 "SELECT seq, event_id, raw_json, verification_status, contract_status, 1368 projection_eligible, inserted_at_ms, updated_at_ms 1369 FROM event_envelopes WHERE seq > ? ORDER BY seq LIMIT ?", 1370 ) 1371 .bind(resume_sequence) 1372 .bind(i64::from(limit) + 1) 1373 .fetch_all(&mut source) 1374 .await 1375 .map_err(|_| Error::LegacyImportStagingFailed)?; 1376 source 1377 .close() 1378 .await 1379 .map_err(|_| Error::LegacyImportStagingFailed)?; 1380 let complete = rows.len() <= usize::from(limit); 1381 let rows = rows.into_iter().take(usize::from(limit)); 1382 let mut last_sequence = resume_sequence; 1383 let mut newly_staged = 0_u16; 1384 for row in rows { 1385 let converted = convert_legacy_event_row(&row)?; 1386 sqlx::query( 1387 "INSERT INTO radroots_runtime_legacy_event_staging( 1388 import_id, source_kind, legacy_sequence, event_id, signed_event, 1389 legacy_verification_status, legacy_contract_status, 1390 legacy_projection_eligible, legacy_inserted_at_ms, 1391 legacy_updated_at_ms 1392 ) VALUES (?, 'event_store', ?, ?, ?, ?, ?, ?, ?, ?)", 1393 ) 1394 .bind(classified.import_id().as_bytes().as_slice()) 1395 .bind(converted.sequence) 1396 .bind(converted.event_id.as_slice()) 1397 .bind(converted.signed_event.as_slice()) 1398 .bind(converted.verification_status) 1399 .bind(converted.contract_status) 1400 .bind(converted.projection_eligible) 1401 .bind(converted.inserted_at_ms) 1402 .bind(converted.updated_at_ms) 1403 .execute(&mut *transaction) 1404 .await 1405 .map_err(|_| Error::LegacyImportStagingFailed)?; 1406 last_sequence = converted.sequence; 1407 newly_staged = newly_staged 1408 .checked_add(1) 1409 .ok_or(Error::LegacyImportStagingFailed)?; 1410 } 1411 let total = staged_row_count 1412 .checked_add(u64::from(newly_staged)) 1413 .ok_or(Error::LegacyImportStagingFailed)?; 1414 let cursor = (last_sequence > 0).then(|| encode_event_stage_cursor(last_sequence)); 1415 sqlx::query( 1416 "UPDATE radroots_runtime_legacy_import_members 1417 SET state = ?, resume_cursor = ?, staged_row_count = ?, updated_at_ms = ? 1418 WHERE import_id = ? AND source_kind = 'event_store'", 1419 ) 1420 .bind(if complete { "ready" } else { "staging" }) 1421 .bind(cursor.as_ref().map(<[u8; 8]>::as_slice)) 1422 .bind(i64::try_from(total).map_err(|_| Error::LegacyImportStagingFailed)?) 1423 .bind(updated_at) 1424 .bind(classified.import_id().as_bytes().as_slice()) 1425 .execute(&mut *transaction) 1426 .await 1427 .map_err(|_| Error::LegacyImportStagingFailed)?; 1428 let pending_members = sqlx::query_scalar::<_, i64>( 1429 "SELECT COUNT(*) FROM radroots_runtime_legacy_import_members 1430 WHERE import_id = ? AND state != 'ready'", 1431 ) 1432 .bind(classified.import_id().as_bytes().as_slice()) 1433 .fetch_one(&mut *transaction) 1434 .await 1435 .map_err(|_| Error::LegacyImportStagingFailed)?; 1436 sqlx::query( 1437 "UPDATE radroots_runtime_legacy_imports SET state = ?, updated_at_ms = ? 1438 WHERE import_id = ?", 1439 ) 1440 .bind(if pending_members == 0 { 1441 "ready" 1442 } else { 1443 "staging" 1444 }) 1445 .bind(updated_at) 1446 .bind(classified.import_id().as_bytes().as_slice()) 1447 .execute(&mut *transaction) 1448 .await 1449 .map_err(|_| Error::LegacyImportStagingFailed)?; 1450 transaction 1451 .commit() 1452 .await 1453 .map_err(|_| Error::LegacyImportStagingFailed)?; 1454 Ok(LegacyEventStagePage { 1455 staged_rows: newly_staged, 1456 staged_row_count: total, 1457 resume_cursor: cursor, 1458 complete, 1459 }) 1460 } 1461 1462 /// Converts one bounded table page from an exact legacy outbox graph. 1463 #[cfg_attr(coverage_nightly, coverage(off))] 1464 pub async fn stage_legacy_outbox( 1465 &self, 1466 classified: &ClassifiedLegacyImport, 1467 limit: u16, 1468 updated_at_unix_ms: u64, 1469 ) -> Result<LegacyOutboxStagePage, Error> { 1470 self.require_legacy_import_writer(classified.target_generation())?; 1471 if limit == 0 || limit > LEGACY_STAGE_PAGE_LIMIT_MAX || updated_at_unix_ms == 0 { 1472 return Err(Error::InvalidLegacyImportStageRequest); 1473 } 1474 verify_prepared_evidence(&classified.prepared).await?; 1475 let classification_sha256 = classification_digest(classified); 1476 let journal = self 1477 .legacy_import_journal(classified.import_id()) 1478 .await? 1479 .ok_or(Error::InvalidLegacyImportJournal)?; 1480 if !journal_matches_classified(&journal, classified, classification_sha256) 1481 || !classified.sources().iter().any(|source| { 1482 source.kind() == LegacySourceKind::Outbox 1483 && source.schema() == LegacySchema::OutboxV1 1484 }) 1485 { 1486 return Err(Error::LegacyImportConflict); 1487 } 1488 let snapshot = classified 1489 .prepared 1490 .snapshots() 1491 .iter() 1492 .find(|snapshot| snapshot.kind() == LegacySourceKind::Outbox) 1493 .ok_or(Error::LegacyImportConflict)?; 1494 let source_path = classified.bundle_path().join(snapshot.relative_path()); 1495 let updated_at = i64::try_from(updated_at_unix_ms) 1496 .map_err(|_| Error::InvalidLegacyImportStageRequest)?; 1497 let mut transaction = self 1498 .pool 1499 .begin_with("BEGIN IMMEDIATE") 1500 .await 1501 .map_err(|_| Error::LegacyImportStagingFailed)?; 1502 let import_row = sqlx::query( 1503 "SELECT state, updated_at_ms FROM radroots_runtime_legacy_imports 1504 WHERE import_id = ? AND target_generation = ? 1505 AND manifest_sha256 = ? AND classification_sha256 = ?", 1506 ) 1507 .bind(classified.import_id().as_bytes().as_slice()) 1508 .bind(classified.target_generation().as_bytes().as_slice()) 1509 .bind(classified.prepared.manifest_sha256().as_bytes().as_slice()) 1510 .bind(classification_sha256.as_bytes().as_slice()) 1511 .fetch_optional(&mut *transaction) 1512 .await 1513 .map_err(|_| Error::LegacyImportStagingFailed)? 1514 .ok_or(Error::LegacyImportConflict)?; 1515 let import_state = parse_import_state( 1516 import_row 1517 .try_get::<String, _>("state") 1518 .map_err(|_| Error::InvalidLegacyImportJournal)? 1519 .as_str(), 1520 )?; 1521 let import_updated_at = import_row 1522 .try_get::<i64, _>("updated_at_ms") 1523 .map_err(|_| Error::InvalidLegacyImportJournal)?; 1524 if updated_at < import_updated_at 1525 || !matches!( 1526 import_state, 1527 LegacyImportState::Classified 1528 | LegacyImportState::Staging 1529 | LegacyImportState::Ready 1530 ) 1531 { 1532 return Err(Error::LegacyImportConflict); 1533 } 1534 let member_row = sqlx::query( 1535 "SELECT state, resume_cursor, staged_row_count, updated_at_ms 1536 FROM radroots_runtime_legacy_import_members 1537 WHERE import_id = ? AND source_kind = 'outbox'", 1538 ) 1539 .bind(classified.import_id().as_bytes().as_slice()) 1540 .fetch_optional(&mut *transaction) 1541 .await 1542 .map_err(|_| Error::LegacyImportStagingFailed)? 1543 .ok_or(Error::LegacyImportConflict)?; 1544 let member_state = parse_member_state( 1545 member_row 1546 .try_get::<String, _>("state") 1547 .map_err(|_| Error::InvalidLegacyImportJournal)? 1548 .as_str(), 1549 )?; 1550 let durable_cursor = member_row 1551 .try_get::<Option<Vec<u8>>, _>("resume_cursor") 1552 .map_err(|_| Error::InvalidLegacyImportJournal)?; 1553 let (table, after) = decode_outbox_stage_cursor(durable_cursor.as_deref())?; 1554 let staged_row_count = u64::try_from( 1555 member_row 1556 .try_get::<i64, _>("staged_row_count") 1557 .map_err(|_| Error::InvalidLegacyImportJournal)?, 1558 ) 1559 .map_err(|_| Error::InvalidLegacyImportJournal)?; 1560 let member_updated_at = member_row 1561 .try_get::<i64, _>("updated_at_ms") 1562 .map_err(|_| Error::InvalidLegacyImportJournal)?; 1563 if updated_at < member_updated_at { 1564 return Err(Error::LegacyImportConflict); 1565 } 1566 if member_state == LegacyImportMemberState::Ready { 1567 let cursor = durable_cursor 1568 .as_deref() 1569 .map(decode_exact_outbox_stage_cursor) 1570 .transpose()? 1571 .ok_or(Error::InvalidLegacyImportJournal)?; 1572 transaction 1573 .commit() 1574 .await 1575 .map_err(|_| Error::LegacyImportStagingFailed)?; 1576 return Ok(LegacyOutboxStagePage { 1577 table: LegacyOutboxTable::DeliveryAttempts, 1578 staged_rows: 0, 1579 staged_row_count, 1580 resume_cursor: cursor, 1581 complete: true, 1582 }); 1583 } 1584 if !matches!( 1585 member_state, 1586 LegacyImportMemberState::Pending | LegacyImportMemberState::Staging 1587 ) { 1588 return Err(Error::LegacyImportConflict); 1589 } 1590 if import_state == LegacyImportState::Classified { 1591 sqlx::query( 1592 "UPDATE radroots_runtime_legacy_imports 1593 SET state = 'staging', updated_at_ms = ? WHERE import_id = ?", 1594 ) 1595 .bind(updated_at) 1596 .bind(classified.import_id().as_bytes().as_slice()) 1597 .execute(&mut *transaction) 1598 .await 1599 .map_err(|_| Error::LegacyImportStagingFailed)?; 1600 } 1601 if member_state == LegacyImportMemberState::Pending { 1602 sqlx::query( 1603 "UPDATE radroots_runtime_legacy_import_members 1604 SET state = 'staging', updated_at_ms = ? 1605 WHERE import_id = ? AND source_kind = 'outbox'", 1606 ) 1607 .bind(updated_at) 1608 .bind(classified.import_id().as_bytes().as_slice()) 1609 .execute(&mut *transaction) 1610 .await 1611 .map_err(|_| Error::LegacyImportStagingFailed)?; 1612 } 1613 let mut source = SqliteConnection::connect_with( 1614 &SqliteConnectOptions::new() 1615 .filename(&source_path) 1616 .read_only(true), 1617 ) 1618 .await 1619 .map_err(|_| Error::LegacyImportStagingFailed)?; 1620 let rows = sqlx::query(outbox_stage_query(table)) 1621 .bind(after) 1622 .bind(i64::from(limit) + 1) 1623 .fetch_all(&mut source) 1624 .await 1625 .map_err(|_| Error::LegacyImportStagingFailed)?; 1626 source 1627 .close() 1628 .await 1629 .map_err(|_| Error::LegacyImportStagingFailed)?; 1630 let table_complete = rows.len() <= usize::from(limit); 1631 let mut last_id = after; 1632 let mut newly_staged = 0_u16; 1633 for row in rows.into_iter().take(usize::from(limit)) { 1634 let legacy_id = row 1635 .try_get::<i64, _>("legacy_id") 1636 .map_err(|_| Error::LegacyImportStagingFailed)?; 1637 let parent_legacy_id = row 1638 .try_get::<Option<i64>, _>("parent_legacy_id") 1639 .map_err(|_| Error::LegacyImportStagingFailed)?; 1640 let related_legacy_id = row 1641 .try_get::<Option<i64>, _>("related_legacy_id") 1642 .map_err(|_| Error::LegacyImportStagingFailed)?; 1643 let record_json = row 1644 .try_get::<Vec<u8>, _>("record_json") 1645 .map_err(|_| Error::LegacyImportStagingFailed)?; 1646 if legacy_id <= last_id || record_json.is_empty() { 1647 return Err(Error::LegacyImportRowInvalid { 1648 source_kind: "outbox", 1649 legacy_sequence: legacy_id, 1650 }); 1651 } 1652 sqlx::query( 1653 "INSERT INTO radroots_runtime_legacy_outbox_staging( 1654 import_id, source_kind, table_kind, legacy_id, 1655 parent_legacy_id, related_legacy_id, record_json 1656 ) VALUES (?, 'outbox', ?, ?, ?, ?, ?)", 1657 ) 1658 .bind(classified.import_id().as_bytes().as_slice()) 1659 .bind(table.as_str()) 1660 .bind(legacy_id) 1661 .bind(parent_legacy_id) 1662 .bind(related_legacy_id) 1663 .bind(record_json) 1664 .execute(&mut *transaction) 1665 .await 1666 .map_err(|_| Error::LegacyImportStagingFailed)?; 1667 last_id = legacy_id; 1668 newly_staged += 1; 1669 } 1670 let complete = table_complete && table.next().is_none(); 1671 let next_cursor = if table_complete { 1672 encode_outbox_stage_cursor( 1673 table.next().unwrap_or(table), 1674 if complete { last_id } else { 0 }, 1675 ) 1676 } else { 1677 encode_outbox_stage_cursor(table, last_id) 1678 }; 1679 let total = staged_row_count 1680 .checked_add(u64::from(newly_staged)) 1681 .ok_or(Error::LegacyImportStagingFailed)?; 1682 sqlx::query( 1683 "UPDATE radroots_runtime_legacy_import_members 1684 SET state = ?, resume_cursor = ?, staged_row_count = ?, updated_at_ms = ? 1685 WHERE import_id = ? AND source_kind = 'outbox'", 1686 ) 1687 .bind(if complete { "ready" } else { "staging" }) 1688 .bind(next_cursor.as_slice()) 1689 .bind(i64::try_from(total).map_err(|_| Error::LegacyImportStagingFailed)?) 1690 .bind(updated_at) 1691 .bind(classified.import_id().as_bytes().as_slice()) 1692 .execute(&mut *transaction) 1693 .await 1694 .map_err(|_| Error::LegacyImportStagingFailed)?; 1695 let pending_members = sqlx::query_scalar::<_, i64>( 1696 "SELECT COUNT(*) FROM radroots_runtime_legacy_import_members 1697 WHERE import_id = ? AND state != 'ready'", 1698 ) 1699 .bind(classified.import_id().as_bytes().as_slice()) 1700 .fetch_one(&mut *transaction) 1701 .await 1702 .map_err(|_| Error::LegacyImportStagingFailed)?; 1703 sqlx::query( 1704 "UPDATE radroots_runtime_legacy_imports SET state = ?, updated_at_ms = ? 1705 WHERE import_id = ?", 1706 ) 1707 .bind(if pending_members == 0 { 1708 "ready" 1709 } else { 1710 "staging" 1711 }) 1712 .bind(updated_at) 1713 .bind(classified.import_id().as_bytes().as_slice()) 1714 .execute(&mut *transaction) 1715 .await 1716 .map_err(|_| Error::LegacyImportStagingFailed)?; 1717 transaction 1718 .commit() 1719 .await 1720 .map_err(|_| Error::LegacyImportStagingFailed)?; 1721 Ok(LegacyOutboxStagePage { 1722 table, 1723 staged_rows: newly_staged, 1724 staged_row_count: total, 1725 resume_cursor: next_cursor, 1726 complete, 1727 }) 1728 } 1729 1730 /// Stages one recoverable page of an exact predecessor private store. 1731 #[cfg_attr(coverage_nightly, coverage(off))] 1732 pub async fn stage_legacy_private( 1733 &self, 1734 classified: &ClassifiedLegacyImport, 1735 limit: u16, 1736 updated_at_unix_ms: u64, 1737 ) -> Result<LegacyPrivateStagePage, Error> { 1738 self.require_legacy_import_writer(classified.target_generation())?; 1739 if limit == 0 || limit > LEGACY_STAGE_PAGE_LIMIT_MAX || updated_at_unix_ms == 0 { 1740 return Err(Error::InvalidLegacyImportStageRequest); 1741 } 1742 verify_prepared_evidence(&classified.prepared).await?; 1743 let classification_sha256 = classification_digest(classified); 1744 let journal = self 1745 .legacy_import_journal(classified.import_id()) 1746 .await? 1747 .ok_or(Error::InvalidLegacyImportJournal)?; 1748 if !journal_matches_classified(&journal, classified, classification_sha256) 1749 || !classified.sources().iter().any(|source| { 1750 source.kind() == LegacySourceKind::Private 1751 && source.schema() == LegacySchema::PrivateV1 1752 }) 1753 { 1754 return Err(Error::LegacyImportConflict); 1755 } 1756 let member = journal 1757 .members() 1758 .iter() 1759 .find(|member| member.classification().kind() == LegacySourceKind::Private) 1760 .ok_or(Error::LegacyImportConflict)?; 1761 let (table, after) = decode_private_stage_cursor(member.resume_cursor())?; 1762 if member.state() == LegacyImportMemberState::Ready { 1763 return Ok(LegacyPrivateStagePage { 1764 table: LegacyPrivateTable::RotationProgress, 1765 staged_rows: 0, 1766 staged_row_count: member.staged_row_count(), 1767 resume_cursor: member 1768 .resume_cursor() 1769 .ok_or(Error::InvalidLegacyImportJournal)? 1770 .to_vec(), 1771 complete: true, 1772 }); 1773 } 1774 if !matches!( 1775 member.state(), 1776 LegacyImportMemberState::Pending | LegacyImportMemberState::Staging 1777 ) || updated_at_unix_ms < member.updated_at_unix_ms() 1778 || updated_at_unix_ms < journal.updated_at_unix_ms() 1779 { 1780 return Err(Error::LegacyImportConflict); 1781 } 1782 let updated_at = i64::try_from(updated_at_unix_ms) 1783 .map_err(|_| Error::InvalidLegacyImportStageRequest)?; 1784 if member.state() == LegacyImportMemberState::Pending { 1785 let mut tx = self 1786 .pool 1787 .begin_with("BEGIN IMMEDIATE") 1788 .await 1789 .map_err(|_| Error::LegacyImportStagingFailed)?; 1790 if journal.state() == LegacyImportState::Classified { 1791 sqlx::query("UPDATE radroots_runtime_legacy_imports SET state = 'staging', updated_at_ms = ? WHERE import_id = ? AND state = 'classified'") 1792 .bind(updated_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 1793 } 1794 let changed = sqlx::query("UPDATE radroots_runtime_legacy_import_members SET state = 'staging', updated_at_ms = ? WHERE import_id = ? AND source_kind = 'private' AND state = 'pending'") 1795 .bind(updated_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 1796 if changed.rows_affected() != 1 { 1797 return Err(Error::LegacyImportConflict); 1798 } 1799 tx.commit() 1800 .await 1801 .map_err(|_| Error::LegacyImportStagingFailed)?; 1802 } 1803 let snapshot = classified 1804 .prepared 1805 .snapshots() 1806 .iter() 1807 .find(|snapshot| snapshot.kind() == LegacySourceKind::Private) 1808 .ok_or(Error::LegacyImportConflict)?; 1809 let mut source = SqliteConnection::connect_with( 1810 &SqliteConnectOptions::new() 1811 .filename(classified.bundle_path().join(snapshot.relative_path())) 1812 .read_only(true), 1813 ) 1814 .await 1815 .map_err(|_| Error::LegacyImportStagingFailed)?; 1816 let rows = sqlx::query(private_stage_query(table)) 1817 .bind(after.as_str()) 1818 .bind(i64::from(limit) + 1) 1819 .fetch_all(&mut source) 1820 .await 1821 .map_err(|_| Error::LegacyImportStagingFailed)?; 1822 source 1823 .close() 1824 .await 1825 .map_err(|_| Error::LegacyImportStagingFailed)?; 1826 let table_complete = rows.len() <= usize::from(limit); 1827 let page_rows = rows 1828 .into_iter() 1829 .take(usize::from(limit)) 1830 .collect::<Vec<_>>(); 1831 let mut last_key = after.clone(); 1832 let mut private_tx = self 1833 .private_pool 1834 .begin_with("BEGIN IMMEDIATE") 1835 .await 1836 .map_err(|_| Error::LegacyImportStagingFailed)?; 1837 for row in &page_rows { 1838 let key = row 1839 .try_get::<String, _>("key_cursor") 1840 .map_err(|_| Error::LegacyImportStagingFailed)?; 1841 let parent = row 1842 .try_get::<Option<i64>, _>("parent_key_version") 1843 .map_err(|_| Error::LegacyImportStagingFailed)?; 1844 let record = row 1845 .try_get::<Vec<u8>, _>("record_json") 1846 .map_err(|_| Error::LegacyImportStagingFailed)?; 1847 if key <= last_key || key.len() > 1024 || record.is_empty() { 1848 return Err(Error::LegacyImportRowInvalid { 1849 source_kind: "private", 1850 legacy_sequence: 0, 1851 }); 1852 } 1853 sqlx::query("INSERT OR IGNORE INTO radroots_private_legacy_import_staging(import_id, table_kind, key_cursor, parent_key_version, record_json) VALUES (?, ?, ?, ?, ?)") 1854 .bind(classified.import_id().as_bytes().as_slice()).bind(table.as_str()).bind(&key).bind(parent).bind(&record).execute(&mut *private_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 1855 let existing = sqlx::query("SELECT parent_key_version, record_json FROM radroots_private_legacy_import_staging WHERE import_id = ? AND table_kind = ? AND key_cursor = ?") 1856 .bind(classified.import_id().as_bytes().as_slice()).bind(table.as_str()).bind(&key).fetch_one(&mut *private_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 1857 if existing 1858 .try_get::<Option<i64>, _>("parent_key_version") 1859 .map_err(|_| Error::LegacyImportStagingFailed)? 1860 != parent 1861 || existing 1862 .try_get::<Vec<u8>, _>("record_json") 1863 .map_err(|_| Error::LegacyImportStagingFailed)? 1864 != record 1865 { 1866 return Err(Error::LegacyImportConflict); 1867 } 1868 last_key = key; 1869 } 1870 private_tx 1871 .commit() 1872 .await 1873 .map_err(|_| Error::LegacyImportStagingFailed)?; 1874 let complete = table_complete && table.next().is_none(); 1875 let next_cursor = if table_complete { 1876 encode_private_stage_cursor( 1877 table.next().unwrap_or(table), 1878 if complete { &last_key } else { "" }, 1879 ) 1880 } else { 1881 encode_private_stage_cursor(table, &last_key) 1882 }; 1883 let total = member 1884 .staged_row_count() 1885 .checked_add( 1886 u64::try_from(page_rows.len()).map_err(|_| Error::LegacyImportStagingFailed)?, 1887 ) 1888 .ok_or(Error::LegacyImportStagingFailed)?; 1889 let mut runtime_tx = self 1890 .pool 1891 .begin_with("BEGIN IMMEDIATE") 1892 .await 1893 .map_err(|_| Error::LegacyImportStagingFailed)?; 1894 let changed = sqlx::query("UPDATE radroots_runtime_legacy_import_members SET state = ?, resume_cursor = ?, staged_row_count = ?, updated_at_ms = ? WHERE import_id = ? AND source_kind = 'private' AND staged_row_count = ? AND resume_cursor IS ?") 1895 .bind(if complete { "ready" } else { "staging" }).bind(&next_cursor).bind(i64::try_from(total).map_err(|_| Error::LegacyImportStagingFailed)?).bind(updated_at).bind(classified.import_id().as_bytes().as_slice()).bind(i64::try_from(member.staged_row_count()).map_err(|_| Error::LegacyImportStagingFailed)?).bind(member.resume_cursor()).execute(&mut *runtime_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 1896 if changed.rows_affected() != 1 { 1897 return Err(Error::LegacyImportConflict); 1898 } 1899 let pending = sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_runtime_legacy_import_members WHERE import_id = ? AND state != 'ready'").bind(classified.import_id().as_bytes().as_slice()).fetch_one(&mut *runtime_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 1900 sqlx::query("UPDATE radroots_runtime_legacy_imports SET state = ?, updated_at_ms = ? WHERE import_id = ?") 1901 .bind(if pending == 0 { "ready" } else { "staging" }).bind(updated_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *runtime_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 1902 runtime_tx 1903 .commit() 1904 .await 1905 .map_err(|_| Error::LegacyImportStagingFailed)?; 1906 Ok(LegacyPrivateStagePage { 1907 table, 1908 staged_rows: u16::try_from(page_rows.len()) 1909 .map_err(|_| Error::LegacyImportStagingFailed)?, 1910 staged_row_count: total, 1911 resume_cursor: next_cursor, 1912 complete, 1913 }) 1914 } 1915 1916 /// Revalidates and describes a Studio predecessor snapshot for its host. 1917 #[cfg_attr(coverage_nightly, coverage(off))] 1918 pub async fn prepare_legacy_studio_handoff( 1919 &self, 1920 classified: &ClassifiedLegacyImport, 1921 ) -> Result<LegacyStudioHandoff, Error> { 1922 self.require_legacy_import_writer(classified.target_generation())?; 1923 verify_prepared_evidence(&classified.prepared).await?; 1924 let classification_sha256 = classification_digest(classified); 1925 let journal = self 1926 .legacy_import_journal(classified.import_id()) 1927 .await? 1928 .ok_or(Error::InvalidLegacyImportJournal)?; 1929 if !journal_matches_classified(&journal, classified, classification_sha256) { 1930 return Err(Error::LegacyImportConflict); 1931 } 1932 let classification = classified 1933 .sources() 1934 .iter() 1935 .find(|source| source.kind() == LegacySourceKind::Studio) 1936 .filter(|source| { 1937 source.schema() == LegacySchema::StudioV1HostHandoff 1938 && source.schema().disposition() == LegacyImportDisposition::HostHandoff 1939 }) 1940 .ok_or(Error::LegacyImportConflict)?; 1941 let member = journal 1942 .members() 1943 .iter() 1944 .find(|member| member.classification().kind() == LegacySourceKind::Studio) 1945 .ok_or(Error::LegacyImportConflict)?; 1946 if !matches!( 1947 member.state(), 1948 LegacyImportMemberState::Pending 1949 | LegacyImportMemberState::Staging 1950 | LegacyImportMemberState::Ready 1951 ) || member.staged_row_count() != 0 1952 { 1953 return Err(Error::LegacyImportConflict); 1954 } 1955 let snapshot = classified 1956 .prepared 1957 .snapshots() 1958 .iter() 1959 .find(|snapshot| snapshot.kind() == LegacySourceKind::Studio) 1960 .ok_or(Error::LegacyImportConflict)?; 1961 let evidence_path = classified.bundle_path().join(snapshot.relative_path()); 1962 Ok(LegacyStudioHandoff { 1963 import_id: classified.import_id(), 1964 evidence_path, 1965 byte_length: snapshot.byte_length(), 1966 source_sha256: snapshot.sha256(), 1967 catalog_sha256: classification.catalog_sha256(), 1968 handoff_sha256: studio_handoff_digest(classified, snapshot, classification), 1969 }) 1970 } 1971 1972 /// Records an exact host-owned Studio handoff acknowledgement without importing it. 1973 #[cfg_attr(coverage_nightly, coverage(off))] 1974 pub async fn acknowledge_legacy_studio_handoff( 1975 &self, 1976 classified: &ClassifiedLegacyImport, 1977 receipt: LegacyStudioHandoffReceipt, 1978 updated_at_unix_ms: u64, 1979 ) -> Result<LegacyImportJournal, Error> { 1980 if updated_at_unix_ms == 0 { 1981 return Err(Error::InvalidLegacyImportStageRequest); 1982 } 1983 let handoff = self.prepare_legacy_studio_handoff(classified).await?; 1984 if receipt.handoff_sha256() != handoff.handoff_sha256() { 1985 return Err(Error::LegacyImportConflict); 1986 } 1987 let receipt_cursor = studio_handoff_receipt_cursor(receipt); 1988 let journal = self 1989 .legacy_import_journal(classified.import_id()) 1990 .await? 1991 .ok_or(Error::InvalidLegacyImportJournal)?; 1992 let member = journal 1993 .members() 1994 .iter() 1995 .find(|member| member.classification().kind() == LegacySourceKind::Studio) 1996 .ok_or(Error::LegacyImportConflict)?; 1997 if member.state() == LegacyImportMemberState::Ready { 1998 return if member.resume_cursor() == Some(receipt_cursor.as_slice()) { 1999 Ok(journal) 2000 } else { 2001 Err(Error::LegacyImportConflict) 2002 }; 2003 } 2004 if !matches!( 2005 member.state(), 2006 LegacyImportMemberState::Pending | LegacyImportMemberState::Staging 2007 ) || member.resume_cursor().is_some() 2008 || member.staged_row_count() != 0 2009 || updated_at_unix_ms < member.updated_at_unix_ms() 2010 || updated_at_unix_ms < journal.updated_at_unix_ms() 2011 { 2012 return Err(Error::LegacyImportConflict); 2013 } 2014 let updated_at = i64::try_from(updated_at_unix_ms) 2015 .map_err(|_| Error::InvalidLegacyImportStageRequest)?; 2016 let mut transaction = self 2017 .pool 2018 .begin_with("BEGIN IMMEDIATE") 2019 .await 2020 .map_err(|_| Error::LegacyImportStagingFailed)?; 2021 if journal.state() == LegacyImportState::Classified { 2022 let changed = sqlx::query("UPDATE radroots_runtime_legacy_imports SET state = 'staging', updated_at_ms = ? WHERE import_id = ? AND state = 'classified'") 2023 .bind(updated_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *transaction).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2024 if changed.rows_affected() != 1 { 2025 return Err(Error::LegacyImportConflict); 2026 } 2027 } 2028 if member.state() == LegacyImportMemberState::Pending { 2029 let changed = sqlx::query("UPDATE radroots_runtime_legacy_import_members SET state = 'staging', updated_at_ms = ? WHERE import_id = ? AND source_kind = 'studio' AND state = 'pending' AND resume_cursor IS NULL AND staged_row_count = 0") 2030 .bind(updated_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *transaction).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2031 if changed.rows_affected() != 1 { 2032 return Err(Error::LegacyImportConflict); 2033 } 2034 } 2035 let changed = sqlx::query("UPDATE radroots_runtime_legacy_import_members SET state = 'ready', resume_cursor = ?, updated_at_ms = ? WHERE import_id = ? AND source_kind = 'studio' AND state = 'staging' AND resume_cursor IS NULL AND staged_row_count = 0") 2036 .bind(receipt_cursor.as_slice()).bind(updated_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *transaction).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2037 if changed.rows_affected() != 1 { 2038 return Err(Error::LegacyImportConflict); 2039 } 2040 let pending = sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_runtime_legacy_import_members WHERE import_id = ? AND state != 'ready'") 2041 .bind(classified.import_id().as_bytes().as_slice()).fetch_one(&mut *transaction).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2042 if pending == 0 { 2043 sqlx::query("UPDATE radroots_runtime_legacy_imports SET state = 'ready', updated_at_ms = ? WHERE import_id = ? AND state = 'staging'") 2044 .bind(updated_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *transaction).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2045 } 2046 transaction 2047 .commit() 2048 .await 2049 .map_err(|_| Error::LegacyImportStagingFailed)?; 2050 self.legacy_import_journal(classified.import_id()) 2051 .await? 2052 .ok_or(Error::InvalidLegacyImportJournal) 2053 } 2054 2055 /// Proves every classified source is completely staged or acknowledged. 2056 #[cfg_attr(coverage_nightly, coverage(off))] 2057 pub async fn validate_legacy_import( 2058 &self, 2059 classified: &ClassifiedLegacyImport, 2060 ) -> Result<LegacyImportValidation, Error> { 2061 self.require_legacy_import_writer(classified.target_generation())?; 2062 verify_prepared_evidence(&classified.prepared).await?; 2063 let classification_sha256 = classification_digest(classified); 2064 let journal = self 2065 .legacy_import_journal(classified.import_id()) 2066 .await? 2067 .ok_or(Error::InvalidLegacyImportJournal)?; 2068 if !journal_matches_classified(&journal, classified, classification_sha256) 2069 || journal.state() != LegacyImportState::Ready 2070 || journal.members().iter().any(|member| { 2071 member.state() != LegacyImportMemberState::Ready 2072 || (member.classification().kind() == LegacySourceKind::Studio 2073 && member.staged_row_count() != 0) 2074 }) 2075 { 2076 return Err(Error::LegacyImportConflict); 2077 } 2078 2079 let mut source_counts = Vec::with_capacity(classified.sources().len()); 2080 for source in classified.sources() { 2081 source_counts.push(( 2082 source.kind(), 2083 source_import_row_count(classified, source.kind()).await?, 2084 )); 2085 } 2086 2087 let mut runtime_tx = self 2088 .pool 2089 .begin_with("BEGIN IMMEDIATE") 2090 .await 2091 .map_err(|_| Error::LegacyImportStagingFailed)?; 2092 let mut private_tx = self 2093 .private_pool 2094 .begin_with("BEGIN IMMEDIATE") 2095 .await 2096 .map_err(|_| Error::LegacyImportStagingFailed)?; 2097 let current_state = sqlx::query_scalar::<_, String>( 2098 "SELECT state FROM radroots_runtime_legacy_imports WHERE import_id = ?", 2099 ) 2100 .bind(classified.import_id().as_bytes().as_slice()) 2101 .fetch_one(&mut *runtime_tx) 2102 .await 2103 .map_err(|_| Error::LegacyImportStagingFailed)?; 2104 let member_rows = sqlx::query( 2105 "SELECT source_kind, state, resume_cursor, staged_row_count 2106 FROM radroots_runtime_legacy_import_members 2107 WHERE import_id = ? ORDER BY source_kind", 2108 ) 2109 .bind(classified.import_id().as_bytes().as_slice()) 2110 .fetch_all(&mut *runtime_tx) 2111 .await 2112 .map_err(|_| Error::LegacyImportStagingFailed)?; 2113 if current_state != "ready" || member_rows.len() != classified.sources().len() { 2114 return Err(Error::LegacyImportConflict); 2115 } 2116 2117 let mut digest = Sha256::new(); 2118 for field in [ 2119 b"radroots.legacy.import.validation.v1".as_slice(), 2120 classified.import_id().as_bytes().as_slice(), 2121 classified.target_generation().as_bytes().as_slice(), 2122 classified.prepared.manifest_sha256().as_bytes().as_slice(), 2123 classification_sha256.as_bytes().as_slice(), 2124 ] { 2125 update_framed_digest(&mut digest, field)?; 2126 } 2127 let mut imported_row_count = 0_u64; 2128 for row in member_rows { 2129 let kind_value = row 2130 .try_get::<String, _>("source_kind") 2131 .map_err(|_| Error::InvalidLegacyImportJournal)?; 2132 let kind = parse_source_kind(kind_value.as_str())?; 2133 let state = row 2134 .try_get::<String, _>("state") 2135 .map_err(|_| Error::InvalidLegacyImportJournal)?; 2136 let cursor = row 2137 .try_get::<Option<Vec<u8>>, _>("resume_cursor") 2138 .map_err(|_| Error::InvalidLegacyImportJournal)?; 2139 let staged = u64::try_from( 2140 row.try_get::<i64, _>("staged_row_count") 2141 .map_err(|_| Error::InvalidLegacyImportJournal)?, 2142 ) 2143 .map_err(|_| Error::InvalidLegacyImportJournal)?; 2144 let source_count = source_counts 2145 .iter() 2146 .find_map(|(source_kind, count)| (*source_kind == kind).then_some(*count)) 2147 .ok_or(Error::LegacyImportConflict)?; 2148 if state != "ready" 2149 || staged != source_count 2150 || cursor.is_none() 2151 || (kind == LegacySourceKind::Studio && staged != 0) 2152 { 2153 return Err(Error::LegacyImportConflict); 2154 } 2155 update_framed_digest(&mut digest, kind_value.as_bytes())?; 2156 update_framed_digest(&mut digest, cursor.as_deref().unwrap_or_default())?; 2157 update_framed_digest(&mut digest, &staged.to_be_bytes())?; 2158 imported_row_count = imported_row_count 2159 .checked_add(staged) 2160 .ok_or(Error::LegacyImportStagingFailed)?; 2161 } 2162 hash_runtime_legacy_staging(&mut runtime_tx, classified.import_id(), &mut digest).await?; 2163 hash_private_legacy_staging(&mut private_tx, classified.import_id(), &mut digest).await?; 2164 runtime_tx 2165 .commit() 2166 .await 2167 .map_err(|_| Error::LegacyImportStagingFailed)?; 2168 private_tx 2169 .commit() 2170 .await 2171 .map_err(|_| Error::LegacyImportStagingFailed)?; 2172 Ok(LegacyImportValidation { 2173 imported_row_count, 2174 validation_sha256: MemberDigest::new(digest.finalize().into()), 2175 }) 2176 } 2177 2178 /// Seals validated legacy staging through a private-first recovery protocol. 2179 #[cfg_attr(coverage_nightly, coverage(off))] 2180 pub async fn finalize_legacy_import( 2181 &self, 2182 classified: &ClassifiedLegacyImport, 2183 expected: LegacyImportValidation, 2184 completed_at_unix_ms: u64, 2185 ) -> Result<LegacyImportCommitReceipt, Error> { 2186 self.require_legacy_import_writer(classified.target_generation())?; 2187 if completed_at_unix_ms == 0 { 2188 return Err(Error::InvalidLegacyImportStageRequest); 2189 } 2190 let classification_sha256 = classification_digest(classified); 2191 let journal = self 2192 .legacy_import_journal(classified.import_id()) 2193 .await? 2194 .ok_or(Error::InvalidLegacyImportJournal)?; 2195 if !journal_matches_classified(&journal, classified, classification_sha256) { 2196 return Err(Error::LegacyImportConflict); 2197 } 2198 if journal.state() == LegacyImportState::Complete { 2199 return self 2200 .completed_legacy_import_receipt(classified.import_id(), expected) 2201 .await; 2202 } 2203 if journal.state() != LegacyImportState::Ready 2204 || completed_at_unix_ms < journal.updated_at_unix_ms() 2205 { 2206 return Err(Error::LegacyImportConflict); 2207 } 2208 let actual = self.validate_legacy_import(classified).await?; 2209 if actual != expected { 2210 return Err(Error::LegacyImportConflict); 2211 } 2212 let completed_at = i64::try_from(completed_at_unix_ms) 2213 .map_err(|_| Error::InvalidLegacyImportStageRequest)?; 2214 let imported_row_count = i64::try_from(expected.imported_row_count()) 2215 .map_err(|_| Error::LegacyImportStagingFailed)?; 2216 2217 let mut private_tx = self 2218 .private_pool 2219 .begin_with("BEGIN IMMEDIATE") 2220 .await 2221 .map_err(|_| Error::LegacyImportStagingFailed)?; 2222 sqlx::query("INSERT OR IGNORE INTO radroots_private_legacy_import_commits(import_id, validation_sha256, imported_row_count, committed_at_ms) VALUES (?, ?, ?, ?)") 2223 .bind(classified.import_id().as_bytes().as_slice()).bind(expected.validation_sha256().as_bytes().as_slice()).bind(imported_row_count).bind(completed_at).execute(&mut *private_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2224 let private_record = sqlx::query("SELECT validation_sha256, imported_row_count, committed_at_ms FROM radroots_private_legacy_import_commits WHERE import_id = ?") 2225 .bind(classified.import_id().as_bytes().as_slice()).fetch_one(&mut *private_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2226 let private_committed_at = private_record 2227 .try_get::<i64, _>("committed_at_ms") 2228 .map_err(|_| Error::LegacyImportStagingFailed)?; 2229 if decode_digest( 2230 private_record 2231 .try_get("validation_sha256") 2232 .map_err(|_| Error::LegacyImportStagingFailed)?, 2233 )? != expected.validation_sha256() 2234 || private_record 2235 .try_get::<i64, _>("imported_row_count") 2236 .map_err(|_| Error::LegacyImportStagingFailed)? 2237 != imported_row_count 2238 || private_committed_at 2239 < i64::try_from(journal.updated_at_unix_ms()) 2240 .map_err(|_| Error::LegacyImportStagingFailed)? 2241 { 2242 return Err(Error::LegacyImportConflict); 2243 } 2244 private_tx 2245 .commit() 2246 .await 2247 .map_err(|_| Error::LegacyImportStagingFailed)?; 2248 let completed_at = private_committed_at; 2249 let completed_at_unix_ms = 2250 u64::try_from(completed_at).map_err(|_| Error::LegacyImportStagingFailed)?; 2251 2252 let mut runtime_tx = self 2253 .pool 2254 .begin_with("BEGIN IMMEDIATE") 2255 .await 2256 .map_err(|_| Error::LegacyImportStagingFailed)?; 2257 sqlx::query("INSERT INTO radroots_runtime_legacy_import_commits(import_id, validation_sha256, imported_row_count, completed_at_ms) VALUES (?, ?, ?, ?)") 2258 .bind(classified.import_id().as_bytes().as_slice()).bind(expected.validation_sha256().as_bytes().as_slice()).bind(imported_row_count).bind(completed_at).execute(&mut *runtime_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2259 let changed = sqlx::query("UPDATE radroots_runtime_legacy_imports SET state = 'committing', updated_at_ms = ? WHERE import_id = ? AND state = 'ready'") 2260 .bind(completed_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *runtime_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2261 if changed.rows_affected() != 1 { 2262 return Err(Error::LegacyImportConflict); 2263 } 2264 let changed = sqlx::query("UPDATE radroots_runtime_legacy_import_members SET state = 'complete', updated_at_ms = ? WHERE import_id = ? AND state = 'ready'") 2265 .bind(completed_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *runtime_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2266 if usize::try_from(changed.rows_affected()).map_err(|_| Error::LegacyImportStagingFailed)? 2267 != classified.sources().len() 2268 { 2269 return Err(Error::LegacyImportConflict); 2270 } 2271 let changed = sqlx::query("UPDATE radroots_runtime_legacy_imports SET state = 'complete', updated_at_ms = ?, completed_at_ms = ? WHERE import_id = ? AND state = 'committing'") 2272 .bind(completed_at).bind(completed_at).bind(classified.import_id().as_bytes().as_slice()).execute(&mut *runtime_tx).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2273 if changed.rows_affected() != 1 { 2274 return Err(Error::LegacyImportConflict); 2275 } 2276 runtime_tx 2277 .commit() 2278 .await 2279 .map_err(|_| Error::LegacyImportStagingFailed)?; 2280 Ok(LegacyImportCommitReceipt { 2281 validation_sha256: expected.validation_sha256(), 2282 imported_row_count: expected.imported_row_count(), 2283 completed_at_unix_ms, 2284 }) 2285 } 2286 2287 #[cfg_attr(coverage_nightly, coverage(off))] 2288 async fn completed_legacy_import_receipt( 2289 &self, 2290 import_id: LegacyImportId, 2291 expected: LegacyImportValidation, 2292 ) -> Result<LegacyImportCommitReceipt, Error> { 2293 let row = sqlx::query("SELECT validation_sha256, imported_row_count, completed_at_ms FROM radroots_runtime_legacy_import_commits WHERE import_id = ?") 2294 .bind(import_id.as_bytes().as_slice()).fetch_one(&self.pool).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2295 let validation_sha256 = decode_digest( 2296 row.try_get("validation_sha256") 2297 .map_err(|_| Error::LegacyImportStagingFailed)?, 2298 )?; 2299 let imported_row_count = u64::try_from( 2300 row.try_get::<i64, _>("imported_row_count") 2301 .map_err(|_| Error::LegacyImportStagingFailed)?, 2302 ) 2303 .map_err(|_| Error::LegacyImportStagingFailed)?; 2304 let completed_at_unix_ms = decode_positive_time( 2305 row.try_get("completed_at_ms") 2306 .map_err(|_| Error::LegacyImportStagingFailed)?, 2307 )?; 2308 if validation_sha256 != expected.validation_sha256() 2309 || imported_row_count != expected.imported_row_count() 2310 { 2311 return Err(Error::LegacyImportConflict); 2312 } 2313 let private_count = sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_private_legacy_import_commits WHERE import_id = ? AND validation_sha256 = ? AND imported_row_count = ? AND committed_at_ms = ?") 2314 .bind(import_id.as_bytes().as_slice()).bind(validation_sha256.as_bytes().as_slice()).bind(i64::try_from(imported_row_count).map_err(|_| Error::LegacyImportStagingFailed)?).bind(i64::try_from(completed_at_unix_ms).map_err(|_| Error::LegacyImportStagingFailed)?).fetch_one(&self.private_pool).await.map_err(|_| Error::LegacyImportStagingFailed)?; 2315 if private_count != 1 { 2316 return Err(Error::LegacyImportConflict); 2317 } 2318 Ok(LegacyImportCommitReceipt { 2319 validation_sha256, 2320 imported_row_count, 2321 completed_at_unix_ms, 2322 }) 2323 } 2324 2325 fn require_legacy_import_writer( 2326 &self, 2327 target_generation: SourceGeneration, 2328 ) -> Result<(), Error> { 2329 self.lifecycle 2330 .require_open() 2331 .map_err(|_| Error::BackupBackendUnavailable)?; 2332 if self.mode != EventStoreMode::ReadWrite { 2333 return Err(Error::RestoreRequiresWritableStorage); 2334 } 2335 if target_generation != self.generation { 2336 return Err(Error::LegacyImportTargetMismatch); 2337 } 2338 Ok(()) 2339 } 2340 } 2341 2342 struct LegacyBackupLayout { 2343 staging: PathBuf, 2344 finalized: PathBuf, 2345 } 2346 2347 impl LegacyBackupLayout { 2348 fn new(plan: &LegacyImportPlan) -> Self { 2349 let id = encode_id(plan.import_id().as_bytes()); 2350 Self { 2351 staging: plan 2352 .backup_root() 2353 .join(format!(".radroots-legacy-import-{id}.staging")), 2354 finalized: plan 2355 .backup_root() 2356 .join(format!("radroots-legacy-import-{id}")), 2357 } 2358 } 2359 2360 #[cfg_attr(coverage_nightly, coverage(off))] 2361 fn create(&self) -> Result<(), Error> { 2362 for path in [&self.staging, &self.finalized] { 2363 if path 2364 .try_exists() 2365 .map_err(|source| Error::LegacyImportFilesystem { 2366 operation: "inspect legacy import backup path", 2367 source, 2368 })? 2369 { 2370 return Err(Error::LegacyImportBackupAlreadyExists(path.clone())); 2371 } 2372 } 2373 #[cfg(unix)] 2374 let mut builder = fs::DirBuilder::new(); 2375 #[cfg(not(unix))] 2376 let builder = fs::DirBuilder::new(); 2377 #[cfg(unix)] 2378 { 2379 use std::os::unix::fs::DirBuilderExt; 2380 builder.mode(0o700); 2381 } 2382 builder 2383 .create(&self.staging) 2384 .map_err(|source| Error::LegacyImportFilesystem { 2385 operation: "create legacy import staging bundle", 2386 source, 2387 }) 2388 } 2389 } 2390 2391 #[cfg_attr(coverage_nightly, coverage(off))] 2392 async fn capture_legacy_source(source: &LegacySource, destination: &Path) -> Result<(), Error> { 2393 let destination_text = destination 2394 .to_str() 2395 .ok_or_else(|| Error::InvalidLegacySource(destination.to_path_buf()))?; 2396 let mut connection = SqliteConnection::connect_with( 2397 &SqliteConnectOptions::new() 2398 .filename(source.path()) 2399 .read_only(true) 2400 .foreign_keys(true), 2401 ) 2402 .await 2403 .map_err(|_| Error::LegacyImportBackupFailed { 2404 source_kind: source.kind().as_str(), 2405 })?; 2406 sqlx::query("VACUUM INTO ?") 2407 .bind(destination_text) 2408 .execute(&mut connection) 2409 .await 2410 .map_err(|_| Error::LegacyImportBackupFailed { 2411 source_kind: source.kind().as_str(), 2412 })?; 2413 connection 2414 .close() 2415 .await 2416 .map_err(|_| Error::LegacyImportBackupFailed { 2417 source_kind: source.kind().as_str(), 2418 })?; 2419 #[cfg(unix)] 2420 { 2421 use std::os::unix::fs::PermissionsExt; 2422 fs::set_permissions(destination, fs::Permissions::from_mode(0o600)).map_err(|source| { 2423 Error::LegacyImportFilesystem { 2424 operation: "secure legacy import backup member", 2425 source, 2426 } 2427 })?; 2428 } 2429 verify_legacy_snapshot(source.kind(), destination).await 2430 } 2431 2432 #[cfg_attr(coverage_nightly, coverage(off))] 2433 async fn verify_legacy_snapshot(kind: LegacySourceKind, path: &Path) -> Result<(), Error> { 2434 let mut connection = SqliteConnection::connect_with( 2435 &SqliteConnectOptions::new() 2436 .filename(path) 2437 .read_only(true) 2438 .foreign_keys(true), 2439 ) 2440 .await 2441 .map_err(|_| Error::LegacyImportSourceInvalid { 2442 source_kind: kind.as_str(), 2443 })?; 2444 let quick_check = sqlx::query_scalar::<_, String>("PRAGMA quick_check") 2445 .fetch_all(&mut connection) 2446 .await 2447 .map_err(|_| Error::LegacyImportSourceInvalid { 2448 source_kind: kind.as_str(), 2449 })?; 2450 let foreign_key_violation = sqlx::query("PRAGMA foreign_key_check") 2451 .fetch_optional(&mut connection) 2452 .await 2453 .map_err(|_| Error::LegacyImportSourceInvalid { 2454 source_kind: kind.as_str(), 2455 })? 2456 .is_some(); 2457 connection 2458 .close() 2459 .await 2460 .map_err(|_| Error::LegacyImportSourceInvalid { 2461 source_kind: kind.as_str(), 2462 })?; 2463 if quick_check == ["ok"] && !foreign_key_violation { 2464 Ok(()) 2465 } else { 2466 Err(Error::LegacyImportSourceInvalid { 2467 source_kind: kind.as_str(), 2468 }) 2469 } 2470 } 2471 2472 #[cfg_attr(coverage_nightly, coverage(off))] 2473 fn snapshot(kind: LegacySourceKind, path: &Path) -> Result<LegacySourceSnapshot, Error> { 2474 let (byte_length, sha256) = file_digest(path)?; 2475 Ok(LegacySourceSnapshot { 2476 kind, 2477 relative_path: kind.backup_file_name().to_owned(), 2478 byte_length, 2479 sha256, 2480 }) 2481 } 2482 2483 #[cfg_attr(coverage_nightly, coverage(off))] 2484 fn file_digest(path: &Path) -> Result<(u64, MemberDigest), Error> { 2485 let mut file = File::open(path).map_err(|source| Error::LegacyImportFilesystem { 2486 operation: "open legacy import evidence member", 2487 source, 2488 })?; 2489 file.sync_all() 2490 .map_err(|source| Error::LegacyImportFilesystem { 2491 operation: "sync legacy import evidence member", 2492 source, 2493 })?; 2494 let byte_length = file 2495 .metadata() 2496 .map_err(|source| Error::LegacyImportFilesystem { 2497 operation: "inspect legacy import evidence member", 2498 source, 2499 })? 2500 .len(); 2501 let mut digest = Sha256::new(); 2502 let mut buffer = [0_u8; 16 * 1_024]; 2503 loop { 2504 let read = file 2505 .read(&mut buffer) 2506 .map_err(|source| Error::LegacyImportFilesystem { 2507 operation: "hash legacy import evidence member", 2508 source, 2509 })?; 2510 if read == 0 { 2511 break; 2512 } 2513 digest.update(&buffer[..read]); 2514 } 2515 Ok((byte_length, MemberDigest::new(digest.finalize().into()))) 2516 } 2517 2518 #[cfg_attr(coverage_nightly, coverage(off))] 2519 async fn verify_prepared_evidence(prepared: &PreparedLegacyImport) -> Result<(), Error> { 2520 let bundle_metadata = fs::symlink_metadata(prepared.bundle_path()) 2521 .map_err(|_| Error::LegacyImportEvidenceInvalid)?; 2522 if !bundle_metadata.is_dir() || bundle_metadata.file_type().is_symlink() { 2523 return Err(Error::LegacyImportEvidenceInvalid); 2524 } 2525 let mut expected = BTreeSet::from([LEGACY_MANIFEST.to_owned()]); 2526 expected.extend( 2527 prepared 2528 .snapshots() 2529 .iter() 2530 .map(|snapshot| snapshot.relative_path().to_owned()), 2531 ); 2532 let mut actual = BTreeSet::new(); 2533 for entry in 2534 fs::read_dir(prepared.bundle_path()).map_err(|source| Error::LegacyImportFilesystem { 2535 operation: "read legacy import evidence bundle", 2536 source, 2537 })? 2538 { 2539 let entry = entry.map_err(|source| Error::LegacyImportFilesystem { 2540 operation: "read legacy import evidence entry", 2541 source, 2542 })?; 2543 let name = entry 2544 .file_name() 2545 .into_string() 2546 .map_err(|_| Error::LegacyImportEvidenceInvalid)?; 2547 let metadata = 2548 fs::symlink_metadata(entry.path()).map_err(|_| Error::LegacyImportEvidenceInvalid)?; 2549 if !metadata.is_file() || metadata.file_type().is_symlink() || !actual.insert(name) { 2550 return Err(Error::LegacyImportEvidenceInvalid); 2551 } 2552 } 2553 if actual != expected { 2554 return Err(Error::LegacyImportEvidenceInvalid); 2555 } 2556 let (manifest_length, manifest_digest) = 2557 file_digest(&prepared.bundle_path().join(LEGACY_MANIFEST))?; 2558 if manifest_length != prepared.manifest_byte_length() 2559 || manifest_digest != prepared.manifest_sha256() 2560 { 2561 return Err(Error::LegacyImportEvidenceInvalid); 2562 } 2563 for evidence in prepared.snapshots() { 2564 let path = prepared.bundle_path().join(evidence.relative_path()); 2565 if snapshot(evidence.kind(), &path)? != *evidence { 2566 return Err(Error::LegacyImportEvidenceInvalid); 2567 } 2568 verify_legacy_snapshot(evidence.kind(), &path).await?; 2569 } 2570 Ok(()) 2571 } 2572 2573 #[derive(Clone, Debug, Eq, PartialEq)] 2574 struct CatalogRow { 2575 object_type: String, 2576 name: String, 2577 table_name: String, 2578 sql: Option<String>, 2579 } 2580 2581 #[cfg_attr(coverage_nightly, coverage(off))] 2582 async fn classify_snapshot( 2583 kind: LegacySourceKind, 2584 path: &Path, 2585 ) -> Result<LegacySourceClassification, Error> { 2586 let mut connection = SqliteConnection::connect_with( 2587 &SqliteConnectOptions::new() 2588 .filename(path) 2589 .read_only(true) 2590 .foreign_keys(true), 2591 ) 2592 .await 2593 .map_err(|_| Error::LegacyImportSourceInvalid { 2594 source_kind: kind.as_str(), 2595 })?; 2596 let raw_user_version = sqlx::query_scalar::<_, i64>("PRAGMA user_version") 2597 .fetch_one(&mut connection) 2598 .await 2599 .map_err(|_| Error::LegacyImportSourceInvalid { 2600 source_kind: kind.as_str(), 2601 })?; 2602 let catalog = read_catalog(&mut connection, kind).await?; 2603 let (schema, governed_catalog) = match kind { 2604 LegacySourceKind::EventStore => { 2605 classify_event_store(&mut connection, raw_user_version, &catalog).await? 2606 } 2607 LegacySourceKind::Outbox => ( 2608 classify_fixed_catalog(kind, raw_user_version, &catalog, 0, OUTBOX_CATALOG_SHA256)?, 2609 catalog, 2610 ), 2611 LegacySourceKind::Private => ( 2612 classify_fixed_catalog(kind, raw_user_version, &catalog, 1, PRIVATE_CATALOG_SHA256)?, 2613 catalog, 2614 ), 2615 LegacySourceKind::Studio => ( 2616 classify_fixed_catalog(kind, raw_user_version, &catalog, 0, STUDIO_CATALOG_SHA256)?, 2617 catalog, 2618 ), 2619 }; 2620 let catalog_sha256 = catalog_fingerprint(&governed_catalog); 2621 connection 2622 .close() 2623 .await 2624 .map_err(|_| Error::LegacyImportSourceInvalid { 2625 source_kind: kind.as_str(), 2626 })?; 2627 let user_version = u32::try_from(raw_user_version) 2628 .map_err(|_| unsupported_schema(kind, raw_user_version, catalog_sha256))?; 2629 Ok(LegacySourceClassification { 2630 kind, 2631 schema, 2632 user_version, 2633 catalog_sha256, 2634 }) 2635 } 2636 2637 #[cfg_attr(coverage_nightly, coverage(off))] 2638 async fn read_catalog( 2639 connection: &mut SqliteConnection, 2640 kind: LegacySourceKind, 2641 ) -> Result<Vec<CatalogRow>, Error> { 2642 sqlx::query("SELECT type, name, tbl_name, sql FROM main.sqlite_schema") 2643 .fetch_all(connection) 2644 .await 2645 .map_err(|_| Error::LegacyImportSourceInvalid { 2646 source_kind: kind.as_str(), 2647 })? 2648 .into_iter() 2649 .map(|row| { 2650 let name = 2651 row.try_get::<String, _>("name") 2652 .map_err(|_| Error::LegacyImportSourceInvalid { 2653 source_kind: kind.as_str(), 2654 })?; 2655 Ok(CatalogRow { 2656 object_type: row 2657 .try_get("type") 2658 .map_err(|_| Error::LegacyImportSourceInvalid { 2659 source_kind: kind.as_str(), 2660 })?, 2661 table_name: row.try_get("tbl_name").map_err(|_| { 2662 Error::LegacyImportSourceInvalid { 2663 source_kind: kind.as_str(), 2664 } 2665 })?, 2666 sql: row 2667 .try_get("sql") 2668 .map_err(|_| Error::LegacyImportSourceInvalid { 2669 source_kind: kind.as_str(), 2670 })?, 2671 name, 2672 }) 2673 }) 2674 .collect::<Result<Vec<_>, Error>>() 2675 .map(|catalog| { 2676 catalog 2677 .into_iter() 2678 .filter(|row| !row.name.to_ascii_lowercase().starts_with("sqlite_")) 2679 .collect() 2680 }) 2681 } 2682 2683 fn classify_fixed_catalog( 2684 kind: LegacySourceKind, 2685 user_version: i64, 2686 catalog: &[CatalogRow], 2687 expected_user_version: i64, 2688 expected_catalog_sha256: &str, 2689 ) -> Result<LegacySchema, Error> { 2690 let fingerprint = catalog_fingerprint(catalog); 2691 if user_version != expected_user_version 2692 || encode_digest(fingerprint.as_bytes()) != expected_catalog_sha256 2693 { 2694 return Err(unsupported_schema(kind, user_version, fingerprint)); 2695 } 2696 Ok(match kind { 2697 LegacySourceKind::Outbox => LegacySchema::OutboxV1, 2698 LegacySourceKind::Private => LegacySchema::PrivateV1, 2699 LegacySourceKind::Studio => LegacySchema::StudioV1HostHandoff, 2700 LegacySourceKind::EventStore => return Err(Error::LegacyImportMigrationHistoryInvalid), 2701 }) 2702 } 2703 2704 #[cfg_attr(coverage_nightly, coverage(off))] 2705 async fn classify_event_store( 2706 connection: &mut SqliteConnection, 2707 user_version: i64, 2708 catalog: &[CatalogRow], 2709 ) -> Result<(LegacySchema, Vec<CatalogRow>), Error> { 2710 if user_version != 0 { 2711 return Err(unsupported_schema( 2712 LegacySourceKind::EventStore, 2713 user_version, 2714 catalog_fingerprint(catalog), 2715 )); 2716 } 2717 let ledger_rows = catalog 2718 .iter() 2719 .filter(|row| { 2720 row.name.eq_ignore_ascii_case(EVENT_STORE_LEDGER) 2721 || row.table_name.eq_ignore_ascii_case(EVENT_STORE_LEDGER) 2722 }) 2723 .collect::<Vec<_>>(); 2724 let governed = catalog 2725 .iter() 2726 .filter(|row| !row.name.eq_ignore_ascii_case(EVENT_STORE_LEDGER)) 2727 .cloned() 2728 .collect::<Vec<_>>(); 2729 let fingerprint = catalog_fingerprint(&governed); 2730 let version = if ledger_rows.is_empty() { 2731 if encode_digest(fingerprint.as_bytes()) != EVENT_STORE_MIGRATIONS[0].schema_sha256 { 2732 return Err(unsupported_schema( 2733 LegacySourceKind::EventStore, 2734 user_version, 2735 fingerprint, 2736 )); 2737 } 2738 1 2739 } else { 2740 if ledger_rows.len() != 1 { 2741 return Err(Error::LegacyImportMigrationHistoryInvalid); 2742 } 2743 let ledger = ledger_rows[0]; 2744 if ledger.object_type != "table" 2745 || ledger.name != EVENT_STORE_LEDGER 2746 || ledger.table_name != EVENT_STORE_LEDGER 2747 || ledger.sql.as_deref() != Some(EVENT_STORE_LEDGER_DDL) 2748 { 2749 return Err(Error::LegacyImportMigrationHistoryInvalid); 2750 } 2751 validate_event_history(connection).await? 2752 }; 2753 let expected = EVENT_STORE_MIGRATIONS 2754 .get(usize::try_from(version - 1).map_err(|_| Error::LegacyImportMigrationHistoryInvalid)?) 2755 .ok_or(Error::LegacyImportMigrationHistoryInvalid)?; 2756 if encode_digest(fingerprint.as_bytes()) != expected.schema_sha256 { 2757 return Err(unsupported_schema( 2758 LegacySourceKind::EventStore, 2759 user_version, 2760 fingerprint, 2761 )); 2762 } 2763 let schema = match version { 2764 1 => LegacySchema::EventStoreV1, 2765 2 => LegacySchema::EventStoreV2, 2766 3 => LegacySchema::EventStoreV3, 2767 4 => LegacySchema::EventStoreV4, 2768 _ => return Err(Error::LegacyImportMigrationHistoryInvalid), 2769 }; 2770 Ok((schema, governed)) 2771 } 2772 2773 #[cfg_attr(coverage_nightly, coverage(off))] 2774 async fn validate_event_history(connection: &mut SqliteConnection) -> Result<u32, Error> { 2775 let rows = sqlx::query( 2776 "SELECT version, name, up_sha256, down_sha256, schema_sha256 FROM main.radroots_event_store_schema_migrations ORDER BY version", 2777 ) 2778 .fetch_all(connection) 2779 .await 2780 .map_err(|_| Error::LegacyImportMigrationHistoryInvalid)?; 2781 if rows.is_empty() || rows.len() > EVENT_STORE_MIGRATIONS.len() { 2782 return Err(Error::LegacyImportMigrationHistoryInvalid); 2783 } 2784 for (index, row) in rows.iter().enumerate() { 2785 let expected = &EVENT_STORE_MIGRATIONS[index]; 2786 if row.try_get::<i64, _>("version").ok() != Some(i64::from(expected.version)) 2787 || row.try_get::<String, _>("name").ok().as_deref() != Some(expected.name) 2788 || row.try_get::<String, _>("up_sha256").ok().as_deref() != Some(expected.up_sha256) 2789 || row.try_get::<String, _>("down_sha256").ok().as_deref() != Some(expected.down_sha256) 2790 || row.try_get::<String, _>("schema_sha256").ok().as_deref() 2791 != Some(expected.schema_sha256) 2792 { 2793 return Err(Error::LegacyImportMigrationHistoryInvalid); 2794 } 2795 } 2796 u32::try_from(rows.len()).map_err(|_| Error::LegacyImportMigrationHistoryInvalid) 2797 } 2798 2799 fn catalog_fingerprint(catalog: &[CatalogRow]) -> MemberDigest { 2800 let mut rows = catalog.to_vec(); 2801 rows.sort_by(|left, right| { 2802 ( 2803 left.object_type.as_bytes(), 2804 left.name.as_bytes(), 2805 left.table_name.as_bytes(), 2806 left.sql.as_deref().unwrap_or("").as_bytes(), 2807 ) 2808 .cmp(&( 2809 right.object_type.as_bytes(), 2810 right.name.as_bytes(), 2811 right.table_name.as_bytes(), 2812 right.sql.as_deref().unwrap_or("").as_bytes(), 2813 )) 2814 }); 2815 let mut digest = Sha256::new(); 2816 for row in rows { 2817 for field in [ 2818 row.object_type.as_str(), 2819 row.name.as_str(), 2820 row.table_name.as_str(), 2821 row.sql.as_deref().unwrap_or(""), 2822 ] { 2823 digest.update(field.as_bytes()); 2824 digest.update([0]); 2825 } 2826 } 2827 MemberDigest::new(digest.finalize().into()) 2828 } 2829 2830 fn unsupported_schema( 2831 kind: LegacySourceKind, 2832 user_version: i64, 2833 catalog_sha256: MemberDigest, 2834 ) -> Error { 2835 Error::UnsupportedLegacySchema { 2836 source_kind: kind.as_str(), 2837 user_version, 2838 catalog_sha256: encode_digest(catalog_sha256.as_bytes()), 2839 } 2840 } 2841 2842 fn classification_digest(classified: &ClassifiedLegacyImport) -> MemberDigest { 2843 let mut digest = Sha256::new(); 2844 for field in [ 2845 classified.import_id().as_bytes().as_slice(), 2846 classified.target_generation().as_bytes().as_slice(), 2847 classified.prepared.manifest_sha256().as_bytes().as_slice(), 2848 ] { 2849 digest.update(field); 2850 digest.update([0]); 2851 } 2852 for source in classified.sources() { 2853 for field in [ 2854 source.kind().as_str().as_bytes(), 2855 source.schema().as_str().as_bytes(), 2856 source.schema().disposition().as_str().as_bytes(), 2857 source.catalog_sha256().as_bytes().as_slice(), 2858 ] { 2859 digest.update(field); 2860 digest.update([0]); 2861 } 2862 digest.update(source.user_version().to_be_bytes()); 2863 digest.update([0]); 2864 } 2865 MemberDigest::new(digest.finalize().into()) 2866 } 2867 2868 fn studio_handoff_digest( 2869 classified: &ClassifiedLegacyImport, 2870 snapshot: &LegacySourceSnapshot, 2871 classification: &LegacySourceClassification, 2872 ) -> MemberDigest { 2873 let mut digest = Sha256::new(); 2874 for field in [ 2875 b"radroots.legacy.studio.handoff.v1".as_slice(), 2876 classified.import_id().as_bytes().as_slice(), 2877 classified.target_generation().as_bytes().as_slice(), 2878 classified.prepared.manifest_sha256().as_bytes().as_slice(), 2879 snapshot.relative_path().as_bytes(), 2880 snapshot.sha256().as_bytes().as_slice(), 2881 classification.catalog_sha256().as_bytes().as_slice(), 2882 ] { 2883 digest.update(field); 2884 digest.update([0]); 2885 } 2886 digest.update(snapshot.byte_length().to_be_bytes()); 2887 MemberDigest::new(digest.finalize().into()) 2888 } 2889 2890 fn studio_handoff_receipt_cursor(receipt: LegacyStudioHandoffReceipt) -> [u8; 64] { 2891 let mut cursor = [0_u8; 64]; 2892 cursor[..32].copy_from_slice(receipt.handoff_sha256().as_bytes()); 2893 cursor[32..].copy_from_slice(receipt.host_commitment_sha256().as_bytes()); 2894 cursor 2895 } 2896 2897 fn update_framed_digest(digest: &mut Sha256, value: &[u8]) -> Result<(), Error> { 2898 let length = u64::try_from(value.len()).map_err(|_| Error::LegacyImportStagingFailed)?; 2899 digest.update(length.to_be_bytes()); 2900 digest.update(value); 2901 Ok(()) 2902 } 2903 2904 #[cfg_attr(coverage_nightly, coverage(off))] 2905 async fn source_import_row_count( 2906 classified: &ClassifiedLegacyImport, 2907 kind: LegacySourceKind, 2908 ) -> Result<u64, Error> { 2909 if kind == LegacySourceKind::Studio { 2910 return Ok(0); 2911 } 2912 let snapshot = classified 2913 .prepared 2914 .snapshots() 2915 .iter() 2916 .find(|snapshot| snapshot.kind() == kind) 2917 .ok_or(Error::LegacyImportConflict)?; 2918 let mut connection = SqliteConnection::connect_with( 2919 &SqliteConnectOptions::new() 2920 .filename(classified.bundle_path().join(snapshot.relative_path())) 2921 .read_only(true), 2922 ) 2923 .await 2924 .map_err(|_| Error::LegacyImportStagingFailed)?; 2925 let query = match kind { 2926 LegacySourceKind::EventStore => "SELECT COUNT(*) FROM event_envelopes", 2927 LegacySourceKind::Outbox => { 2928 "SELECT (SELECT COUNT(*) FROM outbox_operations) 2929 + (SELECT COUNT(*) FROM outbox_event) 2930 + (SELECT COUNT(*) FROM outbox_delivery_plan) 2931 + (SELECT COUNT(*) FROM outbox_delivery_target) 2932 + (SELECT COUNT(*) FROM outbox_delivery_attempt)" 2933 } 2934 LegacySourceKind::Private => { 2935 "SELECT (SELECT COUNT(*) FROM private_metadata) 2936 + (SELECT COUNT(*) FROM wrapped_profile_key) 2937 + (SELECT COUNT(*) FROM wrapped_signing_secret) 2938 + (SELECT COUNT(*) FROM private_farm_location) 2939 + (SELECT COUNT(*) FROM private_trade_artifacts) 2940 + (SELECT COUNT(*) FROM cursor_hmac_key) 2941 + (SELECT COUNT(*) FROM nip46_session_private) 2942 + (SELECT COUNT(*) FROM key_rotation_progress)" 2943 } 2944 LegacySourceKind::Studio => unreachable!("Studio is host-owned"), 2945 }; 2946 let count = sqlx::query_scalar::<_, i64>(query) 2947 .fetch_one(&mut connection) 2948 .await 2949 .map_err(|_| Error::LegacyImportStagingFailed)?; 2950 connection 2951 .close() 2952 .await 2953 .map_err(|_| Error::LegacyImportStagingFailed)?; 2954 u64::try_from(count).map_err(|_| Error::LegacyImportStagingFailed) 2955 } 2956 2957 #[cfg_attr(coverage_nightly, coverage(off))] 2958 async fn hash_runtime_legacy_staging( 2959 transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, 2960 import_id: LegacyImportId, 2961 digest: &mut Sha256, 2962 ) -> Result<(), Error> { 2963 update_framed_digest(digest, b"runtime_events")?; 2964 let event_rows = sqlx::query_scalar::<_, String>( 2965 "SELECT json_array(legacy_sequence, hex(event_id), hex(signed_event), 2966 legacy_verification_status, legacy_contract_status, 2967 legacy_projection_eligible, legacy_inserted_at_ms, 2968 legacy_updated_at_ms) 2969 FROM radroots_runtime_legacy_event_staging 2970 WHERE import_id = ? ORDER BY legacy_sequence", 2971 ) 2972 .bind(import_id.as_bytes().as_slice()) 2973 .fetch_all(&mut **transaction) 2974 .await 2975 .map_err(|_| Error::LegacyImportStagingFailed)?; 2976 for row in event_rows { 2977 update_framed_digest(digest, row.as_bytes())?; 2978 } 2979 2980 update_framed_digest(digest, b"runtime_outbox")?; 2981 let outbox_rows = sqlx::query_scalar::<_, String>( 2982 "SELECT json_array(table_kind, legacy_id, parent_legacy_id, 2983 related_legacy_id, hex(record_json)) 2984 FROM radroots_runtime_legacy_outbox_staging 2985 WHERE import_id = ? ORDER BY table_kind, legacy_id", 2986 ) 2987 .bind(import_id.as_bytes().as_slice()) 2988 .fetch_all(&mut **transaction) 2989 .await 2990 .map_err(|_| Error::LegacyImportStagingFailed)?; 2991 for row in outbox_rows { 2992 update_framed_digest(digest, row.as_bytes())?; 2993 } 2994 Ok(()) 2995 } 2996 2997 #[cfg_attr(coverage_nightly, coverage(off))] 2998 async fn hash_private_legacy_staging( 2999 transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, 3000 import_id: LegacyImportId, 3001 digest: &mut Sha256, 3002 ) -> Result<(), Error> { 3003 update_framed_digest(digest, b"private_records")?; 3004 let rows = sqlx::query_scalar::<_, String>( 3005 "SELECT json_array(table_kind, key_cursor, parent_key_version, 3006 hex(record_json)) 3007 FROM radroots_private_legacy_import_staging 3008 WHERE import_id = ? ORDER BY table_kind, key_cursor", 3009 ) 3010 .bind(import_id.as_bytes().as_slice()) 3011 .fetch_all(&mut **transaction) 3012 .await 3013 .map_err(|_| Error::LegacyImportStagingFailed)?; 3014 for row in rows { 3015 update_framed_digest(digest, row.as_bytes())?; 3016 } 3017 Ok(()) 3018 } 3019 3020 fn outbox_stage_query(table: LegacyOutboxTable) -> &'static str { 3021 match table { 3022 LegacyOutboxTable::Operations => { 3023 "SELECT operation_id AS legacy_id, NULL AS parent_legacy_id, 3024 NULL AS related_legacy_id, 3025 CAST(json_array(operation_kind, expected_pubkey, semantic_scope, 3026 trade_id, mutation_id, canonical_payload_sha256, idempotency_key, 3027 operation_idempotency_digest, status, created_at_ms, updated_at_ms) 3028 AS BLOB) AS record_json 3029 FROM outbox_operations WHERE operation_id > ? 3030 ORDER BY operation_id LIMIT ?" 3031 } 3032 LegacyOutboxTable::Events => { 3033 "SELECT outbox_event_id AS legacy_id, operation_id AS parent_legacy_id, 3034 NULL AS related_legacy_id, 3035 CAST(json_array(event_id, expected_pubkey, draft_json, 3036 signed_event_json, raw_event_json, state, attempt_count, 3037 claim_token, claim_owner, claim_expires_at_ms, 3038 active_delivery_plan_id, next_attempt_after_ms, last_error, 3039 event_store_ingested, event_store_inserted, 3040 event_store_ingested_at_ms, created_at_ms, updated_at_ms) 3041 AS BLOB) AS record_json 3042 FROM outbox_event WHERE outbox_event_id > ? 3043 ORDER BY outbox_event_id LIMIT ?" 3044 } 3045 LegacyOutboxTable::DeliveryPlans => { 3046 "SELECT delivery_plan_id AS legacy_id, outbox_event_id AS parent_legacy_id, 3047 NULL AS related_legacy_id, 3048 CAST(json_array(transport_profile_id, target_policy_fingerprint, 3049 target_policy_version, satisfaction_policy, required_success_count, 3050 delivery_plan_idempotency_digest, status, satisfied_at_ms, 3051 created_at_ms, updated_at_ms) AS BLOB) AS record_json 3052 FROM outbox_delivery_plan WHERE delivery_plan_id > ? 3053 ORDER BY delivery_plan_id LIMIT ?" 3054 } 3055 LegacyOutboxTable::DeliveryTargets => { 3056 "SELECT delivery_target_id AS legacy_id, delivery_plan_id AS parent_legacy_id, 3057 NULL AS related_legacy_id, 3058 CAST(json_array(transport_kind, endpoint_uri, target_scope, 3059 target_label, endpoint_fingerprint, status, last_outcome_kind, 3060 attempt_count, last_attempt_at_ms, completed_at_ms, last_error) 3061 AS BLOB) AS record_json 3062 FROM outbox_delivery_target WHERE delivery_target_id > ? 3063 ORDER BY delivery_target_id LIMIT ?" 3064 } 3065 LegacyOutboxTable::DeliveryAttempts => { 3066 "SELECT delivery_attempt_id AS legacy_id, 3067 delivery_target_id AS parent_legacy_id, 3068 delivery_plan_id AS related_legacy_id, 3069 CAST(json_array(status, outcome_kind, attempted_at_ms, message) 3070 AS BLOB) AS record_json 3071 FROM outbox_delivery_attempt WHERE delivery_attempt_id > ? 3072 ORDER BY delivery_attempt_id LIMIT ?" 3073 } 3074 } 3075 } 3076 3077 fn encode_outbox_stage_cursor(table: LegacyOutboxTable, legacy_id: i64) -> [u8; 9] { 3078 let mut cursor = [0_u8; 9]; 3079 cursor[0] = table.code(); 3080 cursor[1..].copy_from_slice(&legacy_id.to_be_bytes()); 3081 cursor 3082 } 3083 3084 fn decode_outbox_stage_cursor(cursor: Option<&[u8]>) -> Result<(LegacyOutboxTable, i64), Error> { 3085 let Some(cursor) = cursor else { 3086 return Ok((LegacyOutboxTable::Operations, 0)); 3087 }; 3088 let exact = decode_exact_outbox_stage_cursor(cursor)?; 3089 let table = match exact[0] { 3090 1 => LegacyOutboxTable::Operations, 3091 2 => LegacyOutboxTable::Events, 3092 3 => LegacyOutboxTable::DeliveryPlans, 3093 4 => LegacyOutboxTable::DeliveryTargets, 3094 5 => LegacyOutboxTable::DeliveryAttempts, 3095 _ => return Err(Error::InvalidLegacyImportJournal), 3096 }; 3097 let legacy_id = i64::from_be_bytes( 3098 exact[1..] 3099 .try_into() 3100 .map_err(|_| Error::InvalidLegacyImportJournal)?, 3101 ); 3102 if legacy_id < 0 { 3103 return Err(Error::InvalidLegacyImportJournal); 3104 } 3105 Ok((table, legacy_id)) 3106 } 3107 3108 fn decode_exact_outbox_stage_cursor(cursor: &[u8]) -> Result<[u8; 9], Error> { 3109 <[u8; 9]>::try_from(cursor).map_err(|_| Error::InvalidLegacyImportJournal) 3110 } 3111 3112 fn private_stage_query(table: LegacyPrivateTable) -> &'static str { 3113 match table { 3114 LegacyPrivateTable::Metadata => { 3115 "SELECT printf('%020d', singleton) AS key_cursor, NULL AS parent_key_version, CAST(json_array(singleton, schema_version, hex(profile_id), hex(runtime_contract_hash), key_version, sqlite_source_id, created_at_ms, updated_at_ms) AS BLOB) AS record_json FROM private_metadata WHERE printf('%020d', singleton) > ? ORDER BY key_cursor LIMIT ?" 3116 } 3117 LegacyPrivateTable::WrappedProfileKeys => { 3118 "SELECT printf('%020d', key_version) AS key_cursor, NULL AS parent_key_version, CAST(json_array(key_version, credential_backend, hex(wrapped_key), hex(nonce), created_at_ms, retired_at_ms) AS BLOB) AS record_json FROM wrapped_profile_key WHERE printf('%020d', key_version) > ? ORDER BY key_cursor LIMIT ?" 3119 } 3120 LegacyPrivateTable::SigningSecrets => { 3121 "SELECT hex(account_id) AS key_cursor, key_version AS parent_key_version, CAST(json_array(hex(account_id), hex(public_key), key_version, hex(ciphertext), hex(nonce), created_at_ms, updated_at_ms) AS BLOB) AS record_json FROM wrapped_signing_secret WHERE hex(account_id) > ? ORDER BY key_cursor LIMIT ?" 3122 } 3123 LegacyPrivateTable::FarmLocations => { 3124 "SELECT printf('%010d|%s|%s', farm_kind, hex(owner_pubkey), farm_d_tag) AS key_cursor, key_version AS parent_key_version, CAST(json_array(farm_kind, hex(owner_pubkey), farm_d_tag, key_version, hex(ciphertext), hex(nonce), created_at_ms, updated_at_ms) AS BLOB) AS record_json FROM private_farm_location WHERE printf('%010d|%s|%s', farm_kind, hex(owner_pubkey), farm_d_tag) > ? ORDER BY key_cursor LIMIT ?" 3125 } 3126 LegacyPrivateTable::TradeArtifacts => { 3127 "SELECT artifact_id AS key_cursor, key_version AS parent_key_version, CAST(json_array(artifact_id, trade_id, candidate_id, artifact_kind, schema_id, ciphertext_commitment, key_version, hex(ciphertext), hex(encryption_metadata), retention_class, created_at_ms, expires_at_ms, deleted_at_ms) AS BLOB) AS record_json FROM private_trade_artifacts WHERE artifact_id > ? ORDER BY key_cursor LIMIT ?" 3128 } 3129 LegacyPrivateTable::CursorKeys => { 3130 "SELECT hex(key_id) AS key_cursor, key_version AS parent_key_version, CAST(json_array(hex(key_id), key_version, hex(ciphertext), hex(nonce), created_at_ms, retired_at_ms) AS BLOB) AS record_json FROM cursor_hmac_key WHERE hex(key_id) > ? ORDER BY key_cursor LIMIT ?" 3131 } 3132 LegacyPrivateTable::Nip46Sessions => { 3133 "SELECT hex(session_id) AS key_cursor, key_version AS parent_key_version, CAST(json_array(hex(session_id), hex(user_pubkey), hex(remote_signer_pubkey), hex(client_pubkey), key_version, hex(ciphertext), hex(nonce), expires_at_ms, status, created_at_ms, updated_at_ms) AS BLOB) AS record_json FROM nip46_session_private WHERE hex(session_id) > ? ORDER BY key_cursor LIMIT ?" 3134 } 3135 LegacyPrivateTable::RotationProgress => { 3136 "SELECT printf('%020d', singleton) AS key_cursor, NULL AS parent_key_version, CAST(json_array(singleton, from_key_version, to_key_version, table_name, hex(last_primary_key), state, started_at_ms, updated_at_ms, error_code) AS BLOB) AS record_json FROM key_rotation_progress WHERE printf('%020d', singleton) > ? ORDER BY key_cursor LIMIT ?" 3137 } 3138 } 3139 } 3140 3141 fn encode_private_stage_cursor(table: LegacyPrivateTable, key: &str) -> Vec<u8> { 3142 let mut cursor = Vec::with_capacity(key.len() + 1); 3143 cursor.push(table.code()); 3144 cursor.extend_from_slice(key.as_bytes()); 3145 cursor 3146 } 3147 3148 fn decode_private_stage_cursor( 3149 cursor: Option<&[u8]>, 3150 ) -> Result<(LegacyPrivateTable, String), Error> { 3151 let Some(cursor) = cursor else { 3152 return Ok((LegacyPrivateTable::Metadata, String::new())); 3153 }; 3154 if cursor.is_empty() || cursor.len() > 1025 { 3155 return Err(Error::InvalidLegacyImportJournal); 3156 } 3157 let table = match cursor[0] { 3158 1 => LegacyPrivateTable::Metadata, 3159 2 => LegacyPrivateTable::WrappedProfileKeys, 3160 3 => LegacyPrivateTable::SigningSecrets, 3161 4 => LegacyPrivateTable::FarmLocations, 3162 5 => LegacyPrivateTable::TradeArtifacts, 3163 6 => LegacyPrivateTable::CursorKeys, 3164 7 => LegacyPrivateTable::Nip46Sessions, 3165 8 => LegacyPrivateTable::RotationProgress, 3166 _ => return Err(Error::InvalidLegacyImportJournal), 3167 }; 3168 let key = std::str::from_utf8(&cursor[1..]).map_err(|_| Error::InvalidLegacyImportJournal)?; 3169 Ok((table, key.to_owned())) 3170 } 3171 3172 struct ConvertedLegacyEvent { 3173 sequence: i64, 3174 event_id: [u8; 32], 3175 signed_event: Vec<u8>, 3176 verification_status: String, 3177 contract_status: String, 3178 projection_eligible: i64, 3179 inserted_at_ms: i64, 3180 updated_at_ms: i64, 3181 } 3182 3183 fn convert_legacy_event_row(row: &sqlx::sqlite::SqliteRow) -> Result<ConvertedLegacyEvent, Error> { 3184 let sequence = row 3185 .try_get::<i64, _>("seq") 3186 .map_err(|_| Error::LegacyImportStagingFailed)?; 3187 let invalid = || Error::LegacyImportRowInvalid { 3188 source_kind: LegacySourceKind::EventStore.as_str(), 3189 legacy_sequence: sequence, 3190 }; 3191 if sequence <= 0 { 3192 return Err(invalid()); 3193 } 3194 let event_id = row 3195 .try_get::<String, _>("event_id") 3196 .map_err(|_| invalid())?; 3197 let raw_json = row 3198 .try_get::<String, _>("raw_json") 3199 .map_err(|_| invalid())?; 3200 let signed = Codec::decode_signed_event(raw_json.as_str()).map_err(|_| invalid())?; 3201 if signed.id().to_hex() != event_id { 3202 return Err(invalid()); 3203 } 3204 let verification_status = row 3205 .try_get::<String, _>("verification_status") 3206 .map_err(|_| invalid())?; 3207 let contract_status = row 3208 .try_get::<String, _>("contract_status") 3209 .map_err(|_| invalid())?; 3210 let projection_eligible = row 3211 .try_get::<i64, _>("projection_eligible") 3212 .map_err(|_| invalid())?; 3213 let inserted_at_ms = row 3214 .try_get::<i64, _>("inserted_at_ms") 3215 .map_err(|_| invalid())?; 3216 let updated_at_ms = row 3217 .try_get::<i64, _>("updated_at_ms") 3218 .map_err(|_| invalid())?; 3219 if verification_status.is_empty() 3220 || verification_status.len() > 64 3221 || contract_status.is_empty() 3222 || contract_status.len() > 64 3223 || !matches!(projection_eligible, 0 | 1) 3224 || inserted_at_ms <= 0 3225 || updated_at_ms < inserted_at_ms 3226 { 3227 return Err(invalid()); 3228 } 3229 Ok(ConvertedLegacyEvent { 3230 sequence, 3231 event_id: *signed.id().as_bytes(), 3232 signed_event: raw_json.into_bytes(), 3233 verification_status, 3234 contract_status, 3235 projection_eligible, 3236 inserted_at_ms, 3237 updated_at_ms, 3238 }) 3239 } 3240 3241 fn encode_event_stage_cursor(sequence: i64) -> [u8; 8] { 3242 sequence.to_be_bytes() 3243 } 3244 3245 fn decode_event_stage_cursor(cursor: Option<&[u8]>) -> Result<i64, Error> { 3246 cursor.map_or(Ok(0), |bytes| { 3247 decode_exact_event_stage_cursor(bytes).map(i64::from_be_bytes) 3248 }) 3249 } 3250 3251 fn decode_exact_event_stage_cursor(cursor: &[u8]) -> Result<[u8; 8], Error> { 3252 let exact = <[u8; 8]>::try_from(cursor).map_err(|_| Error::InvalidLegacyImportJournal)?; 3253 if i64::from_be_bytes(exact) <= 0 { 3254 return Err(Error::InvalidLegacyImportJournal); 3255 } 3256 Ok(exact) 3257 } 3258 3259 fn journal_matches_classified( 3260 journal: &LegacyImportJournal, 3261 classified: &ClassifiedLegacyImport, 3262 classification_sha256: MemberDigest, 3263 ) -> bool { 3264 let fixed_fields_match = ![ 3265 journal.import_id() == classified.import_id(), 3266 journal.target_generation() == classified.target_generation(), 3267 journal.manifest_sha256() == classified.prepared.manifest_sha256(), 3268 journal.classification_sha256() == classification_sha256, 3269 journal.members().len() == classified.sources().len(), 3270 ] 3271 .contains(&false); 3272 fixed_fields_match 3273 & journal 3274 .members() 3275 .iter() 3276 .zip(classified.sources()) 3277 .all(|(durable, expected)| durable.classification() == expected) 3278 } 3279 3280 fn decode_import_id(bytes: Vec<u8>) -> Result<LegacyImportId, Error> { 3281 LegacyImportId::new(decode_array(bytes)?) 3282 } 3283 3284 fn decode_generation(bytes: Vec<u8>) -> Result<SourceGeneration, Error> { 3285 SourceGeneration::new(decode_array(bytes)?).map_err(|_| Error::InvalidLegacyImportJournal) 3286 } 3287 3288 fn decode_digest(bytes: Vec<u8>) -> Result<MemberDigest, Error> { 3289 Ok(MemberDigest::new(decode_array(bytes)?)) 3290 } 3291 3292 fn decode_array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> { 3293 bytes 3294 .try_into() 3295 .map_err(|_| Error::InvalidLegacyImportJournal) 3296 } 3297 3298 fn decode_positive_time(value: i64) -> Result<u64, Error> { 3299 let value = u64::try_from(value).map_err(|_| Error::InvalidLegacyImportJournal)?; 3300 if value == 0 { 3301 Err(Error::InvalidLegacyImportJournal) 3302 } else { 3303 Ok(value) 3304 } 3305 } 3306 3307 fn parse_source_kind(value: &str) -> Result<LegacySourceKind, Error> { 3308 match value { 3309 "event_store" => Ok(LegacySourceKind::EventStore), 3310 "outbox" => Ok(LegacySourceKind::Outbox), 3311 "private" => Ok(LegacySourceKind::Private), 3312 "studio" => Ok(LegacySourceKind::Studio), 3313 _ => Err(Error::InvalidLegacyImportJournal), 3314 } 3315 } 3316 3317 fn parse_legacy_schema(value: &str) -> Result<LegacySchema, Error> { 3318 match value { 3319 "event_store_v1" => Ok(LegacySchema::EventStoreV1), 3320 "event_store_v2" => Ok(LegacySchema::EventStoreV2), 3321 "event_store_v3" => Ok(LegacySchema::EventStoreV3), 3322 "event_store_v4" => Ok(LegacySchema::EventStoreV4), 3323 "outbox_v1" => Ok(LegacySchema::OutboxV1), 3324 "private_v1" => Ok(LegacySchema::PrivateV1), 3325 "studio_v1_host_handoff" => Ok(LegacySchema::StudioV1HostHandoff), 3326 _ => Err(Error::InvalidLegacyImportJournal), 3327 } 3328 } 3329 3330 const fn expected_user_version(schema: LegacySchema) -> u32 { 3331 match schema { 3332 LegacySchema::PrivateV1 => 1, 3333 LegacySchema::EventStoreV1 3334 | LegacySchema::EventStoreV2 3335 | LegacySchema::EventStoreV3 3336 | LegacySchema::EventStoreV4 3337 | LegacySchema::OutboxV1 3338 | LegacySchema::StudioV1HostHandoff => 0, 3339 } 3340 } 3341 3342 const fn schema_source_kind(schema: LegacySchema) -> LegacySourceKind { 3343 match schema { 3344 LegacySchema::EventStoreV1 3345 | LegacySchema::EventStoreV2 3346 | LegacySchema::EventStoreV3 3347 | LegacySchema::EventStoreV4 => LegacySourceKind::EventStore, 3348 LegacySchema::OutboxV1 => LegacySourceKind::Outbox, 3349 LegacySchema::PrivateV1 => LegacySourceKind::Private, 3350 LegacySchema::StudioV1HostHandoff => LegacySourceKind::Studio, 3351 } 3352 } 3353 3354 fn journal_member_states_are_consistent( 3355 state: LegacyImportState, 3356 members: &[LegacyImportMemberJournal], 3357 ) -> bool { 3358 members.iter().all(|member| match state { 3359 LegacyImportState::Classified => member.state() == LegacyImportMemberState::Pending, 3360 LegacyImportState::Staging => true, 3361 LegacyImportState::Ready => matches!( 3362 member.state(), 3363 LegacyImportMemberState::Ready | LegacyImportMemberState::Complete 3364 ), 3365 LegacyImportState::Committing => matches!( 3366 member.state(), 3367 LegacyImportMemberState::Ready | LegacyImportMemberState::Complete 3368 ), 3369 LegacyImportState::Complete => member.state() == LegacyImportMemberState::Complete, 3370 }) 3371 } 3372 3373 fn parse_import_state(value: &str) -> Result<LegacyImportState, Error> { 3374 match value { 3375 "classified" => Ok(LegacyImportState::Classified), 3376 "staging" => Ok(LegacyImportState::Staging), 3377 "ready" => Ok(LegacyImportState::Ready), 3378 "committing" => Ok(LegacyImportState::Committing), 3379 "complete" => Ok(LegacyImportState::Complete), 3380 _ => Err(Error::InvalidLegacyImportJournal), 3381 } 3382 } 3383 3384 fn parse_member_state(value: &str) -> Result<LegacyImportMemberState, Error> { 3385 match value { 3386 "pending" => Ok(LegacyImportMemberState::Pending), 3387 "staging" => Ok(LegacyImportMemberState::Staging), 3388 "ready" => Ok(LegacyImportMemberState::Ready), 3389 "complete" => Ok(LegacyImportMemberState::Complete), 3390 _ => Err(Error::InvalidLegacyImportJournal), 3391 } 3392 } 3393 3394 #[cfg_attr(coverage_nightly, coverage(off))] 3395 fn write_manifest( 3396 plan: &LegacyImportPlan, 3397 target_generation: SourceGeneration, 3398 snapshots: &[LegacySourceSnapshot], 3399 path: &Path, 3400 ) -> Result<(), Error> { 3401 let mut body = format!( 3402 "schema_version=1\nimport_id={}\ntarget_generation={}\nrequested_at_unix_ms={}\n", 3403 encode_id(plan.import_id().as_bytes()), 3404 encode_digest(target_generation.as_bytes()), 3405 plan.requested_at_unix_ms() 3406 ); 3407 for (source, evidence) in plan.sources().iter().zip(snapshots) { 3408 body.push_str("member="); 3409 body.push_str(evidence.kind().as_str()); 3410 body.push('|'); 3411 body.push_str(&encode_hex(source.path().as_os_str().as_encoded_bytes())); 3412 body.push('|'); 3413 body.push_str(evidence.relative_path()); 3414 body.push('|'); 3415 body.push_str(&evidence.byte_length().to_string()); 3416 body.push('|'); 3417 body.push_str(&encode_digest(evidence.sha256().as_bytes())); 3418 body.push('\n'); 3419 } 3420 let mut options = fs::OpenOptions::new(); 3421 options.create_new(true).write(true); 3422 #[cfg(unix)] 3423 { 3424 use std::os::unix::fs::OpenOptionsExt; 3425 options.mode(0o600); 3426 } 3427 let mut file = options 3428 .open(path) 3429 .map_err(|source| Error::LegacyImportFilesystem { 3430 operation: "create legacy import manifest", 3431 source, 3432 })?; 3433 use std::io::Write; 3434 file.write_all(body.as_bytes()) 3435 .and_then(|()| file.sync_all()) 3436 .map_err(|source| Error::LegacyImportFilesystem { 3437 operation: "persist legacy import manifest", 3438 source, 3439 }) 3440 } 3441 3442 #[cfg_attr(coverage_nightly, coverage(off))] 3443 fn validate_source_path(path: &Path) -> Result<(), Error> { 3444 if !path.is_absolute() 3445 || path.to_str().is_none() 3446 || path 3447 .components() 3448 .any(|component| matches!(component, Component::CurDir | Component::ParentDir)) 3449 { 3450 return Err(Error::InvalidLegacySource(path.to_path_buf())); 3451 } 3452 match fs::symlink_metadata(path) { 3453 Ok(metadata) if metadata.is_file() && !metadata.file_type().is_symlink() => Ok(()), 3454 Ok(_) => Err(Error::InvalidLegacySource(path.to_path_buf())), 3455 Err(_) => Err(Error::InvalidLegacySource(path.to_path_buf())), 3456 } 3457 } 3458 3459 #[cfg_attr(coverage_nightly, coverage(off))] 3460 fn paths_refer_to_same_file(left: &Path, right: &Path) -> Result<bool, Error> { 3461 let left_canonical = 3462 fs::canonicalize(left).map_err(|source| Error::LegacyImportFilesystem { 3463 operation: "resolve legacy import source identity", 3464 source, 3465 })?; 3466 let right_canonical = 3467 fs::canonicalize(right).map_err(|source| Error::LegacyImportFilesystem { 3468 operation: "resolve owned storage identity", 3469 source, 3470 })?; 3471 if left_canonical == right_canonical { 3472 return Ok(true); 3473 } 3474 #[cfg(unix)] 3475 { 3476 use std::os::unix::fs::MetadataExt; 3477 let left_metadata = fs::metadata(left).map_err(|source| Error::LegacyImportFilesystem { 3478 operation: "inspect legacy import source identity", 3479 source, 3480 })?; 3481 let right_metadata = 3482 fs::metadata(right).map_err(|source| Error::LegacyImportFilesystem { 3483 operation: "inspect owned storage identity", 3484 source, 3485 })?; 3486 Ok(left_metadata.dev() == right_metadata.dev() 3487 && left_metadata.ino() == right_metadata.ino()) 3488 } 3489 #[cfg(not(unix))] 3490 Ok(false) 3491 } 3492 3493 #[cfg_attr(coverage_nightly, coverage(off))] 3494 fn sync_directory(path: &Path, operation: &'static str) -> Result<(), Error> { 3495 File::open(path) 3496 .and_then(|directory| directory.sync_all()) 3497 .map_err(|source| Error::LegacyImportFilesystem { operation, source }) 3498 } 3499 3500 fn encode_id(bytes: &[u8; 16]) -> String { 3501 encode_hex(bytes) 3502 } 3503 3504 fn encode_digest(bytes: &[u8; 32]) -> String { 3505 encode_hex(bytes) 3506 } 3507 3508 fn encode_hex(bytes: &[u8]) -> String { 3509 const HEX: &[u8; 16] = b"0123456789abcdef"; 3510 let mut encoded = String::with_capacity(bytes.len() * 2); 3511 for byte in bytes { 3512 encoded.push(char::from(HEX[usize::from(byte >> 4)])); 3513 encoded.push(char::from(HEX[usize::from(byte & 0x0f)])); 3514 } 3515 encoded 3516 } 3517 3518 const fn bytes_are_zero<const N: usize>(bytes: &[u8; N]) -> bool { 3519 let mut index = 0; 3520 while index < N { 3521 if bytes[index] != 0 { 3522 return false; 3523 } 3524 index += 1; 3525 } 3526 true 3527 } 3528 3529 #[cfg(test)] 3530 #[cfg_attr(coverage_nightly, coverage(off))] 3531 mod tests { 3532 use radroots_event::{SignedEvent, wire::Nip01EventWire}; 3533 use radroots_storage::event::SourceGeneration; 3534 use serde::Deserialize; 3535 3536 use crate::{OpenMode, OpenOptions, Paths}; 3537 3538 use super::*; 3539 3540 const POLICY: &str = 3541 include_str!("../../../contracts/storage/legacy_import_backup_policy_v1.toml"); 3542 const CLASSIFICATION_POLICY: &str = 3543 include_str!("../../../contracts/storage/legacy_schema_classification_v1.toml"); 3544 const JOURNAL_POLICY: &str = 3545 include_str!("../../../contracts/storage/legacy_import_journal_policy_v1.toml"); 3546 const EVENT_STAGING_POLICY: &str = 3547 include_str!("../../../contracts/storage/legacy_event_staging_policy_v1.toml"); 3548 const OUTBOX_STAGING_POLICY: &str = 3549 include_str!("../../../contracts/storage/legacy_outbox_staging_policy_v1.toml"); 3550 const EVENT_STORE_V1_SQL: &str = include_str!("fixtures/legacy_event_store_v1.sql"); 3551 const OUTBOX_V1_SQL: &str = include_str!("fixtures/legacy_outbox_v1.sql"); 3552 const PRIVATE_STAGING_POLICY: &str = 3553 include_str!("../../../contracts/storage/legacy_private_staging_policy_v1.toml"); 3554 const STUDIO_HANDOFF_POLICY: &str = 3555 include_str!("../../../contracts/storage/legacy_studio_handoff_policy_v1.toml"); 3556 const IMPORT_VALIDATION_POLICY: &str = 3557 include_str!("../../../contracts/storage/legacy_import_validation_policy_v1.toml"); 3558 const IMPORT_FINALIZE_POLICY: &str = 3559 include_str!("../../../contracts/storage/legacy_import_finalize_policy_v1.toml"); 3560 const IMPORT_QUALIFICATION_POLICY: &str = 3561 include_str!("../../../contracts/storage/legacy_import_qualification_v1.toml"); 3562 const PRIVATE_STORE_V1_SQL: &str = include_str!("fixtures/private_store_v1.sql"); 3563 3564 #[derive(Deserialize)] 3565 struct Policy { 3566 schema_version: u32, 3567 mode: String, 3568 source_kinds: Vec<String>, 3569 source_path: String, 3570 source_cardinality: String, 3571 owned_file_alias: String, 3572 authority: String, 3573 capture: String, 3574 staging: String, 3575 finalized: String, 3576 member_mode: String, 3577 manifest: String, 3578 verification: Vec<String>, 3579 finalization: String, 3580 mutation_before_finalized_backup: bool, 3581 collision: String, 3582 hidden_entropy_or_clock: bool, 3583 } 3584 3585 #[derive(Deserialize)] 3586 struct ClassificationPolicy { 3587 schema_version: u32, 3588 catalog_algorithm: String, 3589 unknown_objects: String, 3590 mixed_source_families: String, 3591 newer_versions: String, 3592 target_mutation: bool, 3593 event_store: EventStorePolicy, 3594 outbox: FixedSchemaPolicy, 3595 private: FixedSchemaPolicy, 3596 studio: StudioSchemaPolicy, 3597 } 3598 3599 #[derive(Deserialize)] 3600 struct EventStorePolicy { 3601 versions: Vec<u32>, 3602 unledgered_version: u32, 3603 ledger: String, 3604 schema_sha256: Vec<String>, 3605 names: Vec<String>, 3606 up_sha256: Vec<String>, 3607 down_sha256: Vec<String>, 3608 } 3609 3610 #[derive(Deserialize)] 3611 struct FixedSchemaPolicy { 3612 version: u32, 3613 user_version: i64, 3614 catalog_sha256: String, 3615 source: String, 3616 schema_sql_sha256: String, 3617 } 3618 3619 #[derive(Deserialize)] 3620 struct StudioSchemaPolicy { 3621 version: u32, 3622 user_version: i64, 3623 catalog_sha256: String, 3624 source: String, 3625 schema_sql_sha256: String, 3626 disposition: String, 3627 } 3628 3629 #[derive(Deserialize)] 3630 struct JournalPolicy { 3631 schema_version: u32, 3632 authority: String, 3633 identity: Vec<String>, 3634 import_states: Vec<String>, 3635 member_states: Vec<String>, 3636 imports_per_target_generation: u32, 3637 begin: String, 3638 resume: String, 3639 source_members: String, 3640 resume_cursor: String, 3641 staged_row_count: String, 3642 host_timestamp: String, 3643 hidden_clock_or_entropy: bool, 3644 legacy_row_conversion: bool, 3645 live_product_row_mutation: bool, 3646 } 3647 3648 #[derive(Deserialize)] 3649 struct EventStagingPolicy { 3650 schema_version: u32, 3651 authority: String, 3652 source_kind: String, 3653 supported_schemas: Vec<String>, 3654 source_table: String, 3655 ordering: String, 3656 cursor: String, 3657 page_limit_max: u16, 3658 conversion: Vec<String>, 3659 transaction: String, 3660 idempotency: String, 3661 staging_rows: String, 3662 evidence_revalidation: String, 3663 live_product_row_mutation: bool, 3664 hidden_clock_or_entropy: bool, 3665 } 3666 3667 #[derive(Deserialize)] 3668 struct OutboxStagingPolicy { 3669 schema_version: u32, 3670 authority: String, 3671 source_kind: String, 3672 supported_schema: String, 3673 table_order: Vec<String>, 3674 cursor: String, 3675 page_limit_max: u16, 3676 record: String, 3677 references: Vec<String>, 3678 transaction: String, 3679 idempotency: String, 3680 staging_rows: String, 3681 evidence_revalidation: String, 3682 live_product_row_mutation: bool, 3683 hidden_clock_or_entropy: bool, 3684 } 3685 3686 #[derive(Deserialize)] 3687 struct PrivateStagingPolicy { 3688 schema_version: u32, 3689 runtime_authority: String, 3690 private_authority: String, 3691 source_kind: String, 3692 supported_schema: String, 3693 table_order: Vec<String>, 3694 cursor: String, 3695 page_limit_max: u16, 3696 record: String, 3697 secret_bearing_staging_database: String, 3698 wrapping_key_reference: String, 3699 recovery: Vec<String>, 3700 crash_before_private_commit: String, 3701 crash_after_private_commit: String, 3702 conflicting_replay: String, 3703 live_private_artifact_mutation: bool, 3704 hidden_clock_or_entropy: bool, 3705 } 3706 3707 #[derive(Deserialize)] 3708 struct StudioHandoffPolicy { 3709 schema_version: u32, 3710 source_kind: String, 3711 supported_schema: String, 3712 disposition: String, 3713 evidence: String, 3714 handoff_identity: String, 3715 receipt: String, 3716 receipt_cursor_bytes: usize, 3717 staged_row_count: u64, 3718 exact_retry: String, 3719 conflicting_retry: String, 3720 sdk_runtime_row_import: bool, 3721 sdk_private_row_import: bool, 3722 sdk_owned_studio_database: bool, 3723 source_deletion: bool, 3724 hidden_clock_or_entropy: bool, 3725 } 3726 3727 #[derive(Deserialize)] 3728 struct ImportValidationPolicy { 3729 schema_version: u32, 3730 required_import_state: String, 3731 required_member_state: String, 3732 source_evidence: String, 3733 source_count_match: String, 3734 studio_source_count: u64, 3735 snapshot: String, 3736 validation_identity: String, 3737 runtime_staging_rows: Vec<String>, 3738 private_staging_rows: Vec<String>, 3739 studio_receipt: String, 3740 validation_mutation: bool, 3741 source_deletion: bool, 3742 dual_write: bool, 3743 hidden_clock_or_entropy: bool, 3744 } 3745 3746 #[derive(Deserialize)] 3747 struct ImportFinalizePolicy { 3748 schema_version: u32, 3749 input: String, 3750 commit_order: Vec<String>, 3751 private_replay: String, 3752 runtime_replay: String, 3753 crash_before_private_commit: String, 3754 crash_after_private_commit: String, 3755 crash_during_runtime_completion: String, 3756 lost_success_response: String, 3757 retained_representation: String, 3758 live_product_dual_write: bool, 3759 source_deletion: bool, 3760 studio_row_import: bool, 3761 host_timestamp: String, 3762 hidden_clock_or_entropy: bool, 3763 } 3764 3765 #[derive(Deserialize)] 3766 struct ImportQualificationPolicy { 3767 schema_version: u32, 3768 source_matrix: Vec<String>, 3769 required_cases: Vec<String>, 3770 mixed_imported_row_count: u64, 3771 mixed_host_handoff_row_count: u64, 3772 exact_retry: bool, 3773 hidden_clock_or_entropy: bool, 3774 } 3775 3776 fn generation(byte: u8) -> SourceGeneration { 3777 SourceGeneration::new([byte; 32]).expect("source generation") 3778 } 3779 3780 async fn target(directory: &Path) -> (Paths, SqliteStorage) { 3781 let paths = Paths::from_directory(directory).expect("target paths"); 3782 let store = SqliteStorage::open( 3783 OpenOptions::new(paths.clone(), OpenMode::Create) 3784 .with_source_generation(generation(121), 12_100) 3785 .expect("source generation"), 3786 ) 3787 .await 3788 .expect("target storage"); 3789 (paths, store) 3790 } 3791 3792 async fn legacy_database(path: &Path, table: &'static str) -> SqliteConnection { 3793 let mut connection = SqliteConnection::connect_with( 3794 &SqliteConnectOptions::new() 3795 .filename(path) 3796 .create_if_missing(true), 3797 ) 3798 .await 3799 .expect("legacy database"); 3800 sqlx::query("PRAGMA journal_mode = WAL") 3801 .execute(&mut connection) 3802 .await 3803 .expect("legacy WAL"); 3804 let (create, insert) = match table { 3805 "event_envelopes" => ( 3806 "CREATE TABLE event_envelopes(value INTEGER NOT NULL)", 3807 "INSERT INTO event_envelopes(value) VALUES (41)", 3808 ), 3809 "sdk_studio_state" => ( 3810 "CREATE TABLE sdk_studio_state(value INTEGER NOT NULL)", 3811 "INSERT INTO sdk_studio_state(value) VALUES (41)", 3812 ), 3813 _ => panic!("unsupported legacy test table"), 3814 }; 3815 sqlx::query(create) 3816 .execute(&mut connection) 3817 .await 3818 .expect("legacy schema"); 3819 sqlx::query(insert) 3820 .execute(&mut connection) 3821 .await 3822 .expect("legacy row"); 3823 connection 3824 } 3825 3826 async fn supported_studio_database(path: &Path) -> SqliteConnection { 3827 let mut connection = SqliteConnection::connect_with( 3828 &SqliteConnectOptions::new() 3829 .filename(path) 3830 .create_if_missing(true), 3831 ) 3832 .await 3833 .expect("supported Studio database"); 3834 sqlx::query( 3835 "CREATE TABLE sdk_studio_state ( 3836 key TEXT PRIMARY KEY NOT NULL, 3837 value_json TEXT NOT NULL, 3838 updated_at_ms INTEGER NOT NULL 3839 )", 3840 ) 3841 .execute(&mut connection) 3842 .await 3843 .expect("supported Studio schema"); 3844 sqlx::query( 3845 "INSERT INTO sdk_studio_state(key, value_json, updated_at_ms) VALUES ('theme', '{}', 1)", 3846 ) 3847 .execute(&mut connection) 3848 .await 3849 .expect("supported Studio row"); 3850 connection 3851 } 3852 3853 fn signed_event(content: &str) -> SignedEvent { 3854 let mut wire = Nip01EventWire { 3855 id: "0".repeat(64), 3856 pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), 3857 created_at: 1_800_000_100, 3858 kind: 1, 3859 tags: vec![], 3860 content: content.to_owned(), 3861 sig: "42".repeat(64), 3862 extra: Default::default(), 3863 }; 3864 wire.id = wire 3865 .computed_event_id() 3866 .expect("canonical event id") 3867 .to_hex(); 3868 let raw_json = serde_json::json!({ 3869 "id": &wire.id, 3870 "pubkey": &wire.pubkey, 3871 "created_at": wire.created_at, 3872 "kind": wire.kind, 3873 "tags": &wire.tags, 3874 "content": &wire.content, 3875 "sig": &wire.sig, 3876 }) 3877 .to_string(); 3878 SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") 3879 } 3880 3881 async fn supported_event_database(path: &Path, events: &[SignedEvent]) -> SqliteConnection { 3882 let mut connection = SqliteConnection::connect_with( 3883 &SqliteConnectOptions::new() 3884 .filename(path) 3885 .create_if_missing(true), 3886 ) 3887 .await 3888 .expect("supported event database"); 3889 sqlx::raw_sql(EVENT_STORE_V1_SQL) 3890 .execute(&mut connection) 3891 .await 3892 .expect("event-store v1 schema"); 3893 for (index, event) in events.iter().enumerate() { 3894 let timestamp = 13_000_i64 + i64::try_from(index).expect("event index"); 3895 sqlx::query( 3896 "INSERT INTO event_envelopes( 3897 event_id, pubkey, created_at, kind, tags_json, content, sig, 3898 raw_json, verification_status, contract_status, contract_id, 3899 event_class, projection_eligible, inserted_at_ms, updated_at_ms 3900 ) VALUES (?, ?, ?, ?, '[]', ?, ?, ?, 'verified', 'admitted', 3901 NULL, NULL, 1, ?, ?)", 3902 ) 3903 .bind(event.id().to_hex()) 3904 .bind(event.envelope().author().to_hex()) 3905 .bind(i64::try_from(event.created_at()).expect("created at")) 3906 .bind(i64::from(event.kind())) 3907 .bind(event.content()) 3908 .bind(event.signature_hex()) 3909 .bind(event.raw_json()) 3910 .bind(timestamp) 3911 .bind(timestamp) 3912 .execute(&mut connection) 3913 .await 3914 .expect("legacy event"); 3915 } 3916 connection 3917 } 3918 3919 async fn supported_outbox_database(path: &Path) -> SqliteConnection { 3920 let mut connection = SqliteConnection::connect_with( 3921 &SqliteConnectOptions::new() 3922 .filename(path) 3923 .create_if_missing(true), 3924 ) 3925 .await 3926 .expect("supported outbox database"); 3927 sqlx::raw_sql(OUTBOX_V1_SQL) 3928 .execute(&mut connection) 3929 .await 3930 .expect("outbox v1 schema"); 3931 sqlx::query( 3932 "INSERT INTO outbox_operations( 3933 operation_kind, expected_pubkey, semantic_scope, trade_id, mutation_id, 3934 canonical_payload_sha256, idempotency_key, operation_idempotency_digest, 3935 status, created_at_ms, updated_at_ms 3936 ) VALUES ('publish', 'author', 'generic_event', NULL, NULL, NULL, 3937 'key', 'operation-digest', 'queued', 1, 1)", 3938 ) 3939 .execute(&mut connection) 3940 .await 3941 .expect("outbox operation"); 3942 sqlx::query( 3943 "INSERT INTO outbox_event( 3944 operation_id, event_id, expected_pubkey, draft_json, signed_event_json, 3945 raw_event_json, state, attempt_count, claim_token, claim_owner, 3946 claim_expires_at_ms, active_delivery_plan_id, next_attempt_after_ms, 3947 last_error, event_store_ingested, event_store_inserted, 3948 event_store_ingested_at_ms, created_at_ms, updated_at_ms 3949 ) VALUES (1, 'event', 'author', '{}', NULL, NULL, 'draft_queued', 0, 3950 NULL, NULL, NULL, NULL, 1, NULL, 0, 0, NULL, 1, 1)", 3951 ) 3952 .execute(&mut connection) 3953 .await 3954 .expect("outbox event"); 3955 sqlx::query( 3956 "INSERT INTO outbox_delivery_plan( 3957 outbox_event_id, transport_profile_id, target_policy_fingerprint, 3958 target_policy_version, satisfaction_policy, required_success_count, 3959 delivery_plan_idempotency_digest, status, satisfied_at_ms, 3960 created_at_ms, updated_at_ms 3961 ) VALUES (1, 'nostr', 'policy', 1, 'all', 1, 'plan-digest', 3962 'queued', NULL, 1, 1)", 3963 ) 3964 .execute(&mut connection) 3965 .await 3966 .expect("outbox plan"); 3967 sqlx::query( 3968 "INSERT INTO outbox_delivery_target( 3969 delivery_plan_id, transport_kind, endpoint_uri, target_scope, 3970 target_label, endpoint_fingerprint, status, last_outcome_kind, 3971 attempt_count, last_attempt_at_ms, completed_at_ms, last_error 3972 ) VALUES (1, 'nostr', 'wss://relay.example', NULL, NULL, 'endpoint', 3973 'pending', NULL, 0, NULL, NULL, NULL)", 3974 ) 3975 .execute(&mut connection) 3976 .await 3977 .expect("outbox target"); 3978 sqlx::query( 3979 "INSERT INTO outbox_delivery_attempt( 3980 delivery_plan_id, delivery_target_id, status, outcome_kind, 3981 attempted_at_ms, message 3982 ) VALUES (1, 1, 'complete', 'accepted', 2, 'accepted')", 3983 ) 3984 .execute(&mut connection) 3985 .await 3986 .expect("outbox attempt"); 3987 connection 3988 } 3989 3990 async fn supported_private_database(path: &Path) -> SqliteConnection { 3991 let mut connection = SqliteConnection::connect_with( 3992 &SqliteConnectOptions::new() 3993 .filename(path) 3994 .create_if_missing(true), 3995 ) 3996 .await 3997 .expect("supported private database"); 3998 sqlx::raw_sql(PRIVATE_STORE_V1_SQL) 3999 .execute(&mut connection) 4000 .await 4001 .expect("private v1 schema"); 4002 sqlx::query("PRAGMA user_version = 1") 4003 .execute(&mut connection) 4004 .await 4005 .expect("private version"); 4006 sqlx::query("INSERT INTO private_metadata VALUES (1,1,?,?,1,'source',1,1)") 4007 .bind([1_u8; 16].as_slice()) 4008 .bind([2_u8; 32].as_slice()) 4009 .execute(&mut connection) 4010 .await 4011 .expect("metadata"); 4012 sqlx::query( 4013 "INSERT INTO wrapped_profile_key VALUES (1,'memory_test_wrapped_v1',?,?,1,NULL)", 4014 ) 4015 .bind([3_u8; 32].as_slice()) 4016 .bind([4_u8; 24].as_slice()) 4017 .execute(&mut connection) 4018 .await 4019 .expect("wrapped key"); 4020 sqlx::query("INSERT INTO wrapped_signing_secret VALUES (?,?,?,?,?,1,1)") 4021 .bind([5_u8; 16].as_slice()) 4022 .bind([6_u8; 32].as_slice()) 4023 .bind(1_i64) 4024 .bind([7_u8; 8].as_slice()) 4025 .bind([8_u8; 24].as_slice()) 4026 .execute(&mut connection) 4027 .await 4028 .expect("signing secret"); 4029 sqlx::query("INSERT INTO private_farm_location VALUES (30340,?,?,1,?,?,1,1)") 4030 .bind([9_u8; 32].as_slice()) 4031 .bind("farm") 4032 .bind([10_u8; 8].as_slice()) 4033 .bind([11_u8; 24].as_slice()) 4034 .execute(&mut connection) 4035 .await 4036 .expect("farm location"); 4037 sqlx::query("INSERT INTO private_trade_artifacts VALUES ('artifact','01234567890123456789012345678901',NULL,'message','schema',?,1,?,?,'retain',1,NULL,NULL)") 4038 .bind("a".repeat(64)).bind([12_u8;8].as_slice()).bind([13_u8;4].as_slice()).execute(&mut connection).await.expect("trade artifact"); 4039 sqlx::query("INSERT INTO cursor_hmac_key VALUES (?,1,?,?,1,NULL)") 4040 .bind([14_u8; 16].as_slice()) 4041 .bind([15_u8; 8].as_slice()) 4042 .bind([16_u8; 24].as_slice()) 4043 .execute(&mut connection) 4044 .await 4045 .expect("cursor key"); 4046 sqlx::query("INSERT INTO nip46_session_private VALUES (?,?,?,?,1,?,?,2,'active',1,1)") 4047 .bind([17_u8; 16].as_slice()) 4048 .bind([18_u8; 32].as_slice()) 4049 .bind([19_u8; 32].as_slice()) 4050 .bind([20_u8; 32].as_slice()) 4051 .bind([21_u8; 8].as_slice()) 4052 .bind([22_u8; 24].as_slice()) 4053 .execute(&mut connection) 4054 .await 4055 .expect("nip46 session"); 4056 sqlx::query( 4057 "INSERT INTO key_rotation_progress VALUES (1,1,2,'done',NULL,'complete',1,1,NULL)", 4058 ) 4059 .execute(&mut connection) 4060 .await 4061 .expect("rotation progress"); 4062 connection 4063 } 4064 4065 #[test] 4066 fn implementation_matches_the_governed_legacy_backup_policy() { 4067 let policy = toml::from_str::<Policy>(POLICY).expect("legacy backup policy"); 4068 assert_eq!(policy.schema_version, 1); 4069 assert_eq!(policy.mode, "explicit_forward_only_one_shot"); 4070 assert_eq!( 4071 policy.source_kinds, 4072 ["event_store", "outbox", "private", "studio"] 4073 ); 4074 assert_eq!( 4075 policy.source_path, 4076 "absolute_existing_utf8_regular_file_no_symlink" 4077 ); 4078 assert_eq!( 4079 policy.source_cardinality, 4080 "one_to_four_unique_kinds_and_paths" 4081 ); 4082 assert_eq!( 4083 policy.owned_file_alias, 4084 "reject_path_canonical_or_file_identity" 4085 ); 4086 assert_eq!(policy.authority, "open_writable_target_backend"); 4087 assert_eq!(policy.capture, "sqlite_vacuum_into"); 4088 assert_eq!( 4089 policy.staging, 4090 ".radroots-legacy-import-{import_id_hex}.staging" 4091 ); 4092 assert_eq!(policy.finalized, "radroots-legacy-import-{import_id_hex}"); 4093 assert_eq!(policy.member_mode, "0600"); 4094 assert_eq!( 4095 policy.manifest, 4096 "manifest.v1_exact_identity_target_generation_timestamp_source_provenance_inventory_lengths_sha256" 4097 ); 4098 assert_eq!( 4099 policy.verification, 4100 [ 4101 "sqlite_quick_check", 4102 "foreign_key_check", 4103 "exact_length", 4104 "sha256" 4105 ] 4106 ); 4107 assert_eq!(policy.finalization, "same_root_atomic_directory_rename"); 4108 assert!(!policy.mutation_before_finalized_backup); 4109 assert_eq!(policy.collision, "reject"); 4110 assert!(!policy.hidden_entropy_or_clock); 4111 } 4112 4113 #[test] 4114 fn implementation_matches_the_governed_schema_classification_policy() { 4115 let policy = toml::from_str::<ClassificationPolicy>(CLASSIFICATION_POLICY) 4116 .expect("legacy classification policy"); 4117 assert_eq!(policy.schema_version, 1); 4118 assert_eq!( 4119 policy.catalog_algorithm, 4120 "sqlite_schema_non_internal_type_name_table_sql_nul_sha256_v1" 4121 ); 4122 assert_eq!(policy.unknown_objects, "reject"); 4123 assert_eq!(policy.mixed_source_families, "reject"); 4124 assert_eq!(policy.newer_versions, "reject"); 4125 assert!(!policy.target_mutation); 4126 assert_eq!(policy.event_store.versions, [1, 2, 3, 4]); 4127 assert_eq!(policy.event_store.unledgered_version, 1); 4128 assert_eq!( 4129 policy.event_store.ledger, 4130 "radroots_event_store_schema_migrations_exact_catalog_and_contiguous_rows" 4131 ); 4132 assert_eq!( 4133 policy.event_store.schema_sha256, 4134 EVENT_STORE_MIGRATIONS 4135 .iter() 4136 .map(|migration| migration.schema_sha256) 4137 .collect::<Vec<_>>() 4138 ); 4139 assert_eq!( 4140 policy.event_store.names, 4141 EVENT_STORE_MIGRATIONS 4142 .iter() 4143 .map(|migration| migration.name) 4144 .collect::<Vec<_>>() 4145 ); 4146 assert_eq!( 4147 policy.event_store.up_sha256, 4148 EVENT_STORE_MIGRATIONS 4149 .iter() 4150 .map(|migration| migration.up_sha256) 4151 .collect::<Vec<_>>() 4152 ); 4153 assert_eq!( 4154 policy.event_store.down_sha256, 4155 EVENT_STORE_MIGRATIONS 4156 .iter() 4157 .map(|migration| migration.down_sha256) 4158 .collect::<Vec<_>>() 4159 ); 4160 assert_fixed_schema_policy( 4161 &policy.outbox, 4162 0, 4163 OUTBOX_CATALOG_SHA256, 4164 "crates/outbox/migrations/0001_outbox.up.sql", 4165 "a7ee775d32c2b9f845961425362e1b1e558ce0d025f7d22dd58f118ba4dab4fa", 4166 ); 4167 assert_fixed_schema_policy( 4168 &policy.private, 4169 1, 4170 PRIVATE_CATALOG_SHA256, 4171 "oss/sdk/crates/sdk/src/private_store.rs::PRIVATE_STORE_MIGRATION_UP", 4172 "f7e71d2cf4347f9b78bafd37980901441b09d111f15d94d260eb7133b626fbe9", 4173 ); 4174 assert_eq!(policy.studio.version, 1); 4175 assert_eq!(policy.studio.user_version, 0); 4176 assert_eq!(policy.studio.catalog_sha256, STUDIO_CATALOG_SHA256); 4177 assert_eq!( 4178 policy.studio.source, 4179 "oss/sdk/crates/sdk/src/studio_store.rs::STUDIO_STORE_MIGRATION_UP" 4180 ); 4181 assert_eq!( 4182 policy.studio.schema_sql_sha256, 4183 "9c5b33810ad9746421dc843651eb02c77d6fc8f00fd630cb835c15b4e36d0590" 4184 ); 4185 assert_eq!(policy.studio.disposition, "host_handoff_not_sdk_import"); 4186 } 4187 4188 #[test] 4189 fn implementation_matches_the_governed_import_journal_policy() { 4190 let policy = 4191 toml::from_str::<JournalPolicy>(JOURNAL_POLICY).expect("legacy journal policy"); 4192 assert_eq!(policy.schema_version, 1); 4193 assert_eq!( 4194 policy.authority, 4195 "runtime_sqlite_owned_forward_migration_v6" 4196 ); 4197 assert_eq!( 4198 policy.identity, 4199 [ 4200 "import_id", 4201 "target_generation", 4202 "manifest_sha256", 4203 "classification_sha256" 4204 ] 4205 ); 4206 assert_eq!( 4207 policy.import_states, 4208 ["classified", "staging", "ready", "committing", "complete"] 4209 ); 4210 assert_eq!( 4211 policy.member_states, 4212 ["pending", "staging", "ready", "complete"] 4213 ); 4214 assert_eq!(policy.imports_per_target_generation, 1); 4215 assert_eq!(policy.begin, "atomic_exact_idempotent_or_conflict"); 4216 assert_eq!(policy.resume, "read_exact_durable_state"); 4217 assert_eq!(policy.source_members, "one_exact_row_per_classified_source"); 4218 assert_eq!(policy.resume_cursor, "opaque_nullable_bytes"); 4219 assert_eq!(policy.staged_row_count, "non_negative"); 4220 assert_eq!(policy.host_timestamp, "positive_monotonic_per_import"); 4221 assert!(!policy.hidden_clock_or_entropy); 4222 assert!(!policy.legacy_row_conversion); 4223 assert!(!policy.live_product_row_mutation); 4224 } 4225 4226 #[test] 4227 fn implementation_matches_the_governed_event_staging_policy() { 4228 let policy = toml::from_str::<EventStagingPolicy>(EVENT_STAGING_POLICY) 4229 .expect("legacy event staging policy"); 4230 assert_eq!(policy.schema_version, 1); 4231 assert_eq!( 4232 policy.authority, 4233 "runtime_sqlite_owned_forward_migration_v7" 4234 ); 4235 assert_eq!(policy.source_kind, "event_store"); 4236 assert_eq!( 4237 policy.supported_schemas, 4238 [ 4239 "event_store_v1", 4240 "event_store_v2", 4241 "event_store_v3", 4242 "event_store_v4" 4243 ] 4244 ); 4245 assert_eq!(policy.source_table, "event_envelopes"); 4246 assert_eq!(policy.ordering, "strict_legacy_sequence_ascending"); 4247 assert_eq!(policy.cursor, "positive_i64_big_endian_8_bytes"); 4248 assert_eq!(policy.page_limit_max, LEGACY_STAGE_PAGE_LIMIT_MAX); 4249 assert_eq!( 4250 policy.conversion, 4251 [ 4252 "decode_id_verified_nip01_signed_event", 4253 "require_legacy_event_id_match", 4254 "event_id_hex_to_32_bytes", 4255 "preserve_exact_signed_event_json_bytes", 4256 "preserve_legacy_admission_evidence_without_trust_upgrade" 4257 ] 4258 ); 4259 assert_eq!( 4260 policy.transaction, 4261 "target_begin_immediate_rows_cursor_count_and_state_atomic" 4262 ); 4263 assert_eq!( 4264 policy.idempotency, 4265 "durable_cursor_exact_resume_completed_retry_noop" 4266 ); 4267 assert_eq!(policy.staging_rows, "append_only_immutable"); 4268 assert_eq!(policy.evidence_revalidation, "before_every_page"); 4269 assert!(!policy.live_product_row_mutation); 4270 assert!(!policy.hidden_clock_or_entropy); 4271 } 4272 4273 #[test] 4274 fn implementation_matches_the_governed_outbox_staging_policy() { 4275 let policy = toml::from_str::<OutboxStagingPolicy>(OUTBOX_STAGING_POLICY) 4276 .expect("legacy outbox staging policy"); 4277 assert_eq!(policy.schema_version, 1); 4278 assert_eq!( 4279 policy.authority, 4280 "runtime_sqlite_owned_forward_migration_v8" 4281 ); 4282 assert_eq!(policy.source_kind, "outbox"); 4283 assert_eq!(policy.supported_schema, "outbox_v1"); 4284 assert_eq!( 4285 policy.table_order, 4286 [ 4287 "operations", 4288 "events", 4289 "delivery_plans", 4290 "delivery_targets", 4291 "delivery_attempts" 4292 ] 4293 ); 4294 assert_eq!( 4295 policy.cursor, 4296 "table_discriminator_u8_plus_non_negative_i64_big_endian" 4297 ); 4298 assert_eq!(policy.page_limit_max, LEGACY_STAGE_PAGE_LIMIT_MAX); 4299 assert_eq!( 4300 policy.record, 4301 "sqlite_json_array_exact_governed_column_order_blob" 4302 ); 4303 assert_eq!( 4304 policy.references, 4305 [ 4306 "event_to_operation", 4307 "delivery_plan_to_event", 4308 "delivery_target_to_plan", 4309 "delivery_attempt_to_target_and_same_plan" 4310 ] 4311 ); 4312 assert_eq!( 4313 policy.transaction, 4314 "target_begin_immediate_rows_cursor_count_and_state_atomic" 4315 ); 4316 assert_eq!( 4317 policy.idempotency, 4318 "durable_table_cursor_exact_resume_completed_retry_noop" 4319 ); 4320 assert_eq!(policy.staging_rows, "append_only_immutable"); 4321 assert_eq!(policy.evidence_revalidation, "before_every_page"); 4322 assert!(!policy.live_product_row_mutation); 4323 assert!(!policy.hidden_clock_or_entropy); 4324 } 4325 4326 #[test] 4327 fn implementation_matches_the_governed_private_staging_policy() { 4328 let policy = toml::from_str::<PrivateStagingPolicy>(PRIVATE_STAGING_POLICY) 4329 .expect("private staging policy"); 4330 assert_eq!(policy.schema_version, 1); 4331 assert_eq!(policy.runtime_authority, "runtime_sqlite_import_journal_v8"); 4332 assert_eq!( 4333 policy.private_authority, 4334 "private_sqlite_forward_migration_v2" 4335 ); 4336 assert_eq!(policy.source_kind, "private"); 4337 assert_eq!(policy.supported_schema, "private_v1"); 4338 assert_eq!( 4339 policy.table_order, 4340 [ 4341 "metadata", 4342 "wrapped_profile_keys", 4343 "signing_secrets", 4344 "farm_locations", 4345 "trade_artifacts", 4346 "cursor_keys", 4347 "nip46_sessions", 4348 "rotation_progress" 4349 ] 4350 ); 4351 assert_eq!( 4352 policy.cursor, 4353 "table_discriminator_u8_plus_utf8_canonical_key_max_1024" 4354 ); 4355 assert_eq!(policy.page_limit_max, LEGACY_STAGE_PAGE_LIMIT_MAX); 4356 assert_eq!( 4357 policy.record, 4358 "sqlite_json_array_governed_column_order_with_blob_hex" 4359 ); 4360 assert_eq!(policy.secret_bearing_staging_database, "private.sqlite"); 4361 assert_eq!( 4362 policy.wrapping_key_reference, 4363 "required_before_dependent_record" 4364 ); 4365 assert_eq!( 4366 policy.recovery, 4367 [ 4368 "runtime_enter_staging", 4369 "private_exact_idempotent_page_commit", 4370 "runtime_cursor_count_commit" 4371 ] 4372 ); 4373 assert_eq!( 4374 policy.crash_before_private_commit, 4375 "no_private_rows_and_old_runtime_cursor" 4376 ); 4377 assert_eq!( 4378 policy.crash_after_private_commit, 4379 "exact_replay_verification_from_old_runtime_cursor" 4380 ); 4381 assert_eq!(policy.conflicting_replay, "reject"); 4382 assert!(!policy.live_private_artifact_mutation); 4383 assert!(!policy.hidden_clock_or_entropy); 4384 } 4385 4386 #[test] 4387 fn implementation_matches_the_governed_studio_handoff_policy() { 4388 let policy = toml::from_str::<StudioHandoffPolicy>(STUDIO_HANDOFF_POLICY) 4389 .expect("Studio handoff policy"); 4390 assert_eq!(policy.schema_version, 1); 4391 assert_eq!(policy.source_kind, "studio"); 4392 assert_eq!(policy.supported_schema, "studio_v1_host_handoff"); 4393 assert_eq!(policy.disposition, "host_handoff"); 4394 assert_eq!(policy.evidence, "immutable_preimport_backup_member"); 4395 assert_eq!( 4396 policy.handoff_identity, 4397 "sha256_domain_import_target_manifest_relative_path_source_catalog_length" 4398 ); 4399 assert_eq!( 4400 policy.receipt, 4401 "handoff_sha256_plus_nonzero_opaque_host_commitment_sha256" 4402 ); 4403 assert_eq!(policy.receipt_cursor_bytes, 64); 4404 assert_eq!(policy.staged_row_count, 0); 4405 assert_eq!(policy.exact_retry, "idempotent"); 4406 assert_eq!(policy.conflicting_retry, "reject"); 4407 assert!(!policy.sdk_runtime_row_import); 4408 assert!(!policy.sdk_private_row_import); 4409 assert!(!policy.sdk_owned_studio_database); 4410 assert!(!policy.source_deletion); 4411 assert!(!policy.hidden_clock_or_entropy); 4412 } 4413 4414 #[test] 4415 fn implementation_matches_the_governed_import_validation_policy() { 4416 let policy = toml::from_str::<ImportValidationPolicy>(IMPORT_VALIDATION_POLICY) 4417 .expect("import validation policy"); 4418 assert_eq!(policy.schema_version, 1); 4419 assert_eq!(policy.required_import_state, "ready"); 4420 assert_eq!(policy.required_member_state, "ready"); 4421 assert_eq!(policy.source_evidence, "reverified_immutable_backup"); 4422 assert_eq!(policy.source_count_match, "exact_per_member"); 4423 assert_eq!(policy.studio_source_count, 0); 4424 assert_eq!( 4425 policy.snapshot, 4426 "runtime_begin_immediate_then_private_begin_immediate" 4427 ); 4428 assert_eq!( 4429 policy.validation_identity, 4430 "sha256_framed_import_target_manifest_classification_members_staged_rows" 4431 ); 4432 assert_eq!(policy.runtime_staging_rows, ["events", "outbox_graph"]); 4433 assert_eq!(policy.private_staging_rows, ["private_records"]); 4434 assert_eq!(policy.studio_receipt, "member_cursor_only"); 4435 assert!(!policy.validation_mutation); 4436 assert!(!policy.source_deletion); 4437 assert!(!policy.dual_write); 4438 assert!(!policy.hidden_clock_or_entropy); 4439 } 4440 4441 #[test] 4442 fn implementation_matches_the_governed_import_finalize_policy() { 4443 let policy = toml::from_str::<ImportFinalizePolicy>(IMPORT_FINALIZE_POLICY) 4444 .expect("import finalize policy"); 4445 assert_eq!(policy.schema_version, 1); 4446 assert_eq!(policy.input, "exact_legacy_import_validation"); 4447 assert_eq!( 4448 policy.commit_order, 4449 ["private_commit_marker", "runtime_atomic_completion"] 4450 ); 4451 assert_eq!(policy.private_replay, "insert_or_ignore_then_exact_verify"); 4452 assert_eq!(policy.runtime_replay, "completed_receipt_exact_verify"); 4453 assert_eq!( 4454 policy.crash_before_private_commit, 4455 "journal_ready_no_private_marker" 4456 ); 4457 assert_eq!( 4458 policy.crash_after_private_commit, 4459 "journal_ready_exact_private_marker_replay" 4460 ); 4461 assert_eq!( 4462 policy.crash_during_runtime_completion, 4463 "runtime_transaction_rolls_back" 4464 ); 4465 assert_eq!( 4466 policy.lost_success_response, 4467 "exact_completed_receipt_reconstructed" 4468 ); 4469 assert_eq!( 4470 policy.retained_representation, 4471 "immutable_owned_legacy_staging" 4472 ); 4473 assert!(!policy.live_product_dual_write); 4474 assert!(!policy.source_deletion); 4475 assert!(!policy.studio_row_import); 4476 assert_eq!(policy.host_timestamp, "positive_monotonic_completion_time"); 4477 assert!(!policy.hidden_clock_or_entropy); 4478 } 4479 4480 #[test] 4481 fn implementation_matches_the_governed_import_qualification_policy() { 4482 let policy = toml::from_str::<ImportQualificationPolicy>(IMPORT_QUALIFICATION_POLICY) 4483 .expect("import qualification policy"); 4484 assert_eq!(policy.schema_version, 1); 4485 assert_eq!( 4486 policy.source_matrix, 4487 [ 4488 "event_store_v1_to_v4", 4489 "outbox_v1", 4490 "private_v1", 4491 "studio_v1_host_handoff" 4492 ] 4493 ); 4494 assert_eq!( 4495 policy.required_cases, 4496 [ 4497 "mandatory_backup", 4498 "unsupported_schema_rejection", 4499 "mixed_source_golden", 4500 "bounded_resume", 4501 "close_reopen", 4502 "invalid_row_rollback", 4503 "private_commit_recovery", 4504 "lost_success_retry", 4505 "conflicting_identity_rejection", 4506 "source_retention", 4507 "no_live_dual_write" 4508 ] 4509 ); 4510 assert_eq!(policy.mixed_imported_row_count, 14); 4511 assert_eq!(policy.mixed_host_handoff_row_count, 0); 4512 assert!(policy.exact_retry); 4513 assert!(!policy.hidden_clock_or_entropy); 4514 } 4515 4516 fn assert_fixed_schema_policy( 4517 policy: &FixedSchemaPolicy, 4518 user_version: i64, 4519 catalog_sha256: &str, 4520 source: &str, 4521 schema_sql_sha256: &str, 4522 ) { 4523 assert_eq!(policy.version, 1); 4524 assert_eq!(policy.user_version, user_version); 4525 assert_eq!(policy.catalog_sha256, catalog_sha256); 4526 assert_eq!(policy.source, source); 4527 assert_eq!(policy.schema_sql_sha256, schema_sql_sha256); 4528 } 4529 4530 #[tokio::test] 4531 async fn preparation_captures_wal_state_and_finalizes_exact_immutable_evidence() { 4532 let target_root = tempfile::tempdir().expect("target root"); 4533 let legacy_root = tempfile::tempdir().expect("legacy root"); 4534 let backup_root = tempfile::tempdir().expect("backup root"); 4535 let event_path = legacy_root.path().join("event_store.sqlite"); 4536 let studio_path = legacy_root.path().join("studio.sqlite"); 4537 let event_connection = legacy_database(&event_path, "event_envelopes").await; 4538 let studio_connection = legacy_database(&studio_path, "sdk_studio_state").await; 4539 let event_source = 4540 LegacySource::new(LegacySourceKind::EventStore, event_path).expect("event source"); 4541 let studio_source = 4542 LegacySource::new(LegacySourceKind::Studio, studio_path).expect("studio source"); 4543 let plan = LegacyImportPlan::new( 4544 LegacyImportId::new([122; 16]).expect("import id"), 4545 vec![event_source, studio_source], 4546 backup_root.path(), 4547 12_200, 4548 ) 4549 .expect("import plan"); 4550 let (target_paths, store) = target(target_root.path()).await; 4551 let owned_alias = legacy_root.path().join("owned-alias.sqlite"); 4552 fs::hard_link(target_paths.runtime(), &owned_alias).expect("owned database hard link"); 4553 let alias_plan = LegacyImportPlan::new( 4554 LegacyImportId::new([124; 16]).expect("alias import id"), 4555 vec![ 4556 LegacySource::new(LegacySourceKind::Private, owned_alias) 4557 .expect("owned database alias source"), 4558 ], 4559 backup_root.path(), 4560 12_201, 4561 ) 4562 .expect("owned alias import plan"); 4563 assert!(matches!( 4564 store.prepare_legacy_import(&alias_plan).await, 4565 Err(Error::InvalidLegacySource(_)) 4566 )); 4567 let prepared = store 4568 .prepare_legacy_import(&plan) 4569 .await 4570 .expect("prepare import"); 4571 assert_eq!(prepared.import_id(), plan.import_id()); 4572 assert_eq!(prepared.target_generation(), generation(121)); 4573 assert!(prepared.bundle_path().is_dir()); 4574 assert_eq!(prepared.snapshots().len(), 2); 4575 let manifest_path = prepared.bundle_path().join(LEGACY_MANIFEST); 4576 let manifest = fs::read_to_string(&manifest_path).expect("legacy import manifest"); 4577 assert_eq!( 4578 file_digest(&manifest_path).expect("manifest digest"), 4579 (prepared.manifest_byte_length(), prepared.manifest_sha256()) 4580 ); 4581 assert_eq!( 4582 manifest, 4583 format!( 4584 concat!( 4585 "schema_version=1\n", 4586 "import_id=7a7a7a7a7a7a7a7a7a7a7a7a7a7a7a7a\n", 4587 "target_generation={}\n", 4588 "requested_at_unix_ms=12200\n", 4589 "member=event_store|{}|event_store.sqlite|{}|{}\n", 4590 "member=studio|{}|studio.sqlite|{}|{}\n" 4591 ), 4592 encode_digest(generation(121).as_bytes()), 4593 encode_hex(plan.sources()[0].path().as_os_str().as_encoded_bytes()), 4594 prepared.snapshots()[0].byte_length(), 4595 encode_digest(prepared.snapshots()[0].sha256().as_bytes()), 4596 encode_hex(plan.sources()[1].path().as_os_str().as_encoded_bytes()), 4597 prepared.snapshots()[1].byte_length(), 4598 encode_digest(prepared.snapshots()[1].sha256().as_bytes()), 4599 ) 4600 ); 4601 #[cfg(unix)] 4602 { 4603 use std::os::unix::fs::PermissionsExt; 4604 assert_eq!( 4605 fs::metadata(&manifest_path) 4606 .expect("manifest permissions") 4607 .permissions() 4608 .mode() 4609 & 0o777, 4610 0o600 4611 ); 4612 } 4613 for evidence in prepared.snapshots() { 4614 let path = prepared.bundle_path().join(evidence.relative_path()); 4615 assert!(path.is_file()); 4616 assert_eq!(snapshot(evidence.kind(), &path).expect("rehash"), *evidence); 4617 let mut backup = SqliteConnection::connect_with( 4618 &SqliteConnectOptions::new().filename(&path).read_only(true), 4619 ) 4620 .await 4621 .expect("open immutable backup"); 4622 let select = match evidence.kind() { 4623 LegacySourceKind::EventStore => "SELECT value FROM event_envelopes", 4624 LegacySourceKind::Studio => "SELECT value FROM sdk_studio_state", 4625 _ => unreachable!("test source kind"), 4626 }; 4627 let row = sqlx::query(select) 4628 .fetch_one(&mut backup) 4629 .await 4630 .expect("latest WAL row"); 4631 assert_eq!(row.get::<i64, _>(0), 41); 4632 backup.close().await.expect("close immutable backup"); 4633 #[cfg(unix)] 4634 { 4635 use std::os::unix::fs::PermissionsExt; 4636 assert_eq!( 4637 fs::metadata(path) 4638 .expect("backup permissions") 4639 .permissions() 4640 .mode() 4641 & 0o777, 4642 0o600 4643 ); 4644 } 4645 } 4646 assert!(matches!( 4647 store.prepare_legacy_import(&plan).await, 4648 Err(Error::LegacyImportBackupAlreadyExists(_)) 4649 )); 4650 assert!(target_paths.runtime().is_file()); 4651 assert!(target_paths.private().is_file()); 4652 event_connection.close().await.expect("close event source"); 4653 studio_connection 4654 .close() 4655 .await 4656 .expect("close studio source"); 4657 store.close().await.expect("close target"); 4658 } 4659 4660 #[tokio::test] 4661 async fn classification_accepts_only_exact_studio_handoff_and_reverifies_evidence() { 4662 let target_root = tempfile::tempdir().expect("target root"); 4663 let legacy_root = tempfile::tempdir().expect("legacy root"); 4664 let backup_root = tempfile::tempdir().expect("backup root"); 4665 let studio_path = legacy_root.path().join("studio.sqlite"); 4666 let studio_connection = supported_studio_database(&studio_path).await; 4667 let plan = LegacyImportPlan::new( 4668 LegacyImportId::new([125; 16]).expect("import id"), 4669 vec![LegacySource::new(LegacySourceKind::Studio, &studio_path).expect("Studio source")], 4670 backup_root.path(), 4671 12_500, 4672 ) 4673 .expect("Studio import plan"); 4674 let (target_paths, store) = target(target_root.path()).await; 4675 let prepared = store 4676 .prepare_legacy_import(&plan) 4677 .await 4678 .expect("prepared Studio import"); 4679 let classified = store 4680 .classify_legacy_import(&prepared) 4681 .await 4682 .expect("classified Studio import"); 4683 assert_eq!(classified.import_id(), plan.import_id()); 4684 assert_eq!(classified.target_generation(), generation(121)); 4685 assert_eq!(classified.bundle_path(), prepared.bundle_path()); 4686 assert_eq!(classified.sources().len(), 1); 4687 let source = &classified.sources()[0]; 4688 assert_eq!(source.kind(), LegacySourceKind::Studio); 4689 assert_eq!(source.schema(), LegacySchema::StudioV1HostHandoff); 4690 assert_eq!( 4691 source.schema().disposition(), 4692 LegacyImportDisposition::HostHandoff 4693 ); 4694 assert_eq!(source.user_version(), 0); 4695 assert_eq!( 4696 encode_digest(source.catalog_sha256().as_bytes()), 4697 STUDIO_CATALOG_SHA256 4698 ); 4699 let journal = store 4700 .begin_legacy_import(&classified, 12_501) 4701 .await 4702 .expect("begin durable import"); 4703 assert_eq!(journal.import_id(), plan.import_id()); 4704 assert_eq!(journal.target_generation(), generation(121)); 4705 assert_eq!(journal.manifest_sha256(), prepared.manifest_sha256()); 4706 assert_eq!( 4707 journal.classification_sha256(), 4708 classification_digest(&classified) 4709 ); 4710 assert_eq!(journal.state(), LegacyImportState::Classified); 4711 assert_eq!(journal.started_at_unix_ms(), 12_501); 4712 assert_eq!(journal.updated_at_unix_ms(), 12_501); 4713 assert_eq!(journal.completed_at_unix_ms(), None); 4714 assert_eq!(journal.members().len(), 1); 4715 assert_eq!( 4716 journal.members()[0].classification(), 4717 &classified.sources()[0] 4718 ); 4719 assert_eq!( 4720 journal.members()[0].state(), 4721 LegacyImportMemberState::Pending 4722 ); 4723 assert_eq!(journal.members()[0].resume_cursor(), None); 4724 assert_eq!(journal.members()[0].staged_row_count(), 0); 4725 assert_eq!(journal.members()[0].updated_at_unix_ms(), 12_501); 4726 assert_eq!( 4727 store 4728 .begin_legacy_import(&classified, 12_599) 4729 .await 4730 .expect("idempotent begin"), 4731 journal 4732 ); 4733 assert_eq!( 4734 store 4735 .legacy_import_journal(plan.import_id()) 4736 .await 4737 .expect("read journal"), 4738 Some(journal.clone()) 4739 ); 4740 for statement in [ 4741 "UPDATE radroots_runtime_legacy_imports SET state = 'ready'", 4742 "UPDATE radroots_runtime_legacy_imports SET import_id = zeroblob(16)", 4743 "DELETE FROM radroots_runtime_legacy_imports", 4744 "UPDATE radroots_runtime_legacy_import_members SET state = 'ready'", 4745 "UPDATE radroots_runtime_legacy_import_members SET source_kind = 'outbox'", 4746 "DELETE FROM radroots_runtime_legacy_import_members", 4747 ] { 4748 assert!( 4749 sqlx::query(statement).execute(&store.pool).await.is_err(), 4750 "journal guard accepted `{statement}`" 4751 ); 4752 } 4753 4754 let conflicting_backup_root = tempfile::tempdir().expect("conflicting backup root"); 4755 let conflicting_plan = LegacyImportPlan::new( 4756 LegacyImportId::new([127; 16]).expect("conflicting import id"), 4757 vec![ 4758 LegacySource::new(LegacySourceKind::Studio, &studio_path) 4759 .expect("conflicting Studio source"), 4760 ], 4761 conflicting_backup_root.path(), 4762 12_700, 4763 ) 4764 .expect("conflicting import plan"); 4765 let conflicting_prepared = store 4766 .prepare_legacy_import(&conflicting_plan) 4767 .await 4768 .expect("conflicting prepared import"); 4769 let conflicting_classified = store 4770 .classify_legacy_import(&conflicting_prepared) 4771 .await 4772 .expect("conflicting classified import"); 4773 assert!(matches!( 4774 store 4775 .begin_legacy_import(&conflicting_classified, 12_701) 4776 .await, 4777 Err(Error::LegacyImportConflict) 4778 )); 4779 4780 let other_root = tempfile::tempdir().expect("other target root"); 4781 let other_paths = Paths::from_directory(other_root.path()).expect("other target paths"); 4782 let other_store = SqliteStorage::open( 4783 OpenOptions::new(other_paths, OpenMode::Create) 4784 .with_source_generation(generation(126), 12_600) 4785 .expect("other source generation"), 4786 ) 4787 .await 4788 .expect("other target storage"); 4789 assert!(matches!( 4790 other_store.classify_legacy_import(&prepared).await, 4791 Err(Error::LegacyImportTargetMismatch) 4792 )); 4793 assert!(matches!( 4794 other_store.begin_legacy_import(&classified, 12_601).await, 4795 Err(Error::LegacyImportTargetMismatch) 4796 )); 4797 other_store.close().await.expect("close other target"); 4798 4799 fs::write(prepared.bundle_path().join("unexpected"), b"unsupported") 4800 .expect("unexpected evidence member"); 4801 assert!(matches!( 4802 store.classify_legacy_import(&prepared).await, 4803 Err(Error::LegacyImportEvidenceInvalid) 4804 )); 4805 studio_connection 4806 .close() 4807 .await 4808 .expect("close Studio source"); 4809 store.close().await.expect("close target"); 4810 let reopened = 4811 SqliteStorage::open(OpenOptions::new(target_paths, OpenMode::ReadWriteExisting)) 4812 .await 4813 .expect("reopen target"); 4814 assert_eq!( 4815 reopened 4816 .legacy_import_journal(plan.import_id()) 4817 .await 4818 .expect("read journal after reopen"), 4819 Some(journal) 4820 ); 4821 reopened.close().await.expect("close reopened target"); 4822 } 4823 4824 #[tokio::test] 4825 async fn event_staging_is_bounded_resumable_atomic_and_isolated_from_live_rows() { 4826 let target_root = tempfile::tempdir().expect("target root"); 4827 let legacy_root = tempfile::tempdir().expect("legacy root"); 4828 let backup_root = tempfile::tempdir().expect("backup root"); 4829 let event_path = legacy_root.path().join("event_store.sqlite"); 4830 let events = [ 4831 signed_event("one"), 4832 signed_event("two"), 4833 signed_event("three"), 4834 ]; 4835 let event_connection = supported_event_database(&event_path, &events).await; 4836 let plan = LegacyImportPlan::new( 4837 LegacyImportId::new([128; 16]).expect("import id"), 4838 vec![ 4839 LegacySource::new(LegacySourceKind::EventStore, &event_path).expect("event source"), 4840 ], 4841 backup_root.path(), 4842 12_800, 4843 ) 4844 .expect("event import plan"); 4845 let (target_paths, store) = target(target_root.path()).await; 4846 let prepared = store 4847 .prepare_legacy_import(&plan) 4848 .await 4849 .expect("prepared event import"); 4850 let classified = store 4851 .classify_legacy_import(&prepared) 4852 .await 4853 .expect("classified event import"); 4854 assert_eq!(classified.sources()[0].schema(), LegacySchema::EventStoreV1); 4855 store 4856 .begin_legacy_import(&classified, 12_801) 4857 .await 4858 .expect("begin event import"); 4859 assert!(matches!( 4860 store.stage_legacy_events(&classified, 0, 12_802).await, 4861 Err(Error::InvalidLegacyImportStageRequest) 4862 )); 4863 4864 let first = store 4865 .stage_legacy_events(&classified, 2, 12_802) 4866 .await 4867 .expect("first event page"); 4868 assert_eq!(first.staged_rows(), 2); 4869 assert_eq!(first.staged_row_count(), 2); 4870 assert_eq!(first.resume_cursor(), Some(&encode_event_stage_cursor(2))); 4871 assert!(!first.is_complete()); 4872 assert_eq!( 4873 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_runtime_events") 4874 .fetch_one(&store.pool) 4875 .await 4876 .expect("live event count"), 4877 0 4878 ); 4879 assert_eq!( 4880 sqlx::query_scalar::<_, i64>( 4881 "SELECT COUNT(*) FROM radroots_runtime_legacy_event_staging", 4882 ) 4883 .fetch_one(&store.pool) 4884 .await 4885 .expect("staged event count"), 4886 2 4887 ); 4888 let journal = store 4889 .legacy_import_journal(plan.import_id()) 4890 .await 4891 .expect("staging journal") 4892 .expect("durable staging journal"); 4893 assert_eq!(journal.state(), LegacyImportState::Staging); 4894 assert_eq!( 4895 journal.members()[0].state(), 4896 LegacyImportMemberState::Staging 4897 ); 4898 assert_eq!(journal.members()[0].staged_row_count(), 2); 4899 assert_eq!( 4900 journal.members()[0].resume_cursor(), 4901 Some(encode_event_stage_cursor(2).as_slice()) 4902 ); 4903 4904 store.close().await.expect("close target between pages"); 4905 let reopened = 4906 SqliteStorage::open(OpenOptions::new(target_paths, OpenMode::ReadWriteExisting)) 4907 .await 4908 .expect("reopen target between pages"); 4909 let second = reopened 4910 .stage_legacy_events(&classified, 2, 12_803) 4911 .await 4912 .expect("second event page"); 4913 assert_eq!(second.staged_rows(), 1); 4914 assert_eq!(second.staged_row_count(), 3); 4915 assert_eq!(second.resume_cursor(), Some(&encode_event_stage_cursor(3))); 4916 assert!(second.is_complete()); 4917 let complete = reopened 4918 .legacy_import_journal(plan.import_id()) 4919 .await 4920 .expect("ready journal") 4921 .expect("durable ready journal"); 4922 assert_eq!(complete.state(), LegacyImportState::Ready); 4923 assert_eq!( 4924 complete.members()[0].state(), 4925 LegacyImportMemberState::Ready 4926 ); 4927 let retry = reopened 4928 .stage_legacy_events(&classified, 2, 12_803) 4929 .await 4930 .expect("completed retry"); 4931 assert_eq!(retry.staged_rows(), 0); 4932 assert_eq!(retry.staged_row_count(), 3); 4933 assert_eq!(retry.resume_cursor(), second.resume_cursor()); 4934 assert!(retry.is_complete()); 4935 let validation = reopened 4936 .validate_legacy_import(&classified) 4937 .await 4938 .expect("validate event import"); 4939 assert_eq!(validation.imported_row_count(), 3); 4940 assert!(!bytes_are_zero(validation.validation_sha256().as_bytes())); 4941 for statement in [ 4942 "UPDATE radroots_runtime_legacy_event_staging SET legacy_sequence = 4 WHERE legacy_sequence = 1", 4943 "DELETE FROM radroots_runtime_legacy_event_staging WHERE legacy_sequence = 1", 4944 ] { 4945 assert!( 4946 sqlx::query(statement) 4947 .execute(&reopened.pool) 4948 .await 4949 .is_err(), 4950 "staging guard accepted `{statement}`" 4951 ); 4952 } 4953 event_connection.close().await.expect("close event source"); 4954 reopened.close().await.expect("close reopened target"); 4955 } 4956 4957 #[tokio::test] 4958 async fn invalid_legacy_event_rolls_back_rows_cursor_count_and_state() { 4959 let target_root = tempfile::tempdir().expect("target root"); 4960 let legacy_root = tempfile::tempdir().expect("legacy root"); 4961 let backup_root = tempfile::tempdir().expect("backup root"); 4962 let event_path = legacy_root.path().join("event_store.sqlite"); 4963 let events = [signed_event("valid"), signed_event("invalid identity")]; 4964 let mut event_connection = supported_event_database(&event_path, &events).await; 4965 sqlx::query("UPDATE event_envelopes SET event_id = ? WHERE seq = 2") 4966 .bind("0".repeat(64)) 4967 .execute(&mut event_connection) 4968 .await 4969 .expect("corrupt legacy event identity"); 4970 let plan = LegacyImportPlan::new( 4971 LegacyImportId::new([129; 16]).expect("import id"), 4972 vec![ 4973 LegacySource::new(LegacySourceKind::EventStore, &event_path).expect("event source"), 4974 ], 4975 backup_root.path(), 4976 12_900, 4977 ) 4978 .expect("event import plan"); 4979 let (_target_paths, store) = target(target_root.path()).await; 4980 let prepared = store 4981 .prepare_legacy_import(&plan) 4982 .await 4983 .expect("prepared event import"); 4984 let classified = store 4985 .classify_legacy_import(&prepared) 4986 .await 4987 .expect("classified event import"); 4988 store 4989 .begin_legacy_import(&classified, 12_901) 4990 .await 4991 .expect("begin event import"); 4992 4993 assert!(matches!( 4994 store.stage_legacy_events(&classified, 2, 12_902).await, 4995 Err(Error::LegacyImportRowInvalid { 4996 source_kind: "event_store", 4997 legacy_sequence: 2, 4998 }) 4999 )); 5000 assert_eq!( 5001 sqlx::query_scalar::<_, i64>( 5002 "SELECT COUNT(*) FROM radroots_runtime_legacy_event_staging", 5003 ) 5004 .fetch_one(&store.pool) 5005 .await 5006 .expect("rolled-back staging count"), 5007 0 5008 ); 5009 let journal = store 5010 .legacy_import_journal(plan.import_id()) 5011 .await 5012 .expect("rolled-back journal") 5013 .expect("durable import journal"); 5014 assert_eq!(journal.state(), LegacyImportState::Classified); 5015 assert_eq!( 5016 journal.members()[0].state(), 5017 LegacyImportMemberState::Pending 5018 ); 5019 assert_eq!(journal.members()[0].resume_cursor(), None); 5020 assert_eq!(journal.members()[0].staged_row_count(), 0); 5021 event_connection.close().await.expect("close event source"); 5022 store.close().await.expect("close target"); 5023 } 5024 5025 #[tokio::test] 5026 async fn legacy_event_row_conversion_rejects_each_invalid_scalar_boundary() { 5027 let mut connection = SqliteConnection::connect("sqlite::memory:") 5028 .await 5029 .expect("row conversion database"); 5030 let event = signed_event("row conversion matrix"); 5031 #[allow(clippy::too_many_arguments)] 5032 async fn row( 5033 connection: &mut SqliteConnection, 5034 event: &SignedEvent, 5035 sequence: i64, 5036 verification_status: &str, 5037 contract_status: &str, 5038 projection_eligible: i64, 5039 inserted_at_ms: i64, 5040 updated_at_ms: i64, 5041 ) -> sqlx::sqlite::SqliteRow { 5042 sqlx::query( 5043 "SELECT ? AS seq, ? AS event_id, ? AS raw_json, 5044 ? AS verification_status, ? AS contract_status, 5045 ? AS projection_eligible, ? AS inserted_at_ms, ? AS updated_at_ms", 5046 ) 5047 .bind(sequence) 5048 .bind(event.id().to_hex()) 5049 .bind(event.raw_json()) 5050 .bind(verification_status) 5051 .bind(contract_status) 5052 .bind(projection_eligible) 5053 .bind(inserted_at_ms) 5054 .bind(updated_at_ms) 5055 .fetch_one(connection) 5056 .await 5057 .expect("legacy event row") 5058 } 5059 5060 assert!( 5061 convert_legacy_event_row( 5062 &row( 5063 &mut connection, 5064 &event, 5065 1, 5066 "verified", 5067 "accepted", 5068 1, 5069 10, 5070 11 5071 ) 5072 .await 5073 ) 5074 .is_ok() 5075 ); 5076 let long_verification = "v".repeat(65); 5077 let long_contract = "c".repeat(65); 5078 for (sequence, verification, contract, eligible, inserted, updated) in [ 5079 (0, "verified", "accepted", 1, 10, 11), 5080 (1, "", "accepted", 1, 10, 11), 5081 (1, long_verification.as_str(), "accepted", 1, 10, 11), 5082 (1, "verified", "", 1, 10, 11), 5083 (1, "verified", long_contract.as_str(), 1, 10, 11), 5084 (1, "verified", "accepted", 2, 10, 11), 5085 (1, "verified", "accepted", 1, 0, 11), 5086 (1, "verified", "accepted", 1, 10, 9), 5087 ] { 5088 let candidate = row( 5089 &mut connection, 5090 &event, 5091 sequence, 5092 verification, 5093 contract, 5094 eligible, 5095 inserted, 5096 updated, 5097 ) 5098 .await; 5099 assert!(matches!( 5100 convert_legacy_event_row(&candidate), 5101 Err(Error::LegacyImportRowInvalid { 5102 source_kind: "event_store", 5103 legacy_sequence: _ 5104 }) 5105 )); 5106 } 5107 } 5108 5109 #[tokio::test] 5110 async fn outbox_staging_resumes_across_the_exact_ordered_graph_without_live_mutation() { 5111 let target_root = tempfile::tempdir().expect("target root"); 5112 let legacy_root = tempfile::tempdir().expect("legacy root"); 5113 let backup_root = tempfile::tempdir().expect("backup root"); 5114 let outbox_path = legacy_root.path().join("outbox.sqlite"); 5115 let outbox_connection = supported_outbox_database(&outbox_path).await; 5116 let plan = LegacyImportPlan::new( 5117 LegacyImportId::new([130; 16]).expect("import id"), 5118 vec![LegacySource::new(LegacySourceKind::Outbox, &outbox_path).expect("outbox source")], 5119 backup_root.path(), 5120 13_000, 5121 ) 5122 .expect("outbox import plan"); 5123 let (target_paths, store) = target(target_root.path()).await; 5124 let prepared = store 5125 .prepare_legacy_import(&plan) 5126 .await 5127 .expect("prepared outbox import"); 5128 let classified = store 5129 .classify_legacy_import(&prepared) 5130 .await 5131 .expect("classified outbox import"); 5132 assert_eq!(classified.sources()[0].schema(), LegacySchema::OutboxV1); 5133 store 5134 .begin_legacy_import(&classified, 13_001) 5135 .await 5136 .expect("begin outbox import"); 5137 let expected = [ 5138 LegacyOutboxTable::Operations, 5139 LegacyOutboxTable::Events, 5140 LegacyOutboxTable::DeliveryPlans, 5141 LegacyOutboxTable::DeliveryTargets, 5142 ]; 5143 for (index, table) in expected.into_iter().enumerate() { 5144 let page = store 5145 .stage_legacy_outbox( 5146 &classified, 5147 1, 5148 13_002 + u64::try_from(index).expect("page index"), 5149 ) 5150 .await 5151 .expect("outbox graph page"); 5152 assert_eq!(page.table(), table); 5153 assert_eq!(page.staged_rows(), 1); 5154 assert!(!page.is_complete()); 5155 } 5156 assert_eq!( 5157 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_runtime_outbox_items") 5158 .fetch_one(&store.pool) 5159 .await 5160 .expect("live outbox count"), 5161 0 5162 ); 5163 store 5164 .close() 5165 .await 5166 .expect("close target before attempt page"); 5167 let reopened = 5168 SqliteStorage::open(OpenOptions::new(target_paths, OpenMode::ReadWriteExisting)) 5169 .await 5170 .expect("reopen target before attempt page"); 5171 let final_page = reopened 5172 .stage_legacy_outbox(&classified, 1, 13_006) 5173 .await 5174 .expect("attempt page"); 5175 assert_eq!(final_page.table(), LegacyOutboxTable::DeliveryAttempts); 5176 assert_eq!(final_page.staged_rows(), 1); 5177 assert_eq!(final_page.staged_row_count(), 5); 5178 assert_eq!( 5179 final_page.resume_cursor(), 5180 &encode_outbox_stage_cursor(LegacyOutboxTable::DeliveryAttempts, 1) 5181 ); 5182 assert!(final_page.is_complete()); 5183 assert_eq!( 5184 sqlx::query_scalar::<_, i64>( 5185 "SELECT COUNT(*) FROM radroots_runtime_legacy_outbox_staging", 5186 ) 5187 .fetch_one(&reopened.pool) 5188 .await 5189 .expect("staged graph count"), 5190 5 5191 ); 5192 let journal = reopened 5193 .legacy_import_journal(plan.import_id()) 5194 .await 5195 .expect("outbox journal") 5196 .expect("durable outbox journal"); 5197 assert_eq!(journal.state(), LegacyImportState::Ready); 5198 assert_eq!(journal.members()[0].state(), LegacyImportMemberState::Ready); 5199 assert_eq!(journal.members()[0].staged_row_count(), 5); 5200 let retry = reopened 5201 .stage_legacy_outbox(&classified, 1, 13_006) 5202 .await 5203 .expect("completed outbox retry"); 5204 assert_eq!(retry.staged_rows(), 0); 5205 assert_eq!(retry.staged_row_count(), 5); 5206 assert!(retry.is_complete()); 5207 let validation = reopened 5208 .validate_legacy_import(&classified) 5209 .await 5210 .expect("validate outbox import"); 5211 assert_eq!(validation.imported_row_count(), 5); 5212 assert!(!bytes_are_zero(validation.validation_sha256().as_bytes())); 5213 for statement in [ 5214 "UPDATE radroots_runtime_legacy_outbox_staging SET legacy_id = 2 WHERE legacy_id = 1", 5215 "DELETE FROM radroots_runtime_legacy_outbox_staging WHERE legacy_id = 1", 5216 ] { 5217 assert!( 5218 sqlx::query(statement) 5219 .execute(&reopened.pool) 5220 .await 5221 .is_err(), 5222 "outbox staging guard accepted `{statement}`" 5223 ); 5224 } 5225 outbox_connection 5226 .close() 5227 .await 5228 .expect("close outbox source"); 5229 reopened.close().await.expect("close reopened target"); 5230 } 5231 5232 #[tokio::test] 5233 async fn private_staging_recovers_exact_replay_across_both_databases() { 5234 let target_root = tempfile::tempdir().expect("target root"); 5235 let legacy_root = tempfile::tempdir().expect("legacy root"); 5236 let backup_root = tempfile::tempdir().expect("backup root"); 5237 let private_path = legacy_root.path().join("private.sqlite"); 5238 let private_connection = supported_private_database(&private_path).await; 5239 let plan = LegacyImportPlan::new( 5240 LegacyImportId::new([131; 16]).expect("import id"), 5241 vec![ 5242 LegacySource::new(LegacySourceKind::Private, &private_path) 5243 .expect("private source"), 5244 ], 5245 backup_root.path(), 5246 13_100, 5247 ) 5248 .expect("private import plan"); 5249 let (target_paths, store) = target(target_root.path()).await; 5250 let prepared = store 5251 .prepare_legacy_import(&plan) 5252 .await 5253 .expect("prepared private import"); 5254 let classified = store 5255 .classify_legacy_import(&prepared) 5256 .await 5257 .expect("classified private import"); 5258 store 5259 .begin_legacy_import(&classified, 13_101) 5260 .await 5261 .expect("begin private import"); 5262 5263 let snapshot = prepared 5264 .snapshots() 5265 .iter() 5266 .find(|snapshot| snapshot.kind() == LegacySourceKind::Private) 5267 .expect("private snapshot"); 5268 let mut evidence = SqliteConnection::connect_with( 5269 &SqliteConnectOptions::new() 5270 .filename(prepared.bundle_path().join(snapshot.relative_path())) 5271 .read_only(true), 5272 ) 5273 .await 5274 .expect("private evidence"); 5275 let row = sqlx::query(private_stage_query(LegacyPrivateTable::Metadata)) 5276 .bind("") 5277 .bind(2_i64) 5278 .fetch_one(&mut evidence) 5279 .await 5280 .expect("metadata replay row"); 5281 sqlx::query("INSERT INTO radroots_private_legacy_import_staging(import_id, table_kind, key_cursor, parent_key_version, record_json) VALUES (?, 'metadata', ?, NULL, ?)") 5282 .bind(plan.import_id().as_bytes().as_slice()).bind(row.get::<String,_>("key_cursor")).bind(row.get::<Vec<u8>,_>("record_json")).execute(store.private_pool()).await.expect("simulate private commit before runtime cursor"); 5283 evidence.close().await.expect("close evidence"); 5284 5285 let tables = [ 5286 LegacyPrivateTable::Metadata, 5287 LegacyPrivateTable::WrappedProfileKeys, 5288 LegacyPrivateTable::SigningSecrets, 5289 LegacyPrivateTable::FarmLocations, 5290 LegacyPrivateTable::TradeArtifacts, 5291 LegacyPrivateTable::CursorKeys, 5292 LegacyPrivateTable::Nip46Sessions, 5293 ]; 5294 for (index, table) in tables.into_iter().enumerate() { 5295 let page = store 5296 .stage_legacy_private( 5297 &classified, 5298 1, 5299 13_102 + u64::try_from(index).expect("page index"), 5300 ) 5301 .await 5302 .expect("private page"); 5303 assert_eq!(page.table(), table); 5304 assert_eq!(page.staged_rows(), 1); 5305 assert!(!page.is_complete()); 5306 } 5307 assert_eq!( 5308 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_private_artifacts") 5309 .fetch_one(store.private_pool()) 5310 .await 5311 .expect("live private count"), 5312 0 5313 ); 5314 store 5315 .close() 5316 .await 5317 .expect("close before final private page"); 5318 let reopened = 5319 SqliteStorage::open(OpenOptions::new(target_paths, OpenMode::ReadWriteExisting)) 5320 .await 5321 .expect("reopen private target"); 5322 let final_page = reopened 5323 .stage_legacy_private(&classified, 1, 13_109) 5324 .await 5325 .expect("rotation page"); 5326 assert_eq!(final_page.table(), LegacyPrivateTable::RotationProgress); 5327 assert_eq!(final_page.staged_rows(), 1); 5328 assert_eq!(final_page.staged_row_count(), 8); 5329 assert!(final_page.is_complete()); 5330 let retry = reopened 5331 .stage_legacy_private(&classified, 1, 13_109) 5332 .await 5333 .expect("private completed retry"); 5334 assert_eq!(retry.staged_rows(), 0); 5335 assert_eq!(retry.staged_row_count(), 8); 5336 assert!(retry.is_complete()); 5337 assert_eq!( 5338 sqlx::query_scalar::<_, i64>( 5339 "SELECT COUNT(*) FROM radroots_private_legacy_import_staging" 5340 ) 5341 .fetch_one(reopened.private_pool()) 5342 .await 5343 .expect("private staging count"), 5344 8 5345 ); 5346 let journal = reopened 5347 .legacy_import_journal(plan.import_id()) 5348 .await 5349 .expect("private journal") 5350 .expect("durable private journal"); 5351 assert_eq!(journal.state(), LegacyImportState::Ready); 5352 assert_eq!(journal.members()[0].state(), LegacyImportMemberState::Ready); 5353 let validation = reopened 5354 .validate_legacy_import(&classified) 5355 .await 5356 .expect("validate private import"); 5357 assert_eq!(validation.imported_row_count(), 8); 5358 assert!(!bytes_are_zero(validation.validation_sha256().as_bytes())); 5359 assert_eq!( 5360 reopened 5361 .validate_legacy_import(&classified) 5362 .await 5363 .expect("repeat private validation"), 5364 validation 5365 ); 5366 sqlx::query("INSERT INTO radroots_private_legacy_import_commits(import_id, validation_sha256, imported_row_count, committed_at_ms) VALUES (?, ?, 8, 13110)") 5367 .bind(plan.import_id().as_bytes().as_slice()).bind(validation.validation_sha256().as_bytes().as_slice()).execute(reopened.private_pool()).await.expect("simulate private commit before runtime completion"); 5368 let receipt = reopened 5369 .finalize_legacy_import(&classified, validation, 13_999) 5370 .await 5371 .expect("recover and finalize private import"); 5372 assert_eq!(receipt.validation_sha256(), validation.validation_sha256()); 5373 assert_eq!(receipt.imported_row_count(), 8); 5374 assert_eq!(receipt.completed_at_unix_ms(), 13_110); 5375 let completed = reopened 5376 .legacy_import_journal(plan.import_id()) 5377 .await 5378 .expect("completed journal") 5379 .expect("durable completed journal"); 5380 assert_eq!(completed.state(), LegacyImportState::Complete); 5381 assert_eq!(completed.completed_at_unix_ms(), Some(13_110)); 5382 assert_eq!( 5383 completed.members()[0].state(), 5384 LegacyImportMemberState::Complete 5385 ); 5386 assert_eq!( 5387 reopened 5388 .finalize_legacy_import(&classified, validation, 13_999) 5389 .await 5390 .expect("lost success response retry"), 5391 receipt 5392 ); 5393 assert_eq!( 5394 sqlx::query_scalar::<_, i64>( 5395 "SELECT COUNT(*) FROM radroots_runtime_legacy_import_commits" 5396 ) 5397 .fetch_one(reopened.pool()) 5398 .await 5399 .expect("runtime commit count"), 5400 1 5401 ); 5402 assert_eq!( 5403 sqlx::query_scalar::<_, i64>( 5404 "SELECT COUNT(*) FROM radroots_private_legacy_import_commits" 5405 ) 5406 .fetch_one(reopened.private_pool()) 5407 .await 5408 .expect("private commit count"), 5409 1 5410 ); 5411 for statement in [ 5412 "UPDATE radroots_runtime_legacy_import_commits SET imported_row_count = 9", 5413 "DELETE FROM radroots_runtime_legacy_import_commits", 5414 ] { 5415 assert!( 5416 sqlx::query(statement) 5417 .execute(reopened.pool()) 5418 .await 5419 .is_err(), 5420 "runtime commit guard accepted `{statement}`" 5421 ); 5422 } 5423 for statement in [ 5424 "UPDATE radroots_private_legacy_import_commits SET imported_row_count = 9", 5425 "DELETE FROM radroots_private_legacy_import_commits", 5426 ] { 5427 assert!( 5428 sqlx::query(statement) 5429 .execute(reopened.private_pool()) 5430 .await 5431 .is_err(), 5432 "private commit guard accepted `{statement}`" 5433 ); 5434 } 5435 private_connection 5436 .close() 5437 .await 5438 .expect("close private source"); 5439 reopened.close().await.expect("close private target"); 5440 } 5441 5442 #[tokio::test] 5443 async fn studio_handoff_requires_exact_host_receipt_and_imports_no_rows() { 5444 let target_root = tempfile::tempdir().expect("target root"); 5445 let legacy_root = tempfile::tempdir().expect("legacy root"); 5446 let backup_root = tempfile::tempdir().expect("backup root"); 5447 let host_root = tempfile::tempdir().expect("host root"); 5448 let studio_path = legacy_root.path().join("studio.sqlite"); 5449 let studio_connection = supported_studio_database(&studio_path).await; 5450 let plan = LegacyImportPlan::new( 5451 LegacyImportId::new([132; 16]).expect("import id"), 5452 vec![LegacySource::new(LegacySourceKind::Studio, &studio_path).expect("Studio source")], 5453 backup_root.path(), 5454 13_200, 5455 ) 5456 .expect("Studio import plan"); 5457 let (_, store) = target(target_root.path()).await; 5458 let prepared = store 5459 .prepare_legacy_import(&plan) 5460 .await 5461 .expect("prepared Studio import"); 5462 let classified = store 5463 .classify_legacy_import(&prepared) 5464 .await 5465 .expect("classified Studio import"); 5466 store 5467 .begin_legacy_import(&classified, 13_201) 5468 .await 5469 .expect("begin Studio import"); 5470 5471 let handoff = store 5472 .prepare_legacy_studio_handoff(&classified) 5473 .await 5474 .expect("Studio handoff"); 5475 assert_eq!(handoff.import_id(), plan.import_id()); 5476 assert!(handoff.evidence_path().starts_with(prepared.bundle_path())); 5477 assert_eq!( 5478 handoff.catalog_sha256(), 5479 classified.sources()[0].catalog_sha256() 5480 ); 5481 let host_path = host_root.path().join("studio.sqlite"); 5482 std::fs::copy(handoff.evidence_path(), &host_path).expect("host accepts evidence"); 5483 let (host_length, host_commitment) = file_digest(&host_path).expect("host commitment"); 5484 assert_eq!(host_length, handoff.byte_length()); 5485 assert_eq!(host_commitment, handoff.source_sha256()); 5486 assert!(matches!( 5487 LegacyStudioHandoffReceipt::new(handoff.handoff_sha256(), MemberDigest::new([0; 32])), 5488 Err(Error::InvalidLegacyImportStageRequest) 5489 )); 5490 let receipt = LegacyStudioHandoffReceipt::new(handoff.handoff_sha256(), host_commitment) 5491 .expect("host receipt"); 5492 let acknowledged = store 5493 .acknowledge_legacy_studio_handoff(&classified, receipt, 13_202) 5494 .await 5495 .expect("acknowledge Studio handoff"); 5496 assert_eq!(acknowledged.state(), LegacyImportState::Ready); 5497 assert_eq!( 5498 acknowledged.members()[0].state(), 5499 LegacyImportMemberState::Ready 5500 ); 5501 assert_eq!(acknowledged.members()[0].staged_row_count(), 0); 5502 assert_eq!( 5503 acknowledged.members()[0].resume_cursor().map(<[u8]>::len), 5504 Some(64) 5505 ); 5506 5507 let retry = store 5508 .acknowledge_legacy_studio_handoff(&classified, receipt, 13_202) 5509 .await 5510 .expect("exact receipt retry"); 5511 assert_eq!(retry, acknowledged); 5512 let validation = store 5513 .validate_legacy_import(&classified) 5514 .await 5515 .expect("validate Studio handoff"); 5516 assert_eq!(validation.imported_row_count(), 0); 5517 assert!(!bytes_are_zero(validation.validation_sha256().as_bytes())); 5518 assert_eq!( 5519 store 5520 .validate_legacy_import(&classified) 5521 .await 5522 .expect("repeat Studio validation"), 5523 validation 5524 ); 5525 let conflict = 5526 LegacyStudioHandoffReceipt::new(handoff.handoff_sha256(), MemberDigest::new([99; 32])) 5527 .expect("conflicting host receipt"); 5528 assert!(matches!( 5529 store 5530 .acknowledge_legacy_studio_handoff(&classified, conflict, 13_203) 5531 .await, 5532 Err(Error::LegacyImportConflict) 5533 )); 5534 for pool in [store.pool(), store.private_pool()] { 5535 assert_eq!( 5536 sqlx::query_scalar::<_, i64>( 5537 "SELECT COUNT(*) FROM sqlite_schema WHERE lower(name) LIKE '%studio%'" 5538 ) 5539 .fetch_one(pool) 5540 .await 5541 .expect("Studio schema isolation"), 5542 0 5543 ); 5544 } 5545 assert_eq!( 5546 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_runtime_events") 5547 .fetch_one(store.pool()) 5548 .await 5549 .expect("runtime event isolation"), 5550 0 5551 ); 5552 assert_eq!( 5553 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_private_artifacts") 5554 .fetch_one(store.private_pool()) 5555 .await 5556 .expect("private artifact isolation"), 5557 0 5558 ); 5559 sqlx::query( 5560 "UPDATE radroots_runtime_legacy_import_members 5561 SET staged_row_count = 1, updated_at_ms = 13204 5562 WHERE import_id = ? AND source_kind = 'studio'", 5563 ) 5564 .bind(plan.import_id().as_bytes().as_slice()) 5565 .execute(store.pool()) 5566 .await 5567 .expect("simulate inconsistent Studio count"); 5568 assert!(matches!( 5569 store.validate_legacy_import(&classified).await, 5570 Err(Error::LegacyImportConflict) 5571 )); 5572 studio_connection 5573 .close() 5574 .await 5575 .expect("close Studio source"); 5576 store.close().await.expect("close target"); 5577 } 5578 5579 #[tokio::test] 5580 async fn mixed_source_import_qualifies_backup_resume_validation_and_completion() { 5581 let target_root = tempfile::tempdir().expect("target root"); 5582 let legacy_root = tempfile::tempdir().expect("legacy root"); 5583 let backup_root = tempfile::tempdir().expect("backup root"); 5584 let event_path = legacy_root.path().join("event_store.sqlite"); 5585 let outbox_path = legacy_root.path().join("outbox.sqlite"); 5586 let private_path = legacy_root.path().join("private.sqlite"); 5587 let studio_path = legacy_root.path().join("studio.sqlite"); 5588 let event_connection = 5589 supported_event_database(&event_path, &[signed_event("mixed")]).await; 5590 let outbox_connection = supported_outbox_database(&outbox_path).await; 5591 let private_connection = supported_private_database(&private_path).await; 5592 let studio_connection = supported_studio_database(&studio_path).await; 5593 let plan = LegacyImportPlan::new( 5594 LegacyImportId::new([133; 16]).expect("import id"), 5595 vec![ 5596 LegacySource::new(LegacySourceKind::Studio, &studio_path).expect("Studio source"), 5597 LegacySource::new(LegacySourceKind::Private, &private_path) 5598 .expect("private source"), 5599 LegacySource::new(LegacySourceKind::Outbox, &outbox_path).expect("outbox source"), 5600 LegacySource::new(LegacySourceKind::EventStore, &event_path).expect("event source"), 5601 ], 5602 backup_root.path(), 5603 14_000, 5604 ) 5605 .expect("mixed import plan"); 5606 let (target_paths, store) = target(target_root.path()).await; 5607 let prepared = store 5608 .prepare_legacy_import(&plan) 5609 .await 5610 .expect("prepare mixed import"); 5611 assert_eq!(prepared.snapshots().len(), 4); 5612 let classified = store 5613 .classify_legacy_import(&prepared) 5614 .await 5615 .expect("classify mixed import"); 5616 store 5617 .begin_legacy_import(&classified, 14_001) 5618 .await 5619 .expect("begin mixed import"); 5620 5621 assert!( 5622 store 5623 .stage_legacy_events(&classified, 1, 14_002) 5624 .await 5625 .expect("stage mixed event") 5626 .is_complete() 5627 ); 5628 for page in 0_u64..5 { 5629 let result = store 5630 .stage_legacy_outbox(&classified, 1, 14_003 + page) 5631 .await 5632 .expect("stage mixed outbox"); 5633 assert_eq!(result.is_complete(), page == 4); 5634 } 5635 for page in 0_u64..8 { 5636 let result = store 5637 .stage_legacy_private(&classified, 1, 14_008 + page) 5638 .await 5639 .expect("stage mixed private"); 5640 assert_eq!(result.is_complete(), page == 7); 5641 } 5642 let handoff = store 5643 .prepare_legacy_studio_handoff(&classified) 5644 .await 5645 .expect("prepare mixed Studio handoff"); 5646 let host_receipt = 5647 LegacyStudioHandoffReceipt::new(handoff.handoff_sha256(), MemberDigest::new([88; 32])) 5648 .expect("mixed Studio host receipt"); 5649 store 5650 .acknowledge_legacy_studio_handoff(&classified, host_receipt, 14_016) 5651 .await 5652 .expect("acknowledge mixed Studio handoff"); 5653 let validation = store 5654 .validate_legacy_import(&classified) 5655 .await 5656 .expect("validate mixed import"); 5657 assert_eq!(validation.imported_row_count(), 14); 5658 store.close().await.expect("close before mixed finalize"); 5659 5660 let reopened = 5661 SqliteStorage::open(OpenOptions::new(target_paths, OpenMode::ReadWriteExisting)) 5662 .await 5663 .expect("reopen mixed import"); 5664 let receipt = reopened 5665 .finalize_legacy_import(&classified, validation, 14_017) 5666 .await 5667 .expect("finalize mixed import"); 5668 assert_eq!(receipt.imported_row_count(), 14); 5669 assert_eq!(receipt.validation_sha256(), validation.validation_sha256()); 5670 assert_eq!( 5671 reopened 5672 .finalize_legacy_import(&classified, validation, 14_999) 5673 .await 5674 .expect("retry mixed completion"), 5675 receipt 5676 ); 5677 let journal = reopened 5678 .legacy_import_journal(plan.import_id()) 5679 .await 5680 .expect("mixed journal") 5681 .expect("durable mixed journal"); 5682 assert_eq!(journal.state(), LegacyImportState::Complete); 5683 assert!( 5684 journal 5685 .members() 5686 .iter() 5687 .all(|member| member.state() == LegacyImportMemberState::Complete) 5688 ); 5689 assert_eq!( 5690 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_runtime_events") 5691 .fetch_one(reopened.pool()) 5692 .await 5693 .expect("mixed live event isolation"), 5694 0 5695 ); 5696 assert_eq!( 5697 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM radroots_private_artifacts") 5698 .fetch_one(reopened.private_pool()) 5699 .await 5700 .expect("mixed live private isolation"), 5701 0 5702 ); 5703 for source in plan.sources() { 5704 assert!(source.path().is_file(), "predecessor source was removed"); 5705 } 5706 event_connection.close().await.expect("close event source"); 5707 outbox_connection 5708 .close() 5709 .await 5710 .expect("close outbox source"); 5711 private_connection 5712 .close() 5713 .await 5714 .expect("close private source"); 5715 studio_connection 5716 .close() 5717 .await 5718 .expect("close Studio source"); 5719 reopened.close().await.expect("close mixed target"); 5720 } 5721 5722 #[test] 5723 fn pure_catalog_classification_rejects_mixed_unknown_and_newer_schemas() { 5724 let studio = vec![CatalogRow { 5725 object_type: "table".to_owned(), 5726 name: "sdk_studio_state".to_owned(), 5727 table_name: "sdk_studio_state".to_owned(), 5728 sql: Some( 5729 "CREATE TABLE sdk_studio_state (\n key TEXT PRIMARY KEY NOT NULL,\n value_json TEXT NOT NULL,\n updated_at_ms INTEGER NOT NULL\n)" 5730 .to_owned(), 5731 ), 5732 }]; 5733 assert_eq!( 5734 classify_fixed_catalog( 5735 LegacySourceKind::Studio, 5736 0, 5737 &studio, 5738 0, 5739 STUDIO_CATALOG_SHA256, 5740 ) 5741 .expect("exact Studio schema"), 5742 LegacySchema::StudioV1HostHandoff 5743 ); 5744 assert!(matches!( 5745 classify_fixed_catalog( 5746 LegacySourceKind::Outbox, 5747 0, 5748 &studio, 5749 0, 5750 OUTBOX_CATALOG_SHA256, 5751 ), 5752 Err(Error::UnsupportedLegacySchema { .. }) 5753 )); 5754 let mut unknown = studio.clone(); 5755 unknown.push(CatalogRow { 5756 object_type: "table".to_owned(), 5757 name: "unknown".to_owned(), 5758 table_name: "unknown".to_owned(), 5759 sql: Some("CREATE TABLE unknown(value INTEGER)".to_owned()), 5760 }); 5761 assert!(matches!( 5762 classify_fixed_catalog( 5763 LegacySourceKind::Studio, 5764 0, 5765 &unknown, 5766 0, 5767 STUDIO_CATALOG_SHA256, 5768 ), 5769 Err(Error::UnsupportedLegacySchema { .. }) 5770 )); 5771 assert!(matches!( 5772 classify_fixed_catalog( 5773 LegacySourceKind::Studio, 5774 2, 5775 &studio, 5776 0, 5777 STUDIO_CATALOG_SHA256, 5778 ), 5779 Err(Error::UnsupportedLegacySchema { .. }) 5780 )); 5781 } 5782 5783 #[tokio::test] 5784 async fn event_store_history_requires_the_exact_contiguous_governed_ledger() { 5785 let mut connection = SqliteConnection::connect("sqlite::memory:") 5786 .await 5787 .expect("event-store history database"); 5788 sqlx::query(EVENT_STORE_LEDGER_DDL) 5789 .execute(&mut connection) 5790 .await 5791 .expect("event-store ledger"); 5792 for migration in EVENT_STORE_MIGRATIONS { 5793 sqlx::query( 5794 "INSERT INTO radroots_event_store_schema_migrations(version, name, up_sha256, down_sha256, schema_sha256) VALUES (?, ?, ?, ?, ?)", 5795 ) 5796 .bind(i64::from(migration.version)) 5797 .bind(migration.name) 5798 .bind(migration.up_sha256) 5799 .bind(migration.down_sha256) 5800 .bind(migration.schema_sha256) 5801 .execute(&mut connection) 5802 .await 5803 .expect("event-store history row"); 5804 } 5805 assert_eq!( 5806 validate_event_history(&mut connection) 5807 .await 5808 .expect("exact event-store history"), 5809 4 5810 ); 5811 sqlx::query( 5812 "UPDATE radroots_event_store_schema_migrations SET up_sha256 = ? WHERE version = 4", 5813 ) 5814 .bind("0".repeat(64)) 5815 .execute(&mut connection) 5816 .await 5817 .expect("tamper history"); 5818 assert!(matches!( 5819 validate_event_history(&mut connection).await, 5820 Err(Error::LegacyImportMigrationHistoryInvalid) 5821 )); 5822 connection.close().await.expect("close history database"); 5823 } 5824 5825 #[test] 5826 fn plans_reject_zero_identity_duplicates_missing_paths_and_symlinks() { 5827 let root = tempfile::tempdir().expect("root"); 5828 let backup_root = tempfile::tempdir().expect("backup root"); 5829 let source_path = root.path().join("event.sqlite"); 5830 fs::write(&source_path, b"not yet inspected").expect("source file"); 5831 assert!(matches!( 5832 LegacyImportId::new([0; 16]), 5833 Err(Error::InvalidLegacyImportPlan) 5834 )); 5835 let source = 5836 LegacySource::new(LegacySourceKind::EventStore, &source_path).expect("regular source"); 5837 let import_id = LegacyImportId::new([123; 16]).expect("import id"); 5838 assert!(matches!( 5839 LegacyImportPlan::new(import_id, Vec::new(), backup_root.path(), 12_300), 5840 Err(Error::InvalidLegacyImportPlan) 5841 )); 5842 assert!(matches!( 5843 LegacyImportPlan::new(import_id, vec![source.clone()], backup_root.path(), 0), 5844 Err(Error::InvalidLegacyImportPlan) 5845 )); 5846 assert!(matches!( 5847 LegacyImportPlan::new( 5848 import_id, 5849 vec![source.clone(); LEGACY_SOURCE_MAX + 1], 5850 backup_root.path(), 5851 12_300, 5852 ), 5853 Err(Error::InvalidLegacyImportPlan) 5854 )); 5855 assert!(matches!( 5856 LegacyImportPlan::new( 5857 import_id, 5858 vec![source.clone(), source], 5859 backup_root.path(), 5860 12_300, 5861 ), 5862 Err(Error::InvalidLegacyImportPlan) 5863 )); 5864 assert!(matches!( 5865 LegacySource::new(LegacySourceKind::Outbox, root.path().join("missing.sqlite")), 5866 Err(Error::InvalidLegacySource(_)) 5867 )); 5868 5869 #[cfg(unix)] 5870 { 5871 use std::os::unix::fs::symlink; 5872 let alias = root.path().join("alias.sqlite"); 5873 symlink(&source_path, &alias).expect("source symlink"); 5874 assert!(matches!( 5875 LegacySource::new(LegacySourceKind::Private, alias), 5876 Err(Error::InvalidLegacySource(_)) 5877 )); 5878 } 5879 } 5880 5881 #[test] 5882 fn stage_cursors_reject_every_malformed_boundary() { 5883 assert_eq!( 5884 decode_outbox_stage_cursor(None).expect("initial outbox cursor"), 5885 (LegacyOutboxTable::Operations, 0) 5886 ); 5887 for table in [ 5888 LegacyOutboxTable::Operations, 5889 LegacyOutboxTable::Events, 5890 LegacyOutboxTable::DeliveryPlans, 5891 LegacyOutboxTable::DeliveryTargets, 5892 LegacyOutboxTable::DeliveryAttempts, 5893 ] { 5894 let encoded = encode_outbox_stage_cursor(table, 1); 5895 assert_eq!( 5896 decode_outbox_stage_cursor(Some(&encoded)).expect("outbox cursor"), 5897 (table, 1) 5898 ); 5899 } 5900 for corrupt in [ 5901 Vec::new(), 5902 vec![0; 8], 5903 vec![0; 9], 5904 encode_outbox_stage_cursor(LegacyOutboxTable::Operations, -1).to_vec(), 5905 ] { 5906 assert!(matches!( 5907 decode_outbox_stage_cursor(Some(&corrupt)), 5908 Err(Error::InvalidLegacyImportJournal) 5909 )); 5910 } 5911 5912 assert_eq!( 5913 decode_private_stage_cursor(None).expect("initial private cursor"), 5914 (LegacyPrivateTable::Metadata, String::new()) 5915 ); 5916 for table in [ 5917 LegacyPrivateTable::Metadata, 5918 LegacyPrivateTable::WrappedProfileKeys, 5919 LegacyPrivateTable::SigningSecrets, 5920 LegacyPrivateTable::FarmLocations, 5921 LegacyPrivateTable::TradeArtifacts, 5922 LegacyPrivateTable::CursorKeys, 5923 LegacyPrivateTable::Nip46Sessions, 5924 LegacyPrivateTable::RotationProgress, 5925 ] { 5926 let encoded = encode_private_stage_cursor(table, "cursor"); 5927 assert_eq!( 5928 decode_private_stage_cursor(Some(&encoded)).expect("private cursor"), 5929 (table, "cursor".to_owned()) 5930 ); 5931 assert!(!private_stage_query(table).is_empty()); 5932 } 5933 for corrupt in [Vec::new(), vec![0], vec![1; 1026], vec![1, 0xff]] { 5934 assert!(matches!( 5935 decode_private_stage_cursor(Some(&corrupt)), 5936 Err(Error::InvalidLegacyImportJournal) 5937 )); 5938 } 5939 5940 assert_eq!( 5941 decode_event_stage_cursor(None).expect("initial event cursor"), 5942 0 5943 ); 5944 let event_cursor = encode_event_stage_cursor(1); 5945 assert_eq!( 5946 decode_event_stage_cursor(Some(&event_cursor)).expect("event cursor"), 5947 1 5948 ); 5949 for corrupt in [Vec::new(), vec![0; 8], (-1_i64).to_be_bytes().to_vec()] { 5950 assert!(matches!( 5951 decode_event_stage_cursor(Some(&corrupt)), 5952 Err(Error::InvalidLegacyImportJournal) 5953 )); 5954 } 5955 assert!(matches!( 5956 decode_positive_time(-1), 5957 Err(Error::InvalidLegacyImportJournal) 5958 )); 5959 assert!(matches!( 5960 decode_positive_time(0), 5961 Err(Error::InvalidLegacyImportJournal) 5962 )); 5963 assert_eq!(decode_positive_time(1).expect("positive time"), 1); 5964 } 5965 }