diff --git a/src/callisto/mod.rs b/src/callisto/mod.rs index a87d67b4..ca18ec38 100644 --- a/src/callisto/mod.rs +++ b/src/callisto/mod.rs @@ -66,6 +66,9 @@ pub mod mega_view_root_chain_scan; pub mod mega_webhook; pub mod mega_webhook_delivery; pub mod mega_webhook_event_type; +pub mod mst2_metadata_payload; +pub mod mst2_metadata_prepare; +pub mod mst2_metadata_prepare_page; pub mod mst2_native_head; pub mod mst2_native_publication; pub mod mst2_publication; diff --git a/src/callisto/mst2_metadata_payload.rs b/src/callisto/mst2_metadata_payload.rs new file mode 100644 index 00000000..798aab27 --- /dev/null +++ b/src/callisto/mst2_metadata_payload.rs @@ -0,0 +1,16 @@ +use sea_orm::entity::prelude::*; + +#[derive(Clone, Debug, PartialEq, DeriveEntityModel, Eq)] +#[sea_orm(table_name = "mst2_metadata_payload")] +pub struct Model { + #[sea_orm(primary_key, auto_increment = false)] + pub page_id: Vec, + pub metadata_codec: i16, + pub byte_size: i32, + pub payload: Vec, + pub created_at: DateTimeWithTimeZone, +} + +#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)] +pub enum Relation {} +impl ActiveModelBehavior for ActiveModel {} diff --git a/src/callisto/mst2_metadata_prepare.rs b/src/callisto/mst2_metadata_prepare.rs new file mode 100644 index 00000000..9401282b --- /dev/null +++ b/src/callisto/mst2_metadata_prepare.rs @@ -0,0 +1,32 @@ +use sea_orm::entity::prelude::*; + +#[derive(Clone, Debug, PartialEq, DeriveEntityModel, Eq)] +#[sea_orm(table_name = "mst2_metadata_prepare")] +pub struct Model { + #[sea_orm(primary_key, auto_increment = false)] + pub prepare_id: String, + pub operation_id: String, + pub manifest_digest: Vec, + pub canonical_plan: Vec, + pub source_domain: String, + pub tagged_root_tree_oid: String, + pub scope: String, + pub schema_version: i16, + pub metadata_codec: i16, + pub materialization_policy: i16, + pub fs_semantics: i16, + pub access_projection: i16, + pub verification_revision: i32, + pub projection_revision: i16, + pub metadata_root: Vec, + pub node_count: i32, + pub edge_count: i32, + pub total_bytes: i64, + pub state: String, + pub created_at: DateTimeWithTimeZone, + pub committed_at: Option, +} + +#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)] +pub enum Relation {} +impl ActiveModelBehavior for ActiveModel {} diff --git a/src/callisto/mst2_metadata_prepare_page.rs b/src/callisto/mst2_metadata_prepare_page.rs new file mode 100644 index 00000000..61299f7e --- /dev/null +++ b/src/callisto/mst2_metadata_prepare_page.rs @@ -0,0 +1,15 @@ +use sea_orm::entity::prelude::*; + +#[derive(Clone, Debug, PartialEq, DeriveEntityModel, Eq)] +#[sea_orm(table_name = "mst2_metadata_prepare_page")] +pub struct Model { + #[sea_orm(primary_key, auto_increment = false)] + pub prepare_id: String, + #[sea_orm(primary_key, auto_increment = false)] + pub page_id: Vec, + pub expected_size: i32, +} + +#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)] +pub enum Relation {} +impl ActiveModelBehavior for ActiveModel {} diff --git a/src/ceres/snapshot/metadata_install.rs b/src/ceres/snapshot/metadata_install.rs new file mode 100644 index 00000000..ab946290 --- /dev/null +++ b/src/ceres/snapshot/metadata_install.rs @@ -0,0 +1,463 @@ +//! Storage-owned identity for a durable native metadata installation. +//! This plan grants no publication, lease or content-retention authority. + +use std::collections::{BTreeMap, BTreeSet}; + +use mst2_codec::metapage::{HEADER_LEN, PAGE_MAX_BYTES}; +use sha2::{Digest, Sha256}; + +use super::{ + error::{SnapshotError, SnapshotErrorCode}, + retention_dag::{MetadataDagLimits, MetadataPageId, ValidatedMetadataDag}, +}; + +const DOMAIN: &[u8] = b"mega.mst2.metadata-install.v1\0"; +pub(crate) const MAX_PLAN_BYTES: usize = 2 * 1024 * 1024; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct MetadataInstallIdentity { + pub source_domain: String, + pub tagged_root_tree_oid: String, + pub scope: String, + pub schema_version: u16, + pub metadata_codec: u16, + pub materialization_policy: u16, + pub fs_semantics: u16, + pub access_projection: u16, + pub verification_revision: i32, + pub projection_revision: u16, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct MetadataInstallPlan { + pub identity: MetadataInstallIdentity, + pub root: MetadataPageId, + pub pages: BTreeMap, + pub edges: BTreeSet<(MetadataPageId, MetadataPageId)>, + pub total_bytes: u64, +} + +impl MetadataInstallPlan { + pub(super) fn from_validated( + identity: MetadataInstallIdentity, + dag: &ValidatedMetadataDag, + ) -> Result { + dag.check_limits(MetadataDagLimits::default())?; + let edges = dag + .edges() + .iter() + .map(|edge| Ok((parse_node_id(&edge.parent)?, parse_node_id(&edge.child)?))) + .collect::>()?; + let plan = Self { + identity, + root: dag.root(), + pages: dag + .payloads() + .iter() + .map(|page| (page.id, page.size)) + .collect(), + edges, + total_bytes: dag.payload_bytes(), + }; + plan.validate()?; + Ok(plan) + } + + pub fn digest(&self) -> Result<[u8; 32], SnapshotError> { + Ok(Sha256::digest(self.encode()?).into()) + } + + pub fn encode(&self) -> Result, SnapshotError> { + self.validate()?; + let mut bytes = DOMAIN.to_vec(); + bytes.extend_from_slice(&1u16.to_be_bytes()); + for value in [ + &self.identity.source_domain, + &self.identity.tagged_root_tree_oid, + &self.identity.scope, + ] { + bytes.extend_from_slice(&(value.len() as u32).to_be_bytes()); + bytes.extend_from_slice(value.as_bytes()); + } + for value in [ + self.identity.schema_version, + self.identity.metadata_codec, + self.identity.materialization_policy, + self.identity.fs_semantics, + self.identity.access_projection, + ] { + bytes.extend_from_slice(&value.to_be_bytes()); + } + bytes.extend_from_slice(&self.identity.verification_revision.to_be_bytes()); + bytes.extend_from_slice(&self.identity.projection_revision.to_be_bytes()); + bytes.extend_from_slice(&self.root); + bytes.extend_from_slice(&(self.pages.len() as u32).to_be_bytes()); + for (id, size) in &self.pages { + bytes.extend_from_slice(id); + bytes.extend_from_slice(&size.to_be_bytes()); + } + bytes.extend_from_slice(&(self.edges.len() as u32).to_be_bytes()); + for (parent, child) in &self.edges { + bytes.extend_from_slice(parent); + bytes.extend_from_slice(child); + } + if bytes.len() > MAX_PLAN_BYTES { + return Err(limit("metadata installation plan exceeds its byte budget")); + } + Ok(bytes) + } + + pub fn decode(bytes: &[u8], expected_digest: &[u8; 32]) -> Result { + if bytes.len() > MAX_PLAN_BYTES { + return Err(limit( + "stored metadata installation plan exceeds its byte budget", + )); + } + if Sha256::digest(bytes).as_slice() != expected_digest { + return Err(integrity( + "stored metadata installation plan digest mismatch", + )); + } + let mut reader = PlanReader(bytes); + if reader.take(DOMAIN.len())? != DOMAIN || reader.u16()? != 1 { + return Err(integrity("unsupported metadata installation plan encoding")); + } + let identity = MetadataInstallIdentity { + source_domain: reader.string(64)?, + tagged_root_tree_oid: reader.string(128)?, + scope: reader.string(4096)?, + schema_version: reader.u16()?, + metadata_codec: reader.u16()?, + materialization_policy: reader.u16()?, + fs_semantics: reader.u16()?, + access_projection: reader.u16()?, + verification_revision: i32::from_be_bytes(reader.array()?), + projection_revision: reader.u16()?, + }; + let root = reader.array()?; + let count = reader.count(MetadataDagLimits::default().nodes)?; + let mut pages = BTreeMap::new(); + let mut total_bytes = 0u64; + let mut last = None; + for _ in 0..count { + let id = reader.array()?; + let size = reader.u64()?; + if last.is_some_and(|previous| previous >= id) { + return Err(integrity( + "metadata installation pages are not uniquely ordered", + )); + } + last = Some(id); + total_bytes = total_bytes + .checked_add(size) + .ok_or_else(|| limit("metadata installation byte count overflow"))?; + pages.insert(id, size); + } + let count = reader.count(MetadataDagLimits::default().edges)?; + let mut edges = BTreeSet::new(); + let mut last = None; + for _ in 0..count { + let edge = (reader.array()?, reader.array()?); + if last.is_some_and(|previous| previous >= edge) { + return Err(integrity( + "metadata installation edges are not uniquely ordered", + )); + } + last = Some(edge); + edges.insert(edge); + } + if !reader.0.is_empty() { + return Err(integrity("metadata installation plan has trailing bytes")); + } + let plan = Self { + identity, + root, + pages, + edges, + total_bytes, + }; + plan.validate()?; + if plan.encode()?.as_slice() != bytes { + return Err(integrity("metadata installation plan is not canonical")); + } + Ok(plan) + } + + fn validate(&self) -> Result<(), SnapshotError> { + let identity = &self.identity; + if identity.source_domain != "native-git" + || identity.schema_version != mst2_codec::descriptor::SCHEMA_VERSION + || identity.metadata_codec != mst2_codec::descriptor::METADATA_CODEC + || identity.materialization_policy + != mst2_codec::descriptor::MATERIALIZATION_POLICY_GIT_RAW_V1 + || identity.fs_semantics != mst2_codec::descriptor::FS_SEMANTICS_LINUX_CODE_V1 + || identity.access_projection != mst2_codec::descriptor::ACCESS_PROJECTION_EXACT_FULL + || identity.verification_revision + != crate::jupiter::storage::mono_storage::MST2_VERIFICATION_VERSION + || identity.projection_revision != 1 + { + return Err(integrity( + "unsupported stored native metadata installation profile", + )); + } + let tagged = identity.tagged_root_tree_oid.as_str(); + let (kind, hex) = tagged + .split_once(':') + .ok_or_else(|| integrity("untagged fixed root tree"))?; + let length = match kind { + "sha1" => 40, + "sha256" | "blake3" => 64, + _ => return Err(integrity("unsupported fixed root tree hash kind")), + }; + if hex.len() != length + || !hex + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) + { + return Err(integrity("noncanonical fixed root tree identity")); + } + super::view::validate_scope_relative_path(&identity.scope) + .map_err(|_| integrity("noncanonical metadata installation scope"))?; + let limits = MetadataDagLimits::default(); + if self.pages.is_empty() + || self.pages.len() > limits.nodes + || self.edges.len() > limits.edges + || self.total_bytes > limits.payload_bytes + { + return Err(limit("metadata installation group exceeds its budget")); + } + if !self.pages.contains_key(&self.root) + || self + .pages + .values() + .any(|size| *size < HEADER_LEN as u64 || *size > PAGE_MAX_BYTES as u64) + || self + .pages + .values() + .try_fold(0u64, |total, size| total.checked_add(*size)) + != Some(self.total_bytes) + { + return Err(integrity( + "metadata installation root or page size is invalid", + )); + } + let mut incoming: BTreeMap<_, usize> = self.pages.keys().map(|id| (*id, 0)).collect(); + let mut children: BTreeMap<_, Vec<_>> = BTreeMap::new(); + for (parent, child) in &self.edges { + if parent == child || !self.pages.contains_key(parent) { + return Err(integrity("invalid metadata installation edge")); + } + *incoming + .get_mut(child) + .ok_or_else(|| integrity("metadata installation child is missing"))? += 1; + children.entry(*parent).or_default().push(*child); + } + let mut ready: Vec<_> = incoming + .iter() + .filter_map(|(id, count)| (*count == 0).then_some(*id)) + .collect(); + let mut visited = 0; + while let Some(parent) = ready.pop() { + visited += 1; + for child in children.get(&parent).into_iter().flatten() { + let count = incoming + .get_mut(child) + .ok_or_else(|| integrity("missing installation child"))?; + *count -= 1; + if *count == 0 { + ready.push(*child); + } + } + } + if visited != self.pages.len() { + return Err(integrity("cyclic metadata installation plan")); + } + let mut reachable = BTreeSet::new(); + let mut pending = vec![self.root]; + while let Some(id) = pending.pop() { + if reachable.insert(id) { + pending.extend(children.get(&id).into_iter().flatten()); + } + } + if reachable.len() != self.pages.len() { + return Err(integrity("unreachable metadata installation pages")); + } + Ok(()) + } +} + +struct PlanReader<'a>(&'a [u8]); + +impl<'a> PlanReader<'a> { + fn take(&mut self, count: usize) -> Result<&'a [u8], SnapshotError> { + if count > self.0.len() { + return Err(integrity("truncated metadata installation plan")); + } + let (bytes, rest) = self.0.split_at(count); + self.0 = rest; + Ok(bytes) + } + fn array(&mut self) -> Result<[u8; N], SnapshotError> { + self.take(N)? + .try_into() + .map_err(|_| integrity("invalid installation field")) + } + fn u16(&mut self) -> Result { + Ok(u16::from_be_bytes(self.array()?)) + } + fn u32(&mut self) -> Result { + Ok(u32::from_be_bytes(self.array()?)) + } + fn u64(&mut self) -> Result { + Ok(u64::from_be_bytes(self.array()?)) + } + fn count(&mut self, maximum: usize) -> Result { + let count = self.u32()? as usize; + if count > maximum { + return Err(limit("metadata installation count exceeds its budget")); + } + Ok(count) + } + fn string(&mut self, maximum: usize) -> Result { + let count = self.count(maximum)?; + String::from_utf8(self.take(count)?.to_vec()) + .map_err(|_| integrity("non-UTF8 installation identity")) + } +} + +fn parse_node_id(value: &str) -> Result { + let value = value + .strip_prefix("page:sha256:") + .ok_or_else(|| integrity("non-page installation node"))?; + hex::decode(value) + .map_err(|_| integrity("invalid installation page ID"))? + .try_into() + .map_err(|_| integrity("invalid installation page ID length")) +} + +fn integrity(message: &str) -> SnapshotError { + SnapshotError::new(SnapshotErrorCode::IntegrityError, message) +} +fn limit(message: &str) -> SnapshotError { + SnapshotError::new(SnapshotErrorCode::LimitExceeded, message) +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use mst2_codec::metapage::{Page, page_id}; + + use super::*; + use crate::ceres::snapshot::{ + pages::PreparedNativeMetadataRetention, retention_dag::MetadataDagBuilder, + }; + + fn plan() -> MetadataInstallPlan { + let bytes = Page::build(&[]).unwrap(); + let mut builder = MetadataDagBuilder::new(MetadataDagLimits::default()); + builder.add_directory(&bytes, &[]).unwrap(); + PreparedNativeMetadataRetention::test_installation( + Arc::new(builder.finish(page_id(&bytes)).unwrap()), + "/", + ) + .install_plan() + .unwrap() + } + + #[test] + fn native_installation_plan_round_trip_binds_every_source_and_profile_field() { + let plan = plan(); + let bytes = plan.encode().unwrap(); + let digest = plan.digest().unwrap(); + assert_eq!(MetadataInstallPlan::decode(&bytes, &digest).unwrap(), plan); + let mut changed = plan.clone(); + changed.identity.tagged_root_tree_oid = format!("sha1:{}", "b".repeat(40)); + assert_ne!(changed.digest().unwrap(), digest); + changed = plan.clone(); + changed.identity.scope = "/other".into(); + assert_ne!(changed.digest().unwrap(), digest); + changed = plan; + changed.identity.projection_revision += 1; + assert_eq!( + changed.encode().unwrap_err().code, + SnapshotErrorCode::IntegrityError + ); + } + + #[test] + fn native_installation_stored_plan_rejects_bad_digest_truncation_trailing_bytes_and_profile() { + let plan = plan(); + let bytes = plan.encode().unwrap(); + let digest = plan.digest().unwrap(); + let mut bad_digest = digest; + bad_digest[0] ^= 1; + assert!(MetadataInstallPlan::decode(&bytes, &bad_digest).is_err()); + for length in [0, DOMAIN.len(), bytes.len() - 1] { + let truncated = &bytes[..length]; + let digest = Sha256::digest(truncated).into(); + assert!(MetadataInstallPlan::decode(truncated, &digest).is_err()); + } + let mut trailing = bytes.clone(); + trailing.push(0); + assert!(MetadataInstallPlan::decode(&trailing, &Sha256::digest(&trailing).into()).is_err()); + let mut wrong_domain = bytes; + wrong_domain[0] ^= 1; + assert!( + MetadataInstallPlan::decode(&wrong_domain, &Sha256::digest(&wrong_domain).into()) + .is_err() + ); + } + + #[test] + fn native_installation_plan_rejects_oversized_stored_bytes_before_allocating_fields() { + let bytes = vec![0; MAX_PLAN_BYTES + 1]; + assert_eq!( + MetadataInstallPlan::decode(&bytes, &[0; 32]) + .unwrap_err() + .code, + SnapshotErrorCode::LimitExceeded + ); + } + + #[test] + fn native_installation_posix_scopes_preserve_utf8_backslash_control_and_path_boundaries() { + let original = plan(); + let at_byte_limit = format!("/{}", vec!["a".repeat(255); 16].join("/")); + assert_eq!(at_byte_limit.len(), 4096); + let at_component_limit = format!("/{}", vec!["a"; 256].join("/")); + for scope in [ + "/a\\b".to_owned(), + "/目录/é".to_owned(), + "/line\ncontrol\t".to_owned(), + format!("/{}", "a".repeat(255)), + at_byte_limit.clone(), + at_component_limit, + ] { + let mut value = original.clone(); + value.identity.scope = scope.clone(); + let encoded = value.encode().unwrap(); + assert_eq!( + MetadataInstallPlan::decode(&encoded, &value.digest().unwrap()) + .unwrap() + .identity + .scope, + scope + ); + } + for scope in [ + format!("/{}", "a".repeat(256)), + format!("/{}", vec!["a"; 257].join("/")), + format!("{at_byte_limit}/a"), + "/a\0b".to_owned(), + "/a/..".to_owned(), + ] { + let mut value = original.clone(); + value.identity.scope = scope; + assert_eq!( + value.encode().unwrap_err().code, + SnapshotErrorCode::IntegrityError + ); + } + } +} diff --git a/src/ceres/snapshot/mod.rs b/src/ceres/snapshot/mod.rs index 8dea3c34..97a1a50b 100644 --- a/src/ceres/snapshot/mod.rs +++ b/src/ceres/snapshot/mod.rs @@ -7,6 +7,7 @@ pub mod chunks; pub mod descriptor; pub mod error; pub mod frame_stream; +pub(crate) mod metadata_install; pub mod namespace; pub mod pages; pub(crate) mod projection_observation; diff --git a/src/ceres/snapshot/pages.rs b/src/ceres/snapshot/pages.rs index fe7695b2..05e4e67f 100644 --- a/src/ceres/snapshot/pages.rs +++ b/src/ceres/snapshot/pages.rs @@ -259,6 +259,40 @@ pub struct PreparedNativeMetadataRetention { } impl PreparedNativeMetadataRetention { + #[cfg(test)] + pub(crate) fn test_installation(dag: Arc, scope: &str) -> Self { + let tree_oid = + ObjectHash::from_hex_for_kind(git_internal::hash::HashKind::Sha1, &"a".repeat(40)) + .unwrap(); + Self { + key: NativeRetentionKey { + projection: NativeProjectionKey::new(tree_oid), + scope: scope.to_owned(), + }, + dag, + } + } + pub(crate) fn install_plan( + &self, + ) -> Result { + use super::metadata_install::{MetadataInstallIdentity, MetadataInstallPlan}; + let key = &self.key.projection; + MetadataInstallPlan::from_validated( + MetadataInstallIdentity { + source_domain: key.source_domain.to_owned(), + tagged_root_tree_oid: key.tree_oid.clone(), + scope: self.key.scope.clone(), + schema_version: key.schema_version, + metadata_codec: key.metadata_codec, + materialization_policy: key.materialization_policy, + fs_semantics: key.fs_semantics, + access_projection: key.access_projection, + verification_revision: key.verification_revision, + projection_revision: key.projection_revision, + }, + &self.dag, + ) + } pub fn fixed_root_tree_oid(&self) -> &str { &self.key.projection.tree_oid } diff --git a/src/commands/mod.rs b/src/commands/mod.rs index 92361218..4c7fc6d3 100644 --- a/src/commands/mod.rs +++ b/src/commands/mod.rs @@ -119,7 +119,7 @@ pub(crate) fn builtin_exec(cmd: &str) -> Option { pub(crate) fn load_mode(cmd: &str, args: &ArgMatches) -> Option { match cmd { "service" => match args.subcommand_name() { - Some("init") => Some(LoadMode::ParsedExistingConfig), + Some("init" | "native-publication-init") => Some(LoadMode::ParsedExistingConfig), _ => Some(LoadMode::FullAppContext), }, "debug" => Some(LoadMode::FullAppContext), diff --git a/src/commands/service/mod.rs b/src/commands/service/mod.rs index 939674c3..204fc9ce 100644 --- a/src/commands/service/mod.rs +++ b/src/commands/service/mod.rs @@ -17,12 +17,19 @@ use crate::{ pub mod http; pub mod init; pub mod multi; +pub mod native_publication_init; pub mod ssh; const CONFIG_RELOAD_POLL_INTERVAL: Duration = Duration::from_secs(5); pub fn cli() -> Command { - let subcommands = vec![init::cli(), http::cli(), ssh::cli(), multi::cli()]; + let subcommands = vec![ + init::cli(), + native_publication_init::cli(), + http::cli(), + ssh::cli(), + multi::cli(), + ]; Command::new("service") .about("Start different kinds of server: for example https or ssh") .subcommands(subcommands) @@ -30,6 +37,9 @@ pub fn cli() -> Command { #[tokio::main] pub(crate) async fn exec(ctx: CommandContext, args: &ArgMatches) -> MegaResult { + if let Some(("native-publication-init", subcommand_args)) = args.subcommand() { + return native_publication_init::exec(ctx, subcommand_args).await; + } let config_path = ctx.config_path.clone(); let config_profile_path = ctx.config_profile_path.clone(); let config = require_config(ctx, "service")?; @@ -265,7 +275,10 @@ mod tests { .map(|cmd| cmd.get_name().to_owned()) .collect::>(); - assert_eq!(names, vec!["init", "http", "ssh", "multi"]); + assert_eq!( + names, + vec!["init", "native-publication-init", "http", "ssh", "multi"] + ); } #[test] diff --git a/src/commands/service/native_publication_init.rs b/src/commands/service/native_publication_init.rs new file mode 100644 index 00000000..483f647c --- /dev/null +++ b/src/commands/service/native_publication_init.rs @@ -0,0 +1,198 @@ +use std::sync::Arc; + +use clap::{Arg, ArgAction, ArgMatches, Command}; +use sea_orm_migration::MigratorTrait; + +use crate::{ + commands::{CommandContext, require_config}, + common::errors::{MegaError, MegaResult}, + config::{DbConfig, loader::ConfigSource}, + jupiter::{ + migration::Migrator, + storage::{ + base_storage::{BaseStorage, StorageConnector}, + init::{postgres_connection, read_only_database_connection}, + mono_storage::MonoStorage, + native_publication_storage::NativeRoot, + }, + }, +}; + +#[path = "native_publication_init_preflight.rs"] +mod preflight; +use preflight::InitializationTarget; + +pub fn cli() -> Command { + let mut command = Command::new("native-publication-init") + .about("Prepare native publication in a stopped, paused and drained SHA-1 deployment"); + for (name, help) in [ + ("instance", "Deployment UUID; must match mst2.instance_uuid"), + ( + "expected-root-commit", + "Exact current native root commit object ID", + ), + ( + "expected-root-tree", + "Exact current native root tree object ID", + ), + ] { + command = command.arg(Arg::new(name).long(name).required(true).help(help)); + } + for (name, help) in [ + ("yes", "Confirm preparation of the native publication head"), + ( + "writers-stopped", + "Confirm all old writer processes have been stopped", + ), + ] { + command = command.arg( + Arg::new(name) + .long(name) + .action(ArgAction::SetTrue) + .required(true) + .help(help), + ); + } + command +} + +pub(crate) async fn exec(ctx: CommandContext, args: &ArgMatches) -> MegaResult { + if !ctx + .config_summary + .as_ref() + .is_some_and(|summary| matches!(summary.source, ConfigSource::Cli | ConfigSource::Env)) + { + return Err(MegaError::Other( + "MST2_NATIVE_INIT_CONFIG_REQUIRED: name the deployment with --config or MEGA_CONFIG" + .into(), + )); + } + if !args.get_flag("yes") || !args.get_flag("writers-stopped") { + return Err(MegaError::Other( + "MST2_NATIVE_INIT_CONFIRMATION_REQUIRED: --yes and --writers-stopped are required" + .into(), + )); + } + let config = require_config(ctx, "service native-publication-init")?; + config.validate()?; + if !config.mst2.enabled || !config.mst2.publication_enabled { + return Err(MegaError::Other( + "MST2_NATIVE_INIT_DISABLED: mst2.enabled and mst2.publication_enabled must be true" + .into(), + )); + } + let argument = |name| { + args.get_one::(name) + .map(String::as_str) + .ok_or_else(|| MegaError::Other(format!("--{name} is required"))) + }; + let target = InitializationTarget::parse( + config.monorepo.object_format.as_str(), + config.mst2.instance_uuid.as_deref(), + argument("instance")?, + argument("expected-root-commit")?, + argument("expected-root-tree")?, + ) + .map_err(|error| MegaError::Other(error.to_string()))?; + ensure_schema_current(&config.database).await?; + let db = Arc::new(postgres_connection(&config.database).await?); + let mono = MonoStorage { + base: BaseStorage::new(db.clone()), + }; + let result = mono + .initialize_native_publication_for_maintenance( + &target.instance, + &NativeRoot { + commit: target.commit, + tree: target.tree, + }, + ) + .await + .map_err(|error| MegaError::Other(error.to_string())); + let close = db.close_by_ref().await; + result?; + close?; + tracing::info!( + instance = %target.instance, + state = "INITIALIZING", + "Native publication initialization prepared; queue remains paused" + ); + Ok(()) +} + +async fn ensure_schema_current(config: &DbConfig) -> MegaResult { + let db = read_only_database_connection(config).await?; + let pending = Migrator::get_pending_migrations_read_only(&db).await; + let _ = db.close().await; + let pending = pending.map_err(|error| { + MegaError::Other(format!( + "MST2_NATIVE_INIT_SCHEMA_MISMATCH: cannot confirm the current schema: {error}" + )) + })?; + if !pending.is_empty() { + return Err(MegaError::Other(format!( + "MST2_NATIVE_INIT_SCHEMA_MISMATCH: {} migrations are pending; this command does not migrate", + pending.len() + ))); + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::commands::{LoadMode, load_mode}; + + #[test] + fn maintenance_cli_requires_each_confirmation_and_fixed_target() { + let arguments = [ + "native-publication-init", + "--yes", + "--writers-stopped", + "--instance", + "instance", + "--expected-root-commit", + "commit", + "--expected-root-tree", + "tree", + ]; + assert!(cli().try_get_matches_from(arguments).is_ok()); + for remove in [1, 2, 3, 5, 7] { + let count = if remove <= 2 { 1 } else { 2 }; + let mut incomplete = arguments.to_vec(); + incomplete.drain(remove..remove + count); + assert!(cli().try_get_matches_from(incomplete).is_err()); + } + let mut service_arguments = vec!["service"]; + service_arguments.extend(arguments); + let args = super::super::cli() + .try_get_matches_from(service_arguments) + .unwrap(); + assert_eq!( + load_mode("service", &args), + Some(LoadMode::ParsedExistingConfig) + ); + } + + #[tokio::test] + async fn maintenance_command_refuses_unnamed_config_before_any_database_io() { + let args = cli() + .try_get_matches_from([ + "native-publication-init", + "--yes", + "--writers-stopped", + "--instance", + "instance", + "--expected-root-commit", + "commit", + "--expected-root-tree", + "tree", + ]) + .unwrap(); + let error = exec(CommandContext::default(), &args).await.unwrap_err(); + assert!( + matches!(&error, MegaError::Other(message) if message.starts_with("MST2_NATIVE_INIT_CONFIG_REQUIRED")), + "unexpected error: {error}" + ); + } +} diff --git a/src/commands/service/native_publication_init_preflight.rs b/src/commands/service/native_publication_init_preflight.rs new file mode 100644 index 00000000..360751b6 --- /dev/null +++ b/src/commands/service/native_publication_init_preflight.rs @@ -0,0 +1,162 @@ +#[derive(Debug, thiserror::Error, PartialEq, Eq)] +pub(super) enum InitializationTargetError { + #[error( + "MST2_NATIVE_INIT_UNSUPPORTED_OBJECT_FORMAT: {0}; this maintenance entry supports sha1 only" + )] + UnsupportedObjectFormat(String), + #[error("MST2_NATIVE_INIT_INVALID_INSTANCE: both deployment instances must be non-nil UUIDs")] + InvalidInstance, + #[error("MST2_NATIVE_INIT_INSTANCE_MISMATCH: --instance differs from mst2.instance_uuid")] + InstanceMismatch, + #[error( + "MST2_NATIVE_INIT_INVALID_ROOT: expected commit and tree must be canonical SHA-1 object IDs" + )] + InvalidRoot, +} + +#[derive(Debug, PartialEq, Eq)] +pub(super) struct InitializationTarget { + pub(super) instance: String, + pub(super) commit: String, + pub(super) tree: String, +} + +impl InitializationTarget { + pub(super) fn parse( + object_format: &str, + configured_instance: Option<&str>, + requested_instance: &str, + commit: &str, + tree: &str, + ) -> Result { + if object_format != "sha1" { + return Err(InitializationTargetError::UnsupportedObjectFormat( + object_format.into(), + )); + } + let parse_instance = |value: &str| { + uuid::Uuid::parse_str(value) + .ok() + .filter(|id| !id.is_nil()) + .ok_or(InitializationTargetError::InvalidInstance) + }; + let configured = + parse_instance(configured_instance.ok_or(InitializationTargetError::InvalidInstance)?)?; + let requested = parse_instance(requested_instance)?; + if configured != requested { + return Err(InitializationTargetError::InstanceMismatch); + } + if [commit, tree].iter().any(|value| { + value.len() != 40 + || !value + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) + }) { + return Err(InitializationTargetError::InvalidRoot); + } + Ok(Self { + instance: configured.to_string(), + commit: commit.into(), + tree: tree.into(), + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const INSTANCE: &str = "6ab219b0-4275-45ba-9d7b-7b0b633018cd"; + + #[test] + fn target_is_canonical_and_bound_to_the_configured_instance() { + let target = InitializationTarget::parse( + "sha1", + Some(INSTANCE), + &INSTANCE.to_uppercase(), + &"a".repeat(40), + &"b".repeat(40), + ) + .unwrap(); + assert_eq!(target.instance, INSTANCE); + assert_eq!(target.commit, "a".repeat(40)); + assert_eq!(target.tree, "b".repeat(40)); + } + + #[test] + fn target_rejects_other_formats_without_inferring_from_equal_width() { + for format in ["sha256", "blake3", "unknown"] { + assert_eq!( + InitializationTarget::parse( + format, + Some(INSTANCE), + INSTANCE, + &"a".repeat(64), + &"b".repeat(64), + ), + Err(InitializationTargetError::UnsupportedObjectFormat( + format.into() + )) + ); + } + } + + #[test] + fn target_rejects_nil_missing_invalid_or_different_instances() { + for instance in [ + None, + Some("invalid"), + Some("00000000-0000-0000-0000-000000000000"), + ] { + assert_eq!( + InitializationTarget::parse( + "sha1", + instance, + INSTANCE, + &"a".repeat(40), + &"b".repeat(40) + ), + Err(InitializationTargetError::InvalidInstance) + ); + } + for instance in ["invalid", "00000000-0000-0000-0000-000000000000"] { + assert_eq!( + InitializationTarget::parse( + "sha1", + Some(INSTANCE), + instance, + &"a".repeat(40), + &"b".repeat(40) + ), + Err(InitializationTargetError::InvalidInstance) + ); + } + assert_eq!( + InitializationTarget::parse( + "sha1", + Some(INSTANCE), + "11111111-2222-4333-8444-555555555555", + &"a".repeat(40), + &"b".repeat(40) + ), + Err(InitializationTargetError::InstanceMismatch) + ); + } + + #[test] + fn target_rejects_noncanonical_or_wrong_width_object_ids() { + for value in [ + "a".repeat(39), + "a".repeat(64), + "A".repeat(40), + "z".repeat(40), + ] { + for (commit, tree) in [(&value, &"b".repeat(40)), (&"a".repeat(40), &value)] { + assert_eq!( + InitializationTarget::parse("sha1", Some(INSTANCE), INSTANCE, commit, tree), + Err(InitializationTargetError::InvalidRoot) + ); + } + } + } +} diff --git a/src/jupiter/migration/m20261005_000300_add_mst2_metadata_install.rs b/src/jupiter/migration/m20261005_000300_add_mst2_metadata_install.rs new file mode 100644 index 00000000..d0e8d168 --- /dev/null +++ b/src/jupiter/migration/m20261005_000300_add_mst2_metadata_install.rs @@ -0,0 +1,76 @@ +//! Additive, forward-only native metadata installation records. + +use mst2_codec::metapage::{HEADER_LEN, PAGE_MAX_BYTES}; +use sea_orm_migration::prelude::*; + +#[derive(DeriveMigrationName)] +pub struct Migration; + +#[async_trait::async_trait] +impl MigrationTrait for Migration { + async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> { + manager.get_connection().execute_unprepared(&format!( + "CREATE TABLE mst2_metadata_storage_scope ( + singleton smallint PRIMARY KEY CHECK (singleton = 1), + storage_uuid text NOT NULL UNIQUE + ); + CREATE TABLE mst2_metadata_payload ( + page_id bytea PRIMARY KEY CHECK (octet_length(page_id) = 32), + metadata_codec smallint NOT NULL CHECK (metadata_codec = 1), + byte_size integer NOT NULL CHECK (byte_size BETWEEN {HEADER_LEN} AND {PAGE_MAX_BYTES}), + payload bytea NOT NULL CHECK (octet_length(payload) = byte_size), + created_at timestamptz NOT NULL DEFAULT now() + ); + CREATE OR REPLACE FUNCTION mst2_metadata_payload_immutable() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN RAISE EXCEPTION 'MST2 metadata payloads cannot be updated or deleted'; END $$; + CREATE TRIGGER mst2_metadata_payload_immutable BEFORE UPDATE OR DELETE + ON mst2_metadata_payload FOR EACH ROW EXECUTE FUNCTION mst2_metadata_payload_immutable(); + CREATE TRIGGER mst2_metadata_scope_immutable BEFORE UPDATE OR DELETE + ON mst2_metadata_storage_scope FOR EACH ROW EXECUTE FUNCTION mst2_metadata_payload_immutable(); + CREATE TABLE mst2_metadata_prepare ( + prepare_id text PRIMARY KEY CHECK (prepare_id ~ '^[0-9a-f]{{8}}-[0-9a-f]{{4}}-4[0-9a-f]{{3}}-[89ab][0-9a-f]{{3}}-[0-9a-f]{{12}}$'), + operation_id text NOT NULL UNIQUE CHECK (octet_length(operation_id) BETWEEN 1 AND 255), + manifest_digest bytea NOT NULL CHECK (octet_length(manifest_digest) = 32), + canonical_plan bytea NOT NULL CHECK (octet_length(canonical_plan) <= 2097152), + source_domain text NOT NULL CHECK (source_domain = 'native-git'), + tagged_root_tree_oid text NOT NULL, + scope text NOT NULL CHECK (octet_length(scope) <= 4096), + schema_version smallint NOT NULL, + metadata_codec smallint NOT NULL CHECK (metadata_codec = 1), + materialization_policy smallint NOT NULL, + fs_semantics smallint NOT NULL, + access_projection smallint NOT NULL, + verification_revision integer NOT NULL, + projection_revision smallint NOT NULL, + metadata_root bytea NOT NULL CHECK (octet_length(metadata_root) = 32), + node_count integer NOT NULL CHECK (node_count BETWEEN 1 AND 4096), + edge_count integer NOT NULL CHECK (edge_count BETWEEN 0 AND 16384), + total_bytes bigint NOT NULL CHECK (total_bytes BETWEEN 0 AND 67108864), + state text NOT NULL CHECK (state IN ('PREPARING', 'COMMITTED')), + created_at timestamptz NOT NULL DEFAULT now(), + committed_at timestamptz, + CHECK ((state = 'COMMITTED') = (committed_at IS NOT NULL)) + ); + CREATE TABLE mst2_metadata_prepare_page ( + prepare_id text NOT NULL REFERENCES mst2_metadata_prepare(prepare_id), + page_id bytea NOT NULL CHECK (octet_length(page_id) = 32), + expected_size integer NOT NULL CHECK (expected_size BETWEEN {HEADER_LEN} AND {PAGE_MAX_BYTES}), + PRIMARY KEY (prepare_id, page_id) + )" + )).await?; + manager + .get_connection() + .execute_raw(sea_orm::Statement::from_sql_and_values( + sea_orm::DbBackend::Postgres, + "INSERT INTO mst2_metadata_storage_scope(singleton,storage_uuid) VALUES(1,$1)", + [uuid::Uuid::new_v4().to_string().into()], + )) + .await + .map(|_| ()) + } + + async fn down(&self, _manager: &SchemaManager) -> Result<(), DbErr> { + // Preserve immutable payloads, prepare pins and recovery evidence. + Ok(()) + } +} diff --git a/src/jupiter/migration/mod.rs b/src/jupiter/migration/mod.rs index 81f1d7f5..7f964e4a 100644 --- a/src/jupiter/migration/mod.rs +++ b/src/jupiter/migration/mod.rs @@ -143,6 +143,7 @@ mod m20261005_000100_add_mst2_publication_request_digest; mod m20261005_000100_add_mst2_retention_durability; mod m20261005_000200_add_mst2_native_head; mod m20261005_000200_harden_mst2_retention_graph; +mod m20261005_000300_add_mst2_metadata_install; mod m20261006_000100_add_view_tables; mod runner; pub use m20260905_000100_add_push_queue::ensure_queue_control_seed; @@ -275,6 +276,7 @@ impl MigratorTrait for Migrator { Box::new(m20261005_000200_add_mst2_native_head::Migration), Box::new(m20261005_000100_add_mst2_retention_durability::Migration), Box::new(m20261005_000200_harden_mst2_retention_graph::Migration), + Box::new(m20261005_000300_add_mst2_metadata_install::Migration), Box::new(m20261006_000100_add_view_tables::Migration), ] } @@ -1186,7 +1188,7 @@ mod tests { async fn import_repo_alias_rows_canonicalized() { let names = migration_names(); assert_eq!( - &names[names.len() - 7..], + &names[names.len() - 8..], &[ "m20260923_000200_canonicalize_import_repo_paths".to_string(), "m20260925_000100_media_paging".to_string(), @@ -1194,9 +1196,10 @@ mod tests { "m20261005_000200_add_mst2_native_head".to_string(), "m20261005_000100_add_mst2_retention_durability".to_string(), "m20261005_000200_harden_mst2_retention_graph".to_string(), + "m20261005_000300_add_mst2_metadata_install".to_string(), VIEW_MIGRATION_NAME.to_string(), ], - "native retention and view tables are registered after media paging" + "native retention, metadata installation and view tables follow media paging" ); let db = alias_db().await; diff --git a/src/jupiter/service/native_publication_push_tests.rs b/src/jupiter/service/native_publication_push_tests.rs index 89290152..8f2f9f76 100644 --- a/src/jupiter/service/native_publication_push_tests.rs +++ b/src/jupiter/service/native_publication_push_tests.rs @@ -41,7 +41,15 @@ async fn native_fixture() -> (tempfile::TempDir, crate::jupiter::storage::Storag mono.save_refs(mega_refs::Model::new( path.clone(), MEGA_BRANCH_NAME.to_owned(), tip.id.to_string(), child.id.to_string(), false, ), None).await.unwrap(); - storage.mono_storage().initialize_native_publication(NATIVE_INSTANCE).await.unwrap(); + storage.push_queue_service.push_queue_storage.set_control_flags(Some(true), None, None).await.unwrap(); + mono.initialize_native_publication_for_maintenance( + NATIVE_INSTANCE, + &crate::jupiter::storage::native_publication_storage::NativeRoot { + commit: root_commit.id.to_string(), tree: root_tree.id.to_string(), + }, + ).await.unwrap(); + assert!(mono.read_native_publication_head(NATIVE_INSTANCE).await.is_err()); + storage.push_queue_service.push_queue_storage.set_control_flags(Some(false), None, None).await.unwrap(); (temp, storage, tip, path) } @@ -82,6 +90,43 @@ async fn real_same_tree_push_advances_one_global_certificate_and_n0_replay_stays assert_eq!(mst2_native_publication::Entity::find().count(mono.get_connection()).await.unwrap(),1); } +#[tokio::test] +async fn maintenance_cannot_reinitialize_published_history_after_head_loss() { + let (_temp, storage, tip, path) = native_fixture().await; + let mono = storage.mono_storage(); + let (new, payload) = save_same_tree_commit(&storage, &tip).await; + let id = wh03_enqueue_push(&storage, &path, &tip.id.to_string(), &new, &payload).await; + assert!(matches!(wh03_exec(&storage, id).await, ExecuteOutcome::Done { .. })); + let head = mono.read_native_publication_head(NATIVE_INSTANCE).await.unwrap(); + let certificates = mst2_native_publication::Entity::find().all(mono.get_connection()).await.unwrap(); + let receipts = mst2_publication::Entity::find().all(mono.get_connection()).await.unwrap(); + let outbox = mst2_publication_outbox::Entity::find().all(mono.get_connection()).await.unwrap(); + assert_eq!(certificates.len(), 1); + storage.push_queue_service.push_queue_storage.set_control_flags(Some(true), None, None).await.unwrap(); + mono.get_connection().execute_unprepared("DELETE FROM mst2_native_head").await.unwrap(); + for instance in [NATIVE_INSTANCE, "11111111-2222-4333-8444-555555555555"] { + let error = mono.initialize_native_publication_for_maintenance(instance, &head.root).await.unwrap_err(); + assert!(error.to_string().contains("native publication history exists")); + } + assert_eq!(mst2_native_head::Entity::find().count(mono.get_connection()).await.unwrap(), 0); + assert_eq!(mst2_native_publication::Entity::find().all(mono.get_connection()).await.unwrap(), certificates); + assert_eq!(mst2_publication::Entity::find().all(mono.get_connection()).await.unwrap(), receipts); + assert_eq!(mst2_publication_outbox::Entity::find().all(mono.get_connection()).await.unwrap(), outbox); + assert!(storage.push_queue_service.push_queue_storage.get_control().await.unwrap().paused); + mono.get_connection().execute_unprepared("DELETE FROM mst2_native_publication").await.unwrap(); + for instance in [NATIVE_INSTANCE, "11111111-2222-4333-8444-555555555555"] { + let error = mono.initialize_native_publication_for_maintenance(instance, &head.root).await.unwrap_err(); + assert!(error.to_string().contains("native publication history exists")); + } + assert_eq!(mst2_native_head::Entity::find().count(mono.get_connection()).await.unwrap(), 0); + assert_eq!(mst2_native_publication::Entity::find().count(mono.get_connection()).await.unwrap(), 0); + assert_eq!(mst2_publication::Entity::find().all(mono.get_connection()).await.unwrap(), receipts); + assert_eq!(mst2_publication_outbox::Entity::find().all(mono.get_connection()).await.unwrap(), outbox); + assert!(storage.push_queue_service.push_queue_storage.get_control().await.unwrap().paused); + let root = mono.get_main_ref("/").await.unwrap().unwrap(); + assert_eq!((root.ref_commit_hash, root.ref_tree_hash), (head.root.commit, head.root.tree)); +} + #[tokio::test] async fn real_merge_fails_closed_when_native_publication_is_enabled() { let (_temp, storage, tip, path) = native_fixture().await; diff --git a/src/jupiter/storage/mod.rs b/src/jupiter/storage/mod.rs index 6091852f..655a7d02 100644 --- a/src/jupiter/storage/mod.rs +++ b/src/jupiter/storage/mod.rs @@ -18,6 +18,7 @@ pub mod media_paging_storage; pub mod mono_storage; pub(crate) mod mst2_publication_storage; pub mod mst2_retention; +pub mod native_metadata_install; pub(crate) mod native_publication_storage; pub mod notification_storage; pub mod object_storage; diff --git a/src/jupiter/storage/mst2_retention.rs b/src/jupiter/storage/mst2_retention.rs index fe916803..11c9c456 100644 --- a/src/jupiter/storage/mst2_retention.rs +++ b/src/jupiter/storage/mst2_retention.rs @@ -37,7 +37,7 @@ const MAX_NODES: usize = 4096; const MAX_EDGES: usize = 16_384; const MAX_ROOTS: usize = 16; const MAX_PENDING_BATCH: u64 = 1000; -const RETENTION_LOCK_KEY: i32 = 1_296_717_362; +pub(crate) const RETENTION_LOCK_KEY: i32 = 1_296_717_362; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum GcClaim { diff --git a/src/jupiter/storage/native_metadata_install.rs b/src/jupiter/storage/native_metadata_install.rs new file mode 100644 index 00000000..7f7d80ed --- /dev/null +++ b/src/jupiter/storage/native_metadata_install.rs @@ -0,0 +1,802 @@ +//! Immutable metadata payload CAS and durable native preparation receipts. +//! No runtime, publication, lease or physical collector is enabled here. + +use std::{ + collections::{BTreeMap, BTreeSet}, + time::Duration, +}; + +use mst2_codec::metapage::{HEADER_LEN, PAGE_MAX_BYTES, Page, page_id}; +use sea_orm::{ + ColumnTrait, ConnectionTrait, DatabaseConnection, DatabaseTransaction, DbBackend, EntityTrait, + IsolationLevel, QueryFilter, QuerySelect, Statement, TransactionTrait, +}; +use serde_json::json; + +use super::mst2_retention::{PostgresRetentionRepository, RETENTION_LOCK_KEY}; +use crate::{ + callisto::{ + mst2_metadata_payload, mst2_metadata_prepare, mst2_metadata_prepare_page, + mst2_retention_edge, mst2_retention_node, mst2_retention_root, + }, + ceres::snapshot::{ + error::{SnapshotError, SnapshotErrorCode}, + metadata_install::MetadataInstallPlan, + pages::PreparedNativeMetadataRetention, + retention::RetentionRoot, + retention_dag::{ + MetadataDagCandidate, MetadataDagLimits, MetadataPagePayload, ValidatedMetadataDag, + }, + }, +}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum MetadataCommitPhase { + Intent, + Payload, + Finalize, +} + +#[derive(Debug, thiserror::Error)] +pub enum MetadataInstallError { + #[error(transparent)] + Rejected(#[from] SnapshotError), + #[error("native metadata {phase:?} commit outcome is unknown for operation {operation_id}")] + CommitUncertain { + operation_id: String, + manifest_digest: [u8; 32], + phase: MetadataCommitPhase, + }, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MetadataPrepareIntent { + prepare_id: String, + operation_id: String, + manifest_digest: [u8; 32], +} + +impl MetadataPrepareIntent { + pub fn prepare_id(&self) -> &str { + &self.prepare_id + } + pub fn operation_id(&self) -> &str { + &self.operation_id + } + pub fn manifest_digest(&self) -> [u8; 32] { + self.manifest_digest + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PreparedMetadataReceipt { + intent: MetadataPrepareIntent, + metadata_root: [u8; 32], + payload_bytes: u64, +} + +impl PreparedMetadataReceipt { + pub fn intent(&self) -> &MetadataPrepareIntent { + &self.intent + } + pub fn metadata_root(&self) -> [u8; 32] { + self.metadata_root + } + pub fn payload_bytes(&self) -> u64 { + self.payload_bytes + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum MetadataPrepareObservation { + Absent, + Preparing(MetadataPrepareIntent), + Committed(PreparedMetadataReceipt), +} + +struct StoredPlan { + record: mst2_metadata_prepare::Model, + plan: MetadataInstallPlan, +} + +impl StoredPlan { + fn intent(&self) -> Result { + Ok(MetadataPrepareIntent { + prepare_id: self.record.prepare_id.clone(), + operation_id: self.record.operation_id.clone(), + manifest_digest: self + .record + .manifest_digest + .as_slice() + .try_into() + .map_err(|_| integrity("invalid stored metadata manifest digest"))?, + }) + } + fn receipt(&self) -> Result { + Ok(PreparedMetadataReceipt { + intent: self.intent()?, + metadata_root: self.plan.root, + payload_bytes: self.plan.total_bytes, + }) + } +} + +#[derive(Clone)] +pub struct PostgresMetadataInstallRepository { + connection: DatabaseConnection, + barrier_timeout: Duration, + storage_scope: PrimaryStorageScope, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +struct PrimaryStorageScope { + storage_uuid: String, + database: String, + database_oid: i64, + schema: String, + schema_oid: i64, + server_address: Option, + server_port: Option, +} + +impl PostgresMetadataInstallRepository { + /// The connection must target the deployment's primary PostgreSQL database. + pub async fn new(connection: DatabaseConnection) -> Result { + let storage_scope = read_storage_scope(&connection).await?; + Ok(Self { + connection, + barrier_timeout: Duration::from_secs(5), + storage_scope, + }) + } + + pub async fn begin_intent( + &self, + operation_id: &str, + prepared: &PreparedNativeMetadataRetention, + ) -> Result { + validate_operation_id(operation_id)?; + let plan = prepared.install_plan()?; + let digest = plan.digest()?; + let txn = self.transaction().await?; + let result = async { + self.barrier(&txn).await?; + if let Some(stored) = load_plan(&txn, operation_id, &digest).await? { + return stored.intent(); + } + let id = uuid::Uuid::new_v4().to_string(); + let identity = &plan.identity; + txn.execute_raw(statement( + "INSERT INTO mst2_metadata_prepare (prepare_id, operation_id, manifest_digest, canonical_plan, + source_domain, tagged_root_tree_oid, scope, schema_version, metadata_codec, materialization_policy, + fs_semantics, access_projection, verification_revision, projection_revision, metadata_root, + node_count, edge_count, total_bytes, state) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,'PREPARING')", + [id.clone().into(), operation_id.into(), digest.to_vec().into(), plan.encode()?.into(), + identity.source_domain.clone().into(), identity.tagged_root_tree_oid.clone().into(), identity.scope.clone().into(), + (identity.schema_version as i16).into(), (identity.metadata_codec as i16).into(), + (identity.materialization_policy as i16).into(), (identity.fs_semantics as i16).into(), + (identity.access_projection as i16).into(), identity.verification_revision.into(), + (identity.projection_revision as i16).into(), plan.root.to_vec().into(), (plan.pages.len() as i32).into(), + (plan.edges.len() as i32).into(), (plan.total_bytes as i64).into()], + )).await.map_err(internal)?; + let pages: Vec<_> = plan.pages.iter().map(|(id,size)| json!({"page_id":hex::encode(id),"size":size})).collect(); + txn.execute_raw(statement( + "INSERT INTO mst2_metadata_prepare_page (prepare_id,page_id,expected_size) + SELECT $1,decode(p.page_id,'hex'),p.size FROM jsonb_to_recordset($2::jsonb) AS p(page_id text,size integer)", + [id.clone().into(), serde_json::to_string(&pages).map_err(internal)?.into()], + )).await.map_err(internal)?; + Ok(MetadataPrepareIntent { prepare_id:id, operation_id:operation_id.into(), manifest_digest:digest }) + }.await; + commit( + txn, + result, + operation_id, + digest, + MetadataCommitPhase::Intent, + ) + .await + } + + pub async fn install_page( + &self, + intent: &MetadataPrepareIntent, + payload: &MetadataPagePayload, + ) -> Result<(), MetadataInstallError> { + validate_payload(payload)?; + let txn = self.transaction().await?; + let result = async { + self.barrier(&txn).await?; + let stored = require_plan(&txn, intent).await?; + if stored.plan.pages.get(&payload.id) != Some(&payload.size) { + return Err(integrity( + "metadata payload is not a member of this fixed installation", + )); + } + txn.execute_raw(statement( + "INSERT INTO mst2_metadata_payload(page_id,metadata_codec,byte_size,payload) + VALUES ($1,$2,$3,$4) ON CONFLICT(page_id) DO NOTHING", + [ + payload.id.to_vec().into(), + (stored.plan.identity.metadata_codec as i16).into(), + (payload.size as i32).into(), + payload.bytes.clone().into(), + ], + )) + .await + .map_err(internal)?; + let existing = mst2_metadata_payload::Entity::find_by_id(payload.id.to_vec()) + .one(&txn) + .await + .map_err(internal)? + .ok_or_else(|| integrity("installed metadata payload disappeared"))?; + if existing.payload != payload.bytes + || existing.byte_size as u64 != payload.size + || existing.metadata_codec != stored.plan.identity.metadata_codec as i16 + { + return Err(integrity( + "immutable metadata payload identity conflicts with stored bytes", + )); + } + Ok(()) + } + .await; + commit( + txn, + result, + &intent.operation_id, + intent.manifest_digest, + MetadataCommitPhase::Payload, + ) + .await + } + + /// Verify all stored bytes outside the graph transaction. Immutable payloads + /// cannot change; LIVE state and coverage are checked again under the lock. + pub async fn load_installed_dag( + &self, + intent: &MetadataPrepareIntent, + ) -> Result { + let stored = require_plan(&self.connection, intent).await?; + load_installed_dag(&self.connection, &stored).await + } + + pub async fn finalize( + &self, + intent: &MetadataPrepareIntent, + ) -> Result { + let dag = self.load_installed_dag(intent).await?; + let txn = self.transaction().await?; + let result = self.finalize_in_txn(&txn, intent, &dag).await; + commit( + txn, + result, + &intent.operation_id, + intent.manifest_digest, + MetadataCommitPhase::Finalize, + ) + .await + } + + // This provisional result is private: only finalize's successful outer + // commit or the recovery barrier can issue a durable receipt to a caller. + async fn finalize_in_txn( + &self, + txn: &DatabaseTransaction, + intent: &MetadataPrepareIntent, + dag: &ValidatedMetadataDag, + ) -> Result { + self.barrier(txn).await?; + let stored = require_plan(txn, intent).await?; + let expected: BTreeSet<_> = dag + .payloads() + .iter() + .map(|page| (page.id, page.size)) + .collect(); + if dag.root() != stored.plan.root + || expected + != stored + .plan + .pages + .iter() + .map(|(id, size)| (*id, *size)) + .collect() + || dag + .edges() + .iter() + .map(|edge| (edge.parent.clone(), edge.child.clone())) + .collect::>() + != stored + .plan + .edges + .iter() + .map(|(parent, child)| (node_id(parent), node_id(child))) + .collect() + { + return Err(integrity( + "validated installed DAG differs from durable preparation plan", + )); + } + check_payload_coverage(txn, &stored).await?; + if stored.record.state == "COMMITTED" { + verify_graph(txn, &stored).await?; + return stored.receipt(); + } + PostgresRetentionRepository::retain_group_in_txn( + txn, + dag.nodes(), + dag.edges(), + &[RetentionRoot::Prepare(intent.prepare_id.clone())], + ) + .await?; + verify_graph(txn, &stored).await?; + let result = txn + .execute_raw(statement( + "UPDATE mst2_metadata_prepare SET state='COMMITTED',committed_at=now() + WHERE prepare_id=$1 AND manifest_digest=$2 AND state='PREPARING'", + [ + intent.prepare_id.clone().into(), + intent.manifest_digest.to_vec().into(), + ], + )) + .await + .map_err(internal)?; + if result.rows_affected() != 1 { + return Err(integrity("metadata prepare state changed during finalize")); + } + stored.receipt() + } + + /// Call with a fresh connection to the same primary after a commit error. + /// Lock acquisition is the completion barrier; absence before it proves nothing. + pub async fn inspect_prepare( + &self, + fresh_primary: &DatabaseConnection, + operation_id: &str, + digest: [u8; 32], + phase: MetadataCommitPhase, + ) -> Result { + validate_operation_id(operation_id)?; + let txn = fresh_primary + .begin_with_config(Some(IsolationLevel::ReadCommitted), None) + .await + .map_err(|_| uncertain(operation_id, digest, phase))?; + if self.barrier(&txn).await.is_err() { + let _ = txn.rollback().await; + return Err(uncertain(operation_id, digest, phase)); + } + let result = async { + match load_plan(&txn, operation_id, &digest).await? { + None => Ok(MetadataPrepareObservation::Absent), + Some(stored) if stored.record.state == "COMMITTED" => { + // Recovery rechecks the bounded bytes/DAG under the barrier, + // so corruption cannot turn a stored state into a receipt. + // This may read 64 MiB; normal finalization verifies outside + // its short graph transaction instead. + load_installed_dag(&txn, &stored).await?; + check_payload_coverage(&txn, &stored).await?; + verify_graph(&txn, &stored).await?; + Ok(MetadataPrepareObservation::Committed(stored.receipt()?)) + } + Some(stored) => Ok(MetadataPrepareObservation::Preparing(stored.intent()?)), + } + } + .await; + txn.rollback() + .await + .map_err(|_| uncertain(operation_id, digest, phase))?; + result.map_err(|error: SnapshotError| { + if error.code == SnapshotErrorCode::Internal { + uncertain(operation_id, digest, phase) + } else { + MetadataInstallError::Rejected(error) + } + }) + } + + async fn transaction(&self) -> Result { + if self.connection.get_database_backend() != DbBackend::Postgres { + return Err(internal("metadata installation requires PostgreSQL")); + } + self.connection + .begin_with_config(Some(IsolationLevel::ReadCommitted), None) + .await + .map_err(internal) + } + + async fn barrier(&self, txn: &DatabaseTransaction) -> Result<(), SnapshotError> { + if txn.get_database_backend() != DbBackend::Postgres { + return Err(internal("metadata installation requires PostgreSQL")); + } + if read_storage_scope(txn).await? != self.storage_scope { + return Err(internal( + "metadata recovery connection is outside the captured primary storage scope", + )); + } + let isolation = txn + .query_one_raw(statement("SHOW transaction_isolation", [])) + .await + .map_err(internal)? + .ok_or_else(|| internal("missing transaction isolation"))? + .try_get_by_index::(0) + .map_err(internal)?; + if isolation != "read committed" { + return Err(internal("metadata installation requires READ COMMITTED")); + } + let recovery = txn + .query_one_raw(statement("SELECT pg_is_in_recovery()", [])) + .await + .map_err(internal)? + .ok_or_else(|| internal("missing primary status"))? + .try_get_by_index::(0) + .map_err(internal)?; + if recovery { + return Err(internal( + "metadata installation recovery requires the primary", + )); + } + let timeout = format!("{}ms", self.barrier_timeout.as_millis().clamp(1, 5000)); + txn.execute_raw(statement( + "SELECT set_config('lock_timeout',$1,true)", + [timeout.into()], + )) + .await + .map_err(internal)?; + txn.execute_raw(statement( + "SELECT pg_advisory_xact_lock($1,hashtext(current_schema()))", + [RETENTION_LOCK_KEY.into()], + )) + .await + .map_err(internal)?; + Ok(()) + } +} + +async fn read_storage_scope( + connection: &C, +) -> Result { + if connection.get_database_backend() != DbBackend::Postgres { + return Err(internal("metadata storage scope requires PostgreSQL")); + } + let row = connection + .query_one_raw(statement( + "SELECT s.storage_uuid, current_database() AS database, d.oid::bigint AS database_oid, + current_schema() AS schema, n.oid::bigint AS schema_oid, + inet_server_addr()::text AS server_address, inet_server_port() AS server_port, + pg_is_in_recovery() AS replica FROM mst2_metadata_storage_scope s + JOIN pg_catalog.pg_database d ON d.datname=current_database() + JOIN pg_catalog.pg_namespace n ON n.nspname=current_schema() WHERE s.singleton=1", + [], + )) + .await + .map_err(internal)? + .ok_or_else(|| internal("metadata primary storage scope is missing"))?; + if row.try_get::("", "replica").map_err(internal)? { + return Err(internal("metadata storage scope requires the primary")); + } + let storage_uuid: String = row.try_get("", "storage_uuid").map_err(internal)?; + let id = uuid::Uuid::parse_str(&storage_uuid) + .map_err(|_| internal("invalid metadata storage UUID"))?; + if id.is_nil() || id.to_string() != storage_uuid { + return Err(internal("noncanonical metadata storage UUID")); + } + Ok(PrimaryStorageScope { + storage_uuid, + database: row.try_get("", "database").map_err(internal)?, + database_oid: row.try_get("", "database_oid").map_err(internal)?, + schema: row.try_get("", "schema").map_err(internal)?, + schema_oid: row.try_get("", "schema_oid").map_err(internal)?, + server_address: row.try_get("", "server_address").map_err(internal)?, + server_port: row.try_get("", "server_port").map_err(internal)?, + }) +} + +async fn load_plan( + connection: &C, + operation_id: &str, + digest: &[u8; 32], +) -> Result, SnapshotError> { + let Some(record) = mst2_metadata_prepare::Entity::find() + .filter(mst2_metadata_prepare::Column::OperationId.eq(operation_id)) + .one(connection) + .await + .map_err(internal)? + else { + return Ok(None); + }; + if record.manifest_digest.as_slice() != digest { + return Err(SnapshotError::new( + SnapshotErrorCode::Conflict, + "metadata operation ID is bound to a different manifest", + )); + } + let plan = MetadataInstallPlan::decode(&record.canonical_plan, digest)?; + let identity = &plan.identity; + if record.source_domain != identity.source_domain + || record.tagged_root_tree_oid != identity.tagged_root_tree_oid + || record.scope != identity.scope + || record.schema_version != identity.schema_version as i16 + || record.metadata_codec != identity.metadata_codec as i16 + || record.materialization_policy != identity.materialization_policy as i16 + || record.fs_semantics != identity.fs_semantics as i16 + || record.access_projection != identity.access_projection as i16 + || record.verification_revision != identity.verification_revision + || record.projection_revision != identity.projection_revision as i16 + || record.metadata_root.as_slice() != plan.root + || record.node_count as usize != plan.pages.len() + || record.edge_count as usize != plan.edges.len() + || record.total_bytes as u64 != plan.total_bytes + || !["PREPARING", "COMMITTED"].contains(&record.state.as_str()) + || (record.state == "COMMITTED") != record.committed_at.is_some() + { + return Err(integrity( + "stored metadata preparation fields disagree with their canonical plan", + )); + } + let id = uuid::Uuid::parse_str(&record.prepare_id) + .map_err(|_| integrity("invalid stored preparation identity"))?; + if id.is_nil() || id.to_string() != record.prepare_id { + return Err(integrity("noncanonical stored preparation identity")); + } + let rows = mst2_metadata_prepare_page::Entity::find() + .filter(mst2_metadata_prepare_page::Column::PrepareId.eq(&record.prepare_id)) + .limit((MetadataDagLimits::default().nodes + 1) as u64) + .all(connection) + .await + .map_err(internal)?; + let mut coverage = BTreeSet::new(); + for row in rows { + let page: [u8; 32] = row + .page_id + .as_slice() + .try_into() + .map_err(|_| integrity("invalid stored preparation page ID"))?; + coverage.insert((page, row.expected_size as u64)); + } + if coverage != plan.pages.iter().map(|(id, size)| (*id, *size)).collect() { + return Err(integrity( + "stored metadata preparation coverage differs from its canonical plan", + )); + } + Ok(Some(StoredPlan { record, plan })) +} + +async fn require_plan( + connection: &C, + intent: &MetadataPrepareIntent, +) -> Result { + validate_operation_id(&intent.operation_id)?; + let stored = load_plan(connection, &intent.operation_id, &intent.manifest_digest) + .await? + .ok_or_else(|| unavailable("metadata preparation intent is not durable"))?; + if stored.record.prepare_id != intent.prepare_id { + return Err(integrity("metadata intent identity mismatch")); + } + Ok(stored) +} + +async fn load_installed_dag( + connection: &C, + stored: &StoredPlan, +) -> Result { + let rows = mst2_metadata_payload::Entity::find() + .filter( + mst2_metadata_payload::Column::PageId + .is_in(stored.plan.pages.keys().map(|id| id.to_vec())), + ) + .all(connection) + .await + .map_err(internal)?; + let mut pages = Vec::with_capacity(rows.len()); + for row in rows { + let id: [u8; 32] = row + .page_id + .as_slice() + .try_into() + .map_err(|_| integrity("invalid stored page ID"))?; + if row.metadata_codec != stored.plan.identity.metadata_codec as i16 + || stored.plan.pages.get(&id) != Some(&(row.byte_size as u64)) + { + return Err(integrity( + "stored page profile or size disagrees with its manifest", + )); + } + let page = MetadataPagePayload { + id, + size: row.byte_size as u64, + bytes: row.payload, + }; + validate_payload(&page)?; + pages.push(page); + } + if pages.len() != stored.plan.pages.len() { + return Err(unavailable( + "native metadata installation has missing payloads", + )); + } + ValidatedMetadataDag::validate( + MetadataDagCandidate { + metadata_codec: stored.plan.identity.metadata_codec, + root: stored.plan.root, + pages, + edges: stored.plan.edges.iter().copied().collect(), + }, + MetadataDagLimits::default(), + ) +} + +async fn check_payload_coverage( + connection: &C, + stored: &StoredPlan, +) -> Result<(), SnapshotError> { + let missing=connection.query_one_raw(statement( + "SELECT p.page_id FROM mst2_metadata_prepare_page p LEFT JOIN mst2_metadata_payload b ON b.page_id=p.page_id + WHERE p.prepare_id=$1 AND (b.page_id IS NULL OR b.byte_size<>p.expected_size OR b.metadata_codec<>$2 + OR octet_length(b.payload)<>p.expected_size) LIMIT 1", + [stored.record.prepare_id.clone().into(),stored.record.metadata_codec.into()], + )).await.map_err(internal)?; + if missing.is_some() { + return Err(unavailable( + "durable metadata payload coverage is incomplete", + )); + } + Ok(()) +} + +async fn verify_graph( + connection: &C, + stored: &StoredPlan, +) -> Result<(), SnapshotError> { + let sizes: BTreeMap<_, _> = stored + .plan + .pages + .iter() + .map(|(id, size)| (node_id(id), *size)) + .collect(); + let ids: Vec<_> = sizes.keys().cloned().collect(); + let nodes = mst2_retention_node::Entity::find() + .filter(mst2_retention_node::Column::NodeId.is_in(ids.clone())) + .all(connection) + .await + .map_err(internal)?; + if nodes.len() != ids.len() { + return Err(unavailable("prepared metadata graph has missing nodes")); + } + for node in nodes { + if node.state != "LIVE" || node.kind != "page" { + return Err(unavailable("prepared metadata graph is not LIVE")); + } + if sizes.get(&node.node_id) != Some(&(node.bytes as u64)) { + return Err(integrity( + "prepared metadata graph bytes differ from its plan", + )); + } + } + let edges = mst2_retention_edge::Entity::find() + .filter(mst2_retention_edge::Column::ParentId.is_in(ids.clone())) + .limit((MetadataDagLimits::default().edges + 1) as u64) + .all(connection) + .await + .map_err(internal)?; + let actual: BTreeSet<_> = edges + .into_iter() + .map(|edge| (edge.parent_id, edge.child_id)) + .collect(); + let expected = stored + .plan + .edges + .iter() + .map(|(parent, child)| (node_id(parent), node_id(child))) + .collect(); + if actual != expected { + return Err(integrity( + "prepared metadata retention edges differ from its plan", + )); + } + let roots = mst2_retention_root::Entity::find() + .filter( + mst2_retention_root::Column::RootKey + .eq(format!("prepare:{}", stored.record.prepare_id)), + ) + .limit((MetadataDagLimits::default().nodes + 1) as u64) + .all(connection) + .await + .map_err(internal)?; + if roots.iter().any(|root| root.root_kind != "prepare") + || roots + .into_iter() + .map(|root| root.node_id) + .collect::>() + != ids.iter().cloned().collect() + { + return Err(unavailable("prepared metadata pin coverage is incomplete")); + } + let bad=connection.query_one_raw(statement( + "SELECT n.node_id FROM mst2_retention_node n WHERE n.node_id IN (SELECT jsonb_array_elements_text($1::jsonb)) + AND n.incoming_refs<>(SELECT count(*) FROM mst2_retention_edge e WHERE e.child_id=n.node_id) LIMIT 1", + [serde_json::to_string(&ids).map_err(internal)?.into()], + )).await.map_err(internal)?; + if bad.is_some() { + return Err(integrity("prepared metadata graph counter audit failed")); + } + Ok(()) +} + +fn validate_payload(payload: &MetadataPagePayload) -> Result<(), SnapshotError> { + if !(HEADER_LEN..=PAGE_MAX_BYTES).contains(&payload.bytes.len()) { + return Err(SnapshotError::new( + SnapshotErrorCode::LimitExceeded, + "metadata page exceeds protocol length bounds", + )); + } + if payload.size != payload.bytes.len() as u64 || page_id(&payload.bytes) != payload.id { + return Err(SnapshotError::new( + SnapshotErrorCode::DigestMismatch, + "metadata payload size or protocol page ID mismatch", + )); + } + Page::decode(&payload.bytes) + .map_err(|error| integrity(&format!("invalid canonical metadata page: {error}")))?; + Ok(()) +} + +fn validate_operation_id(operation: &str) -> Result<(), SnapshotError> { + if operation.is_empty() || operation.len() > 255 || operation.contains('\0') { + return Err(SnapshotError::new( + SnapshotErrorCode::InvalidRequest, + "metadata operation ID must be 1..=255 UTF8 bytes", + )); + } + Ok(()) +} + +async fn commit( + txn: DatabaseTransaction, + result: Result, + operation: &str, + digest: [u8; 32], + phase: MetadataCommitPhase, +) -> Result { + match result { + Ok(value) => { + txn.commit() + .await + .map_err(|_| uncertain(operation, digest, phase))?; + Ok(value) + } + Err(error) => { + txn.rollback().await.map_err(internal)?; + Err(error.into()) + } + } +} +fn uncertain( + operation: &str, + digest: [u8; 32], + phase: MetadataCommitPhase, +) -> MetadataInstallError { + MetadataInstallError::CommitUncertain { + operation_id: operation.into(), + manifest_digest: digest, + phase, + } +} +fn node_id(id: &[u8; 32]) -> String { + format!("page:sha256:{}", hex::encode(id)) +} +fn statement(sql: &str, values: [sea_orm::Value; N]) -> Statement { + Statement::from_sql_and_values(DbBackend::Postgres, sql, values) +} +fn internal(error: impl std::fmt::Display) -> SnapshotError { + SnapshotError::new(SnapshotErrorCode::Internal, error.to_string()) +} +fn integrity(message: &str) -> SnapshotError { + SnapshotError::new(SnapshotErrorCode::IntegrityError, message) +} +fn unavailable(message: &str) -> SnapshotError { + SnapshotError::new(SnapshotErrorCode::ObjectUnavailable, message) +} + +#[cfg(test)] +#[path = "native_metadata_install_tests.rs"] +mod tests; diff --git a/src/jupiter/storage/native_metadata_install_tests.rs b/src/jupiter/storage/native_metadata_install_tests.rs new file mode 100644 index 00000000..4b47fe38 --- /dev/null +++ b/src/jupiter/storage/native_metadata_install_tests.rs @@ -0,0 +1,1226 @@ +use std::sync::{ + Arc, + atomic::{AtomicBool, Ordering}, +}; + +use mst2_codec::metapage::{Entry, EntryKind}; +use sea_orm::{Database, PaginatorTrait}; +use sea_orm_migration::MigratorTrait; +use tokio::{ + io::{AsyncReadExt, AsyncWriteExt}, + net::{TcpListener, TcpStream}, + sync::{Notify, watch}, + task::{JoinHandle, JoinSet}, +}; + +use super::*; +use crate::{ + ceres::snapshot::retention_dag::MetadataDagBuilder, + jupiter::{ + migration::Migrator, + tests::{TestSchemaGuard, test_db_config}, + }, +}; + +fn prepared(scope: &str) -> PreparedNativeMetadataRetention { + let child_entries = [Entry::file(EntryKind::Regular, b"file", 3, [42; 32])]; + let child = Page::build(&child_entries).unwrap(); + let entries = [ + Entry::dir(b"one", page_id(&child)), + Entry::dir(b"two", page_id(&child)), + ]; + let root = Page::build(&entries).unwrap(); + let mut builder = MetadataDagBuilder::new(MetadataDagLimits::default()); + builder.add_directory(&child, &child_entries).unwrap(); + builder.add_directory(&root, &entries).unwrap(); + PreparedNativeMetadataRetention::test_installation( + Arc::new(builder.finish(page_id(&root)).unwrap()), + scope, + ) +} + +async fn fixture() -> ( + DatabaseConnection, + DatabaseConnection, + TestSchemaGuard, + String, +) { + let temp = tempfile::TempDir::new().unwrap(); + let (config, schema) = test_db_config(temp.path()).await; + let first = Database::connect(config.db_url.clone()).await.unwrap(); + Migrator::up(&first, None).await.unwrap(); + let second = Database::connect(config.db_url.clone()).await.unwrap(); + (first, second, schema, config.db_url) +} + +async fn install( + repository: &PostgresMetadataInstallRepository, + intent: &MetadataPrepareIntent, + prepared: &PreparedNativeMetadataRetention, +) { + for page in prepared.dag().payloads() { + repository.install_page(intent, page).await.unwrap(); + } +} + +fn rejected(error: MetadataInstallError) -> SnapshotErrorCode { + match error { + MetadataInstallError::Rejected(error) => error.code, + _ => panic!("expected definite rejection"), + } +} + +#[tokio::test] +async fn native_metadata_two_primary_connections_replay_one_receipt_without_recounting() { + let (first, second, _schema, _url) = fixture().await; + let a = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let b = PostgresMetadataInstallRepository::new(second.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let (ia, ib) = tokio::join!( + a.begin_intent("same", &prepared), + b.begin_intent("same", &prepared) + ); + let ia = ia.unwrap(); + assert_eq!(ia, ib.unwrap()); + for page in prepared.dag().payloads() { + let (left, right) = tokio::join!(a.install_page(&ia, page), b.install_page(&ia, page)); + left.unwrap(); + right.unwrap(); + } + let (ra, rb) = tokio::join!(a.finalize(&ia), b.finalize(&ia)); + assert_eq!(ra.unwrap(), rb.unwrap()); + assert_eq!( + mst2_metadata_prepare::Entity::find() + .count(&first) + .await + .unwrap(), + 1 + ); + assert_eq!( + mst2_metadata_payload::Entity::find() + .count(&first) + .await + .unwrap(), + 2 + ); + assert_eq!( + mst2_retention_edge::Entity::find() + .count(&first) + .await + .unwrap(), + 1 + ); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&first) + .await + .unwrap(), + 2 + ); + let child = prepared.dag().edges()[0].child.clone(); + assert_eq!( + PostgresRetentionRepository::new(first) + .node(&child) + .await + .unwrap() + .unwrap() + .incoming_refs, + 1 + ); +} + +#[tokio::test] +async fn native_metadata_operation_conflict_and_multibyte_byte_limit_reject_without_writes() { + let (first, second, _schema, _url) = fixture().await; + let a = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let b = PostgresMetadataInstallRepository::new(second) + .await + .unwrap(); + a.begin_intent("bound", &prepared("/")).await.unwrap(); + assert_eq!( + rejected( + b.begin_intent("bound", &prepared("/scope")) + .await + .unwrap_err() + ), + SnapshotErrorCode::Conflict + ); + assert_eq!( + rejected( + a.begin_intent(&"é".repeat(128), &prepared("/")) + .await + .unwrap_err() + ), + SnapshotErrorCode::InvalidRequest + ); + assert_eq!( + mst2_metadata_prepare::Entity::find() + .count(&first) + .await + .unwrap(), + 1 + ); + assert_eq!( + mst2_metadata_payload::Entity::find() + .count(&first) + .await + .unwrap(), + 0 + ); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&first) + .await + .unwrap(), + 0 + ); +} + +#[tokio::test] +async fn native_metadata_protocol_hash_and_immutable_bytes_conflicts_cannot_be_overwritten() { + let (first, _second, _schema, _url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository.begin_intent("bytes", &prepared).await.unwrap(); + let original = &prepared.dag().payloads()[0]; + let mut invalid = original.clone(); + invalid.id[0] ^= 1; + assert_eq!( + rejected( + repository + .install_page(&intent, &invalid) + .await + .unwrap_err() + ), + SnapshotErrorCode::DigestMismatch + ); + let mut oversized = original.clone(); + oversized.bytes.resize(PAGE_MAX_BYTES + 1, 0); + oversized.size = oversized.bytes.len() as u64; + oversized.id = page_id(&oversized.bytes); + assert_eq!( + rejected( + repository + .install_page(&intent, &oversized) + .await + .unwrap_err() + ), + SnapshotErrorCode::LimitExceeded + ); + // A corrupt existing row is never replaced by a good upload. + let mut corrupt = original.bytes.clone(); + corrupt[0] ^= 1; + first.execute_raw(statement("INSERT INTO mst2_metadata_payload(page_id,metadata_codec,byte_size,payload) VALUES($1,1,$2,$3)", + [original.id.to_vec().into(),(original.size as i32).into(),corrupt.clone().into()])).await.unwrap(); + assert_eq!( + rejected( + repository + .install_page(&intent, original) + .await + .unwrap_err() + ), + SnapshotErrorCode::IntegrityError + ); + assert_eq!( + mst2_metadata_payload::Entity::find_by_id(original.id.to_vec()) + .one(&first) + .await + .unwrap() + .unwrap() + .payload, + corrupt + ); + assert!( + first + .execute_unprepared("UPDATE mst2_metadata_payload SET payload=payload") + .await + .is_err() + ); + assert!( + first + .execute_unprepared("DELETE FROM mst2_metadata_payload") + .await + .is_err() + ); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&first) + .await + .unwrap(), + 0 + ); +} + +#[tokio::test] +async fn native_metadata_partial_installation_survives_repository_and_connection_restart() { + let (first, second, _schema, _url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository.begin_intent("restart", &prepared).await.unwrap(); + repository + .install_page(&intent, &prepared.dag().payloads()[0]) + .await + .unwrap(); + assert_eq!( + mst2_retention_node::Entity::find() + .count(&first) + .await + .unwrap(), + 0 + ); + drop(repository); + first.close().await.unwrap(); + let restarted = PostgresMetadataInstallRepository::new(second.clone()) + .await + .unwrap(); + assert!(matches!( + restarted + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Payload + ) + .await + .unwrap(), + MetadataPrepareObservation::Preparing(_) + )); + assert_eq!( + restarted + .load_installed_dag(&intent) + .await + .unwrap_err() + .code, + SnapshotErrorCode::ObjectUnavailable + ); + install(&restarted, &intent, &prepared).await; + let receipt = restarted.finalize(&intent).await.unwrap(); + assert_eq!(receipt.metadata_root(), prepared.dag().root()); + assert_eq!(receipt.payload_bytes(), prepared.dag().payload_bytes()); +} + +#[tokio::test] +async fn native_metadata_half_installation_survives_real_process_kill() { + const CHILD_DB: &str = "MEGA_MST2_METADATA_CRASH_CHILD_DB"; + const CHECKPOINT: &str = "MEGA_MST2_METADATA_CRASH_CHECKPOINT"; + const OPERATION: &str = "process-crash"; + if let Ok(url) = std::env::var(CHILD_DB) { + let database = Database::connect(url).await.unwrap(); + let repository = PostgresMetadataInstallRepository::new(database) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository.begin_intent(OPERATION, &prepared).await.unwrap(); + repository + .install_page(&intent, &prepared.dag().payloads()[0]) + .await + .unwrap(); + let checkpoint = std::path::PathBuf::from(std::env::var(CHECKPOINT).unwrap()); + let staging = checkpoint.with_extension("staging"); + std::fs::write( + &staging, + format!( + "{}\n{}", + intent.prepare_id(), + hex::encode(intent.manifest_digest()) + ), + ) + .unwrap(); + std::fs::rename(staging, checkpoint).unwrap(); + // Keep repository/connection alive until the parent kills this process. + std::future::pending::<()>().await; + unreachable!(); + } + let (first, second, _schema, url) = fixture().await; + let checkpoint_dir = tempfile::tempdir().unwrap(); + let checkpoint = checkpoint_dir.path().join("partial-installation"); + let full_name = concat!( + module_path!(), + "::native_metadata_half_installation_survives_real_process_kill" + ); + let test_name = full_name.split_once("::").unwrap().1; + let mut worker = tokio::process::Command::new(std::env::current_exe().unwrap()) + .args(["--exact", test_name, "--nocapture"]) + .env(CHILD_DB, url) + .env(CHECKPOINT, &checkpoint) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::inherit()) + .kill_on_drop(true) + .spawn() + .unwrap(); + let marker = tokio::time::timeout(Duration::from_secs(60), async { + loop { + if checkpoint.exists() { + break std::fs::read_to_string(&checkpoint).unwrap(); + } + if let Some(status) = worker.try_wait().unwrap() { + panic!("metadata crash child exited before partial durable commit: {status}"); + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await + .expect("metadata crash child did not reach its durable checkpoint"); + worker.kill().await.unwrap(); + let status = worker.wait().await.unwrap(); + assert!(!status.success()); + #[cfg(unix)] + { + use std::os::unix::process::ExitStatusExt; + assert_eq!(status.signal(), Some(libc::SIGKILL)); + } + first.close().await.unwrap(); + let restarted = PostgresMetadataInstallRepository::new(second.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let digest = prepared.install_plan().unwrap().digest().unwrap(); + let original = match restarted + .inspect_prepare(&second, OPERATION, digest, MetadataCommitPhase::Payload) + .await + .unwrap() + { + MetadataPrepareObservation::Preparing(intent) => intent, + other => panic!("half installation exposed a false committed receipt: {other:?}"), + }; + let mut marker = marker.lines(); + assert_eq!(marker.next(), Some(original.prepare_id())); + assert_eq!(marker.next(), Some(hex::encode(digest).as_str())); + assert_eq!( + restarted.begin_intent(OPERATION, &prepared).await.unwrap(), + original + ); + assert_eq!( + mst2_metadata_payload::Entity::find() + .count(&second) + .await + .unwrap(), + 1 + ); + assert_eq!( + mst2_metadata_prepare_page::Entity::find() + .count(&second) + .await + .unwrap(), + prepared.dag().payloads().len() as u64 + ); + assert_eq!( + rejected(restarted.finalize(&original).await.unwrap_err()), + SnapshotErrorCode::ObjectUnavailable + ); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + 0 + ); + assert_eq!( + mst2_retention_node::Entity::find() + .count(&second) + .await + .unwrap(), + 0 + ); + install(&restarted, &original, &prepared).await; + let receipt = restarted.finalize(&original).await.unwrap(); + assert_eq!(receipt.intent(), &original); + assert_eq!(restarted.finalize(&original).await.unwrap(), receipt); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + prepared.dag().payloads().len() as u64 + ); + assert_eq!( + mst2_retention_edge::Entity::find() + .count(&second) + .await + .unwrap(), + 1 + ); +} + +async fn overwrite_installed_payload_for_test( + connection: &DatabaseConnection, + page: &MetadataPagePayload, + bytes: Vec, +) { + assert_eq!(bytes.len(), page.bytes.len()); + let txn = connection.begin().await.unwrap(); + txn.execute_unprepared( + "ALTER TABLE mst2_metadata_payload DISABLE TRIGGER mst2_metadata_payload_immutable", + ) + .await + .unwrap(); + txn.execute_raw(statement( + "UPDATE mst2_metadata_payload SET payload=$2 WHERE page_id=$1", + [page.id.to_vec().into(), bytes.into()], + )) + .await + .unwrap(); + txn.execute_unprepared( + "ALTER TABLE mst2_metadata_payload ENABLE TRIGGER mst2_metadata_payload_immutable", + ) + .await + .unwrap(); + txn.commit().await.unwrap(); +} + +#[tokio::test] +async fn native_metadata_committed_recovery_rejects_same_length_corruption_and_preserves_pins() { + let (first, second, _schema, _url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository + .begin_intent("corruption", &prepared) + .await + .unwrap(); + install(&repository, &intent, &prepared).await; + let receipt = repository.finalize(&intent).await.unwrap(); + let page = &prepared.dag().payloads()[0]; + let counters = mst2_retention_node::Entity::find() + .all(&first) + .await + .unwrap() + .into_iter() + .map(|node| (node.node_id, node.incoming_refs)) + .collect::>(); + let mut corrupt = page.bytes.clone(); + corrupt[0] ^= 1; + overwrite_installed_payload_for_test(&first, page, corrupt).await; + assert_eq!( + rejected( + repository + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap_err() + ), + SnapshotErrorCode::DigestMismatch + ); + assert_eq!( + rejected(repository.finalize(&intent).await.unwrap_err()), + SnapshotErrorCode::DigestMismatch + ); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + prepared.dag().payloads().len() as u64 + ); + assert_eq!( + mst2_retention_node::Entity::find() + .all(&second) + .await + .unwrap() + .into_iter() + .map(|node| (node.node_id, node.incoming_refs)) + .collect::>(), + counters + ); + assert_eq!( + mst2_metadata_prepare::Entity::find_by_id(intent.prepare_id().to_owned()) + .one(&second) + .await + .unwrap() + .unwrap() + .state, + "COMMITTED" + ); + overwrite_installed_payload_for_test(&first, page, page.bytes.clone()).await; + assert_eq!( + repository + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap(), + MetadataPrepareObservation::Committed(receipt) + ); +} + +#[tokio::test] +async fn native_metadata_outer_rollback_preserves_intent_but_exposes_no_graph_or_receipt() { + let (first, second, _schema, _url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository + .begin_intent("rollback", &prepared) + .await + .unwrap(); + install(&repository, &intent, &prepared).await; + let dag = repository.load_installed_dag(&intent).await.unwrap(); + let txn = repository.transaction().await.unwrap(); + repository + .finalize_in_txn(&txn, &intent, &dag) + .await + .unwrap(); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + 0 + ); + assert_eq!( + mst2_metadata_prepare::Entity::find_by_id(intent.prepare_id().to_owned()) + .one(&second) + .await + .unwrap() + .unwrap() + .state, + "PREPARING" + ); + txn.rollback().await.unwrap(); + assert!(matches!( + repository + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap(), + MetadataPrepareObservation::Preparing(_) + )); + assert_eq!( + mst2_retention_node::Entity::find() + .count(&second) + .await + .unwrap(), + 0 + ); + repository.finalize(&intent).await.unwrap(); +} + +#[tokio::test] +async fn native_metadata_final_lock_rechecks_deleting_after_payload_verification() { + let (first, second, _schema, _url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository + .begin_intent("deleting", &prepared) + .await + .unwrap(); + install(&repository, &intent, &prepared).await; + let dag = repository.load_installed_dag(&intent).await.unwrap(); + let graph = PostgresRetentionRepository::new(second.clone()); + graph + .retain_group(dag.nodes(), dag.edges(), &[]) + .await + .unwrap(); + let root = node_id(&dag.root()); + assert_eq!( + graph.mark_deleting("gc-root", &root).await.unwrap(), + super::super::mst2_retention::GcClaim::Marked + ); + let txn = repository.transaction().await.unwrap(); + assert_eq!( + repository + .finalize_in_txn(&txn, &intent, &dag) + .await + .unwrap_err() + .code, + SnapshotErrorCode::ObjectUnavailable + ); + txn.rollback().await.unwrap(); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + 0 + ); + assert_eq!( + mst2_metadata_prepare::Entity::find_by_id(intent.prepare_id().to_owned()) + .one(&second) + .await + .unwrap() + .unwrap() + .state, + "PREPARING" + ); + assert_eq!(graph.node(&root).await.unwrap().unwrap().state, "DELETING"); +} + +#[tokio::test] +async fn native_metadata_receipt_replay_rejects_lost_pin_without_recreating_it() { + let (first, second, _schema, _url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository + .begin_intent("lost-pin", &prepared) + .await + .unwrap(); + install(&repository, &intent, &prepared).await; + repository.finalize(&intent).await.unwrap(); + PostgresRetentionRepository::new(first) + .release_root(&RetentionRoot::Prepare(intent.prepare_id().into())) + .await + .unwrap(); + assert_eq!( + rejected( + repository + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap_err() + ), + SnapshotErrorCode::ObjectUnavailable + ); + assert_eq!( + rejected(repository.finalize(&intent).await.unwrap_err()), + SnapshotErrorCode::ObjectUnavailable + ); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + 0 + ); + assert_eq!( + mst2_metadata_payload::Entity::find() + .count(&second) + .await + .unwrap(), + 2 + ); +} + +#[tokio::test] +async fn native_metadata_recovery_lock_timeout_is_unknown_and_never_releases_pins() { + let (first, second, _schema, _url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository.begin_intent("timeout", &prepared).await.unwrap(); + install(&repository, &intent, &prepared).await; + repository.finalize(&intent).await.unwrap(); + let blocker = repository.transaction().await.unwrap(); + repository.barrier(&blocker).await.unwrap(); + let mut recovery = PostgresMetadataInstallRepository::new(second.clone()) + .await + .unwrap(); + recovery.barrier_timeout = Duration::from_millis(25); + assert!(matches!( + recovery + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap_err(), + MetadataInstallError::CommitUncertain { + phase: MetadataCommitPhase::Finalize, + .. + } + )); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + 2 + ); + blocker.rollback().await.unwrap(); + assert!(matches!( + recovery + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap(), + MetadataPrepareObservation::Committed(_) + )); +} + +#[tokio::test] +async fn native_metadata_stored_identity_and_plan_tampering_fail_closed_on_restart() { + let (first, second, _schema, _url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let intent = repository + .begin_intent("profile", &prepared("/")) + .await + .unwrap(); + first + .execute_unprepared("UPDATE mst2_metadata_prepare SET projection_revision=2") + .await + .unwrap(); + assert_eq!( + rejected( + repository + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Intent + ) + .await + .unwrap_err() + ), + SnapshotErrorCode::IntegrityError + ); + first.execute_unprepared("UPDATE mst2_metadata_prepare SET projection_revision=1,canonical_plan=canonical_plan || decode('ff','hex')").await.unwrap(); + assert_eq!( + rejected( + repository + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Intent + ) + .await + .unwrap_err() + ), + SnapshotErrorCode::IntegrityError + ); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + 0 + ); +} + +#[derive(Clone, Copy)] +enum Fault { + BeforeCommit, + AfterCommit, + PrepareRead, +} + +struct PgCommitFaultProxy { + url: String, + armed: Arc, + fired: Arc, + commit_observed: Arc, + changed: Arc, + task: JoinHandle<()>, +} + +impl PgCommitFaultProxy { + async fn start(original: &str, fault: Fault) -> Self { + let mut url = url::Url::parse(original).unwrap(); + let host = url.host_str().unwrap().to_owned(); + let port = url.port().unwrap_or(5432); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + url.set_host(Some("127.0.0.1")).unwrap(); + url.set_port(Some(address.port())).unwrap(); + let pairs: Vec<_> = url + .query_pairs() + .filter(|(key, _)| key != "sslmode") + .map(|(key, value)| (key.into_owned(), value.into_owned())) + .collect(); + url.set_query(None); + url.query_pairs_mut() + .extend_pairs(pairs) + .append_pair("sslmode", "disable"); + let armed = Arc::new(AtomicBool::new(false)); + let fired = Arc::new(AtomicBool::new(false)); + let commit_observed = Arc::new(AtomicBool::new(false)); + let changed = Arc::new(Notify::new()); + let a = armed.clone(); + let f = fired.clone(); + let c = commit_observed.clone(); + let n = changed.clone(); + let task = tokio::spawn(async move { + let mut connections = JoinSet::new(); + loop { + tokio::select! { + accepted=listener.accept()=>{ + let (client,_)=accepted.unwrap();let host=host.clone(); + let a=a.clone();let f=f.clone();let c=c.clone();let n=n.clone(); + connections.spawn(async move { + let server=TcpStream::connect((host.as_str(),port)).await?; + forward_connection(client,server,fault,a,f,c,n).await + }); + } + _=connections.join_next(),if !connections.is_empty()=>{} + } + } + }); + Self { + url: url.to_string(), + armed, + fired, + commit_observed, + changed, + task, + } + } + async fn wait_for_fault(&self) { + tokio::time::timeout(Duration::from_secs(10), async { + loop { + let changed = self.changed.notified(); + if self.fired.load(Ordering::SeqCst) { + break; + } + changed.await; + } + }) + .await + .expect("proxy never observed the real COMMIT fault"); + } +} +impl Drop for PgCommitFaultProxy { + fn drop(&mut self) { + self.task.abort(); + } +} + +async fn read_frame( + reader: &mut R, +) -> std::io::Result<(u8, Vec)> { + let kind = reader.read_u8().await?; + let length = reader.read_u32().await?; + if !(4..=4 * 1024 * 1024).contains(&length) { + return Err(std::io::Error::other("invalid PostgreSQL frame length")); + } + let mut body = vec![0; length as usize - 4]; + reader.read_exact(&mut body).await?; + Ok((kind, body)) +} +async fn write_frame( + writer: &mut W, + kind: u8, + body: &[u8], +) -> std::io::Result<()> { + writer.write_u8(kind).await?; + writer.write_u32(body.len() as u32 + 4).await?; + writer.write_all(body).await?; + writer.flush().await +} + +async fn forward_connection( + mut client: TcpStream, + mut server: TcpStream, + fault: Fault, + armed: Arc, + fired: Arc, + commit_observed: Arc, + changed: Arc, +) -> std::io::Result<()> { + // sslmode=disable makes the first message a bounded StartupMessage. + let length = client.read_u32().await?; + if !(8..=65536).contains(&length) { + return Err(std::io::Error::other("invalid PostgreSQL startup length")); + } + let mut startup = vec![0; length as usize - 4]; + client.read_exact(&mut startup).await?; + server.write_u32(length).await?; + server.write_all(&startup).await?; + server.flush().await?; + let (mut cr, mut cw) = client.into_split(); + let (mut sr, mut sw) = server.into_split(); + let suppress = Arc::new(AtomicBool::new(false)); + let sender_suppress = suppress.clone(); + let (stop, mut stopped) = watch::channel(false); + let sender_stop = stop.clone(); + let sender_fired = fired.clone(); + let sender_changed = changed.clone(); + let to_server = async move { + loop { + let frame = tokio::select! {result=read_frame(&mut cr)=>result?, _=stopped.changed()=>return Ok::<_,std::io::Error>(())}; + let is_commit = frame.0 == b'Q' && frame.1.as_slice() == b"COMMIT\0"; + let fault_target = match fault { + Fault::PrepareRead => { + b"PQ".contains(&frame.0) + && frame + .1 + .windows(b"mst2_metadata_prepare".len()) + .any(|bytes| bytes == b"mst2_metadata_prepare") + } + _ => is_commit, + }; + if fault_target && armed.swap(false, Ordering::SeqCst) { + match fault { + Fault::BeforeCommit | Fault::PrepareRead => { + sender_fired.store(true, Ordering::SeqCst); + sender_changed.notify_one(); + sw.shutdown().await?; + sender_stop.send_replace(true); + return Ok(()); + } + Fault::AfterCommit => sender_suppress.store(true, Ordering::SeqCst), + } + } + write_frame(&mut sw, frame.0, &frame.1).await?; + } + }; + let mut stopped = stop.subscribe(); + let to_client = async move { + loop { + let frame = tokio::select! {result=read_frame(&mut sr)=>result?, _=stopped.changed()=>{cw.shutdown().await?;return Ok::<_,std::io::Error>(());}}; + if suppress.load(Ordering::SeqCst) { + if frame.0 == b'C' && frame.1.as_slice() == b"COMMIT\0" { + commit_observed.store(true, Ordering::SeqCst); + fired.store(true, Ordering::SeqCst); + changed.notify_one(); + cw.shutdown().await?; + stop.send_replace(true); + return Ok(()); + } + } else { + write_frame(&mut cw, frame.0, &frame.1).await?; + } + } + }; + tokio::try_join!(to_server, to_client)?; + Ok(()) +} + +async fn commit_fault_case(fault: Fault) { + let (direct, recovery, _schema, url) = fixture().await; + let proxy = PgCommitFaultProxy::start(&url, fault).await; + let mut options = sea_orm::ConnectOptions::new(proxy.url.clone()); + options.max_connections(1).min_connections(1); + let proxied = Database::connect(options).await.unwrap(); + let repository = PostgresMetadataInstallRepository::new(proxied) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository + .begin_intent("commit-fault", &prepared) + .await + .unwrap(); + install(&repository, &intent, &prepared).await; + proxy.armed.store(true, Ordering::SeqCst); + let error = tokio::time::timeout(Duration::from_secs(15), repository.finalize(&intent)) + .await + .unwrap() + .unwrap_err(); + assert!(matches!( + error, + MetadataInstallError::CommitUncertain { + phase: MetadataCommitPhase::Finalize, + .. + } + )); + proxy.wait_for_fault().await; + let observed = repository + .inspect_prepare( + &recovery, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize, + ) + .await + .unwrap(); + match fault { + Fault::AfterCommit => { + assert!(proxy.commit_observed.load(Ordering::SeqCst)); + assert!(matches!(observed, MetadataPrepareObservation::Committed(_))); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&direct) + .await + .unwrap(), + 2 + ); + assert_eq!( + mst2_retention_edge::Entity::find() + .count(&direct) + .await + .unwrap(), + 1 + ); + let restarted = PostgresMetadataInstallRepository::new(direct.clone()) + .await + .unwrap(); + restarted.finalize(&intent).await.unwrap(); + assert_eq!( + mst2_retention_edge::Entity::find() + .count(&direct) + .await + .unwrap(), + 1 + ); + } + Fault::BeforeCommit => { + assert!(!proxy.commit_observed.load(Ordering::SeqCst)); + assert!(matches!(observed, MetadataPrepareObservation::Preparing(_))); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&direct) + .await + .unwrap(), + 0 + ); + assert_eq!( + mst2_retention_node::Entity::find() + .count(&direct) + .await + .unwrap(), + 0 + ); + PostgresMetadataInstallRepository::new(direct) + .await + .unwrap() + .finalize(&intent) + .await + .unwrap(); + } + Fault::PrepareRead => panic!("query fault uses its independent recovery test"), + } +} + +#[tokio::test] +async fn native_metadata_real_commit_response_loss_recovers_receipt_without_releasing_pin() { + commit_fault_case(Fault::AfterCommit).await; +} + +#[tokio::test] +async fn native_metadata_real_connection_loss_before_commit_recovers_preparing_without_false_success() + { + commit_fault_case(Fault::BeforeCommit).await; +} + +#[tokio::test] +async fn native_metadata_wrong_schema_primary_cannot_report_absent_for_another_committed_operation() +{ + let (first, second, _schema, _url) = fixture().await; + let (wrong_primary, _other, _wrong_schema, _wrong_url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository.begin_intent("scope", &prepared).await.unwrap(); + install(&repository, &intent, &prepared).await; + repository.finalize(&intent).await.unwrap(); + assert!(matches!( + repository + .inspect_prepare( + &wrong_primary, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap_err(), + MetadataInstallError::CommitUncertain { .. } + )); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + 2 + ); + assert_eq!( + mst2_metadata_prepare::Entity::find() + .count(&wrong_primary) + .await + .unwrap(), + 0 + ); + assert!(matches!( + repository + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap(), + MetadataPrepareObservation::Committed(_) + )); +} + +#[tokio::test] +async fn native_metadata_query_loss_after_recovery_barrier_preserves_typed_unknown() { + let (first, second, _schema, url) = fixture().await; + let repository = PostgresMetadataInstallRepository::new(first.clone()) + .await + .unwrap(); + let prepared = prepared("/"); + let intent = repository + .begin_intent("read-fault", &prepared) + .await + .unwrap(); + install(&repository, &intent, &prepared).await; + repository.finalize(&intent).await.unwrap(); + let proxy = PgCommitFaultProxy::start(&url, Fault::PrepareRead).await; + let mut options = sea_orm::ConnectOptions::new(proxy.url.clone()); + options.max_connections(1).min_connections(1); + let fresh = Database::connect(options).await.unwrap(); + proxy.armed.store(true, Ordering::SeqCst); + assert!(matches!( + repository + .inspect_prepare( + &fresh, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap_err(), + MetadataInstallError::CommitUncertain { .. } + )); + proxy.wait_for_fault().await; + assert!(!proxy.commit_observed.load(Ordering::SeqCst)); + assert_eq!( + mst2_retention_root::Entity::find() + .count(&second) + .await + .unwrap(), + 2 + ); + assert!(matches!( + repository + .inspect_prepare( + &second, + intent.operation_id(), + intent.manifest_digest(), + MetadataCommitPhase::Finalize + ) + .await + .unwrap(), + MetadataPrepareObservation::Committed(_) + )); +} diff --git a/src/jupiter/storage/native_publication_storage.rs b/src/jupiter/storage/native_publication_storage.rs index 2a5d0ff0..8039a3ea 100644 --- a/src/jupiter/storage/native_publication_storage.rs +++ b/src/jupiter/storage/native_publication_storage.rs @@ -281,7 +281,104 @@ impl MonoStorage { Ok(head) } - /// Maintenance-only, with stopped old writers and a drained queue; never invoked by resolve. + /// The command confirms stopped writers; this transaction checks the + /// paused admission gate, drained queue and exact native root. + pub(crate) async fn initialize_native_publication_for_maintenance( + &self, + instance: &str, + expected_root: &NativeRoot, + ) -> Result<(), PublicationReceiptError> { + let instance = canonical_instance(instance)?; + if uuid::Uuid::parse_str(&instance).is_ok_and(|id| id.is_nil()) { + return Err(integrity("native deployment instance must not be nil")); + } + for oid in [&expected_root.commit, &expected_root.tree] { + if oid.len() != 40 + || !oid + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) + { + return Err(integrity( + "maintenance requires canonical SHA-1 root identities", + )); + } + } + let txn = self.get_connection().begin().await?; + if !PushQueueStorage::try_mono_write_lock(&txn) + .await + .map_err(|error| integrity(&error.to_string()))? + { + return Err(PublicationReceiptError::Conflict( + "native initialization refused while a writer holds the mono write lock".into(), + )); + } + let control = txn + .query_one_raw(Statement::from_string( + txn.get_database_backend(), + "SELECT paused, hard_stopped FROM queue_control WHERE id=1 FOR UPDATE NOWAIT" + .to_owned(), + )) + .await? + .ok_or_else(|| { + integrity( + "queue_control row missing; maintenance cannot establish admission exclusion", + ) + })?; + if !control.try_get::("", "paused")? { + return Err(integrity( + "queue admission must be paused before native initialization", + )); + } + if control.try_get::("", "hard_stopped")? { + return Err(integrity( + "hard-stopped queue must be investigated before native initialization", + )); + } + let history = txn + .query_one_raw(Statement::from_string( + txn.get_database_backend(), + "SELECT EXISTS(SELECT 1 FROM mst2_native_publication) OR \ + EXISTS(SELECT 1 FROM mst2_publication WHERE native_certificate_version IS NOT NULL) AS present" + .to_owned(), + )) + .await? + .ok_or_else(|| integrity("native history observation missing"))?; + if history.try_get::("", "present")? { + return Err(integrity( + "native publication history exists; initialization cannot repair or replace it", + )); + } + let objects = txn + .query_one_raw(Statement::from_sql_and_values( + txn.get_database_backend(), + "SELECT (SELECT count(*) FROM mega_commit WHERE commit_id=$1) AS commits, \ + (SELECT max(tree) FROM mega_commit WHERE commit_id=$1) AS commit_tree, \ + (SELECT count(*) FROM mega_tree WHERE tree_id=$2) AS trees", + [ + expected_root.commit.clone().into(), + expected_root.tree.clone().into(), + ], + )) + .await? + .ok_or_else(|| integrity("native object observation missing"))?; + if objects.try_get::("", "commits")? != 1 + || objects + .try_get::>("", "commit_tree")? + .as_deref() + != Some(expected_root.tree.as_str()) + || objects.try_get::("", "trees")? != 1 + { + return Err(integrity( + "native root requires a unique stored commit with the expected tree and a unique stored tree", + )); + } + self.initialize_native_publication_in_txn(&txn, &instance, Some(expected_root)) + .await?; + txn.commit().await?; + Ok(()) + } + + #[cfg(test)] pub(crate) async fn initialize_native_publication( &self, instance: &str, @@ -293,13 +390,30 @@ impl MonoStorage { .map_err(|error| integrity(&error.to_string()))?; txn.execute_unprepared("SELECT id FROM queue_control WHERE id=1 FOR UPDATE") .await?; - let observation = observe(&txn).await?; + self.initialize_native_publication_in_txn(&txn, &instance, None) + .await?; + txn.commit().await?; + Ok(()) + } + + async fn initialize_native_publication_in_txn( + &self, + txn: &DatabaseTransaction, + instance: &str, + expected_root: Option<&NativeRoot>, + ) -> Result<(), PublicationReceiptError> { + let observation = observe(txn).await?; if observation.head.is_some() { return Err(integrity("native head already initialized")); } let root = observation .root .ok_or_else(|| integrity("native root missing"))?; + if expected_root.is_some_and(|expected| *expected != root) { + return Err(PublicationReceiptError::Conflict( + "native root differs from the maintenance command's expected commit/tree".into(), + )); + } let row = txn.query_one_raw(Statement::from_string(txn.get_database_backend(), "SELECT count(*)::bigint AS active FROM push_queue WHERE status IN ('Queued', 'Running')".to_owned())).await? .ok_or_else(|| integrity("queue observation missing"))?; @@ -331,7 +445,6 @@ impl MonoStorage { VALUES ('/', $1, $2, $3, $4, $5, 'INITIALIZING')", [instance.into(), floor.into(), NATIVE_EPOCH.into(), root.commit.into(), root.tree.into()], )).await?; - txn.commit().await?; Ok(()) } diff --git a/src/jupiter/storage/native_publication_tests.rs b/src/jupiter/storage/native_publication_tests.rs index a225eb09..c3e83323 100644 --- a/src/jupiter/storage/native_publication_tests.rs +++ b/src/jupiter/storage/native_publication_tests.rs @@ -475,3 +475,380 @@ pub(super) fn crash_checkpoint(phase: &str) { panic!("SIGKILL did not terminate the native worker"); } } + +async fn maintenance_fixture() -> (tempfile::TempDir, MonoStorage, PushQueueStorage, NativeRoot) { + use git_internal::{ + hash::HashKind, + internal::object::{ + blob::Blob, + commit::Commit, + tree::{Tree, TreeItem, TreeItemMode}, + }, + }; + + let (temp, mono, queue) = fixture().await; + mono.get_connection() + .execute_unprepared("DELETE FROM mst2_native_head") + .await + .unwrap(); + queue + .set_control_flags(Some(true), None, None) + .await + .unwrap(); + let blob = + Blob::from_content_bytes_with_kind(HashKind::Sha1, b"maintenance root".to_vec()).unwrap(); + let tree = Tree::from_tree_items_with_kind( + HashKind::Sha1, + vec![TreeItem::new( + TreeItemMode::Blob, + blob.id, + ".gitkeep".to_owned(), + )], + ) + .unwrap(); + let commit = + Commit::from_tree_id_with_kind(HashKind::Sha1, tree.id, vec![], "maintenance root") + .unwrap(); + mono.save_mega_trees(vec![tree.clone()], commit.id, None) + .await + .unwrap(); + mono.save_mega_commits(vec![commit.clone()], None) + .await + .unwrap(); + let root = NativeRoot { + commit: commit.id.to_string(), + tree: tree.id.to_string(), + }; + mono.get_connection().execute_raw(Statement::from_sql_and_values( + mono.get_connection().get_database_backend(), + "UPDATE mega_refs SET ref_commit_hash=$1,ref_tree_hash=$2 WHERE path='/' AND ref_name=$3 AND is_cl=false", + [root.commit.clone().into(), root.tree.clone().into(), MEGA_BRANCH_NAME.into()], + )).await.unwrap(); + (temp, mono, queue, root) +} + +async fn assert_no_maintenance_head(mono: &MonoStorage) { + assert_eq!( + mst2_native_head::Entity::find() + .count(mono.get_connection()) + .await + .unwrap(), + 0 + ); + assert_eq!( + mst2_native_publication::Entity::find() + .count(mono.get_connection()) + .await + .unwrap(), + 0 + ); + assert_eq!( + mst2_publication_outbox::Entity::find() + .count(mono.get_connection()) + .await + .unwrap(), + 0 + ); +} + +#[tokio::test] +async fn maintenance_initializes_only_a_floor_and_preserves_the_paused_gate() { + let (_temp, mono, queue, root) = maintenance_fixture().await; + mono.get_connection() + .execute_unprepared( + "INSERT INTO mst2_namespace_seq(namespace,sequence,epoch) VALUES('/older',23,1)", + ) + .await + .unwrap(); + mono.initialize_native_publication_for_maintenance(INSTANCE, &root) + .await + .unwrap(); + let before = mst2_native_head::Entity::find() + .one(mono.get_connection()) + .await + .unwrap() + .unwrap(); + assert_eq!(before.instance_id, INSTANCE); + assert_eq!(before.sequence, 23); + assert_eq!(before.writer_epoch, 1); + assert_eq!( + (&before.root_commit, &before.root_tree), + (&root.commit, &root.tree) + ); + assert_eq!(before.state, "INITIALIZING"); + assert_eq!(before.certificate_receipt_id, None); + assert!(queue.get_control().await.unwrap().paused); + assert!(mono.read_native_publication_head(INSTANCE).await.is_err()); + assert!( + mono.initialize_native_publication_for_maintenance(INSTANCE, &root) + .await + .is_err() + ); + assert!( + mono.initialize_native_publication_for_maintenance( + "11111111-2222-4333-8444-555555555555", + &root + ) + .await + .is_err() + ); + assert_eq!( + mst2_native_head::Entity::find() + .one(mono.get_connection()) + .await + .unwrap() + .unwrap(), + before + ); + assert_eq!( + selected_ref(mono.get_connection(), "/").await.unwrap(), + Some(root) + ); + assert_eq!( + mst2_native_publication::Entity::find() + .count(mono.get_connection()) + .await + .unwrap(), + 0 + ); + assert_eq!( + mst2_publication_outbox::Entity::find() + .count(mono.get_connection()) + .await + .unwrap(), + 0 + ); +} + +#[tokio::test] +async fn maintenance_rejects_unpaused_missing_or_hard_stopped_control_without_writes() { + for state in ["unpaused", "missing", "hard-stopped"] { + let (_temp, mono, queue, root) = maintenance_fixture().await; + match state { + "unpaused" => queue + .set_control_flags(Some(false), None, None) + .await + .unwrap(), + "hard-stopped" => queue + .set_control_flags(None, Some(true), None) + .await + .unwrap(), + "missing" => { + mono.get_connection() + .execute_unprepared("DELETE FROM queue_control") + .await + .unwrap(); + } + _ => unreachable!(), + } + assert!( + mono.initialize_native_publication_for_maintenance(INSTANCE, &root) + .await + .is_err(), + "{state}" + ); + assert_no_maintenance_head(&mono).await; + let control = crate::callisto::queue_control::Entity::find() + .one(mono.get_connection()) + .await + .unwrap(); + match state { + "missing" => assert!(control.is_none()), + "hard-stopped" => assert!(control.unwrap().hard_stopped), + _ => assert!(!control.unwrap().paused), + } + } +} + +#[tokio::test] +async fn maintenance_refuses_queued_and_running_work_until_drained() { + for status in ["Queued", "Running"] { + let (_temp, mono, queue, root) = maintenance_fixture().await; + queue + .set_control_flags(Some(false), None, None) + .await + .unwrap(); + let EnqueueOutcome::Inserted { id } = queue + .enqueue_atomic(EnqueueParams { + kind: PushQueueKindEnum::Push, + operation_id: "maintenance-pending", + path: "/project", + old_id: &"c".repeat(40), + new_id: &"e".repeat(40), + requester: None, + payload: serde_json::json!({"n":1,"commits":["e".repeat(40)]}), + }) + .await + .unwrap() + else { + panic!("fresh pending operation"); + }; + mono.get_connection() + .execute_raw(Statement::from_sql_and_values( + mono.get_connection().get_database_backend(), + "UPDATE push_queue SET status=$1::push_queue_status_enum WHERE id=$2", + [status.into(), id.into()], + )) + .await + .unwrap(); + queue + .set_control_flags(Some(true), None, None) + .await + .unwrap(); + assert!( + mono.initialize_native_publication_for_maintenance(INSTANCE, &root) + .await + .is_err(), + "{status}" + ); + assert_no_maintenance_head(&mono).await; + assert!(queue.get_control().await.unwrap().paused); + } +} + +#[tokio::test] +async fn maintenance_refuses_busy_writer_or_admission_locks_and_can_retry() { + for lock in ["writer", "admission"] { + let (_temp, mono, _queue, root) = maintenance_fixture().await; + let other = mono.get_connection().begin().await.unwrap(); + if lock == "writer" { + PushQueueStorage::acquire_mono_write_lock(&other) + .await + .unwrap(); + } else { + other + .execute_unprepared("SELECT id FROM queue_control WHERE id=1 FOR UPDATE") + .await + .unwrap(); + } + let result = tokio::time::timeout( + std::time::Duration::from_secs(5), + mono.initialize_native_publication_for_maintenance(INSTANCE, &root), + ) + .await; + assert!( + result + .expect("maintenance lock refusal must be bounded") + .is_err(), + "{lock}" + ); + assert_no_maintenance_head(&mono).await; + other.rollback().await.unwrap(); + mono.initialize_native_publication_for_maintenance(INSTANCE, &root) + .await + .unwrap(); + } +} + +#[tokio::test] +async fn maintenance_rejects_wrong_or_ambiguous_root_and_invalid_counters() { + for damage in [ + "commit", + "tree", + "missing", + "ambiguous", + "negative", + "exhausted", + ] { + let (_temp, mono, _queue, mut root) = maintenance_fixture().await; + match damage { + "commit" => root.commit = "e".repeat(40), + "tree" => root.tree = "e".repeat(40), + "missing" => { + mono.get_connection() + .execute_unprepared("DELETE FROM mega_refs WHERE path='/'") + .await + .unwrap(); + } + "ambiguous" => { + // Simulate a corrupted/old schema, independently of the normal + // uniqueness protection on selected refs. + mono.get_connection() + .execute_unprepared("DROP INDEX uniq_mref_path") + .await + .unwrap(); + mono.get_connection().execute_unprepared("INSERT INTO mega_refs(id,path,ref_name,ref_commit_hash,ref_tree_hash,created_at,updated_at,is_cl) SELECT 3,path,ref_name,ref_commit_hash,ref_tree_hash,created_at,updated_at,is_cl FROM mega_refs WHERE path='/'").await.unwrap(); + } + "negative" | "exhausted" => { + let floor = if damage == "negative" { + -1_i64 + } else { + i64::MAX + }; + mono.get_connection().execute_raw(Statement::from_sql_and_values( + mono.get_connection().get_database_backend(), + "INSERT INTO mst2_namespace_seq(namespace,sequence,epoch) VALUES('/older',$1,1)", [floor.into()], + )).await.unwrap(); + } + _ => unreachable!(), + } + assert!( + mono.initialize_native_publication_for_maintenance(INSTANCE, &root) + .await + .is_err(), + "{damage}" + ); + assert_no_maintenance_head(&mono).await; + } +} + +#[tokio::test] +async fn maintenance_rejects_nil_instance_and_non_sha1_expected_roots() { + let (_temp, mono, _queue, root) = maintenance_fixture().await; + for instance in ["invalid", "00000000-0000-0000-0000-000000000000"] { + assert!( + mono.initialize_native_publication_for_maintenance(instance, &root) + .await + .is_err() + ); + } + for commit in ["a".repeat(64), "A".repeat(40)] { + let wrong = NativeRoot { + commit, + tree: root.tree.clone(), + }; + assert!( + mono.initialize_native_publication_for_maintenance(INSTANCE, &wrong) + .await + .is_err() + ); + } + assert_no_maintenance_head(&mono).await; +} + +#[tokio::test] +async fn maintenance_refuses_missing_commit_tree_or_a_mismatched_commit_tree() { + for damage in ["commit", "tree", "wrong-tree"] { + let (_temp, mono, _queue, root) = maintenance_fixture().await; + let sql = match damage { + "commit" => "DELETE FROM mega_commit WHERE commit_id=$1", + "tree" => "DELETE FROM mega_tree WHERE tree_id=$1", + "wrong-tree" => "UPDATE mega_commit SET tree=repeat('e',40) WHERE commit_id=$1", + _ => unreachable!(), + }; + let id = if damage == "tree" { + &root.tree + } else { + &root.commit + }; + mono.get_connection() + .execute_raw(Statement::from_sql_and_values( + mono.get_connection().get_database_backend(), + sql, + [id.clone().into()], + )) + .await + .unwrap(); + assert!( + mono.initialize_native_publication_for_maintenance(INSTANCE, &root) + .await + .is_err(), + "{damage}" + ); + assert_no_maintenance_head(&mono).await; + assert_eq!( + selected_ref(mono.get_connection(), "/").await.unwrap(), + Some(root) + ); + } +} diff --git a/tests/integration_git_cli.rs b/tests/integration_git_cli.rs index 41e5b509..8f35c974 100644 --- a/tests/integration_git_cli.rs +++ b/tests/integration_git_cli.rs @@ -20,7 +20,6 @@ use std::{ collections::BTreeMap, fs, io::{Read, Write}, - net::TcpStream, path::{Path, PathBuf}, process::{Child, Command, ExitStatus, Stdio}, sync::atomic::{AtomicUsize, Ordering}, @@ -204,34 +203,51 @@ impl ServiceProcess { service } - fn wait_until_listening( + fn wait_until_openapi_ready( &mut self, port: u16, timeout: Duration, stdout_path: &Path, stderr_path: &Path, ) { + let client = reqwest::blocking::Client::builder() + .timeout(Duration::from_secs(5)) + .redirect(reqwest::redirect::Policy::none()) + .build() + .expect("build HTTP readiness client"); + let url = format!("http://127.0.0.1:{port}/api/openapi.json"); let deadline = Instant::now() + timeout; loop { - if TcpStream::connect(("127.0.0.1", port)).is_ok() { - return; - } if let Some(status) = self.child.try_wait().expect("poll service") { self.reaped = true; panic!( - "service exited before binding port {port} (status {status})\nstdout:\n{}\nstderr:\n{}", + "service exited before OpenAPI was ready on port {port} (status {status})\nstdout:\n{}\nstderr:\n{}", read_log(stdout_path), read_log(stderr_path), ); } - if Instant::now() >= deadline { + let remaining = deadline.saturating_duration_since(Instant::now()); + if remaining.is_zero() { panic!( - "service did not bind port {port} within {timeout:?}\nstdout:\n{}\nstderr:\n{}", + "service OpenAPI not ready on port {port} within {timeout:?}\nstdout:\n{}\nstderr:\n{}", read_log(stdout_path), read_log(stderr_path), ); } - sleep(Duration::from_millis(200)); + if let Ok(response) = client + .get(&url) + .timeout(remaining.min(Duration::from_secs(5))) + .send() + && response.status() == reqwest::StatusCode::OK + && Instant::now() <= deadline + { + return; + } + sleep( + deadline + .saturating_duration_since(Instant::now()) + .min(Duration::from_millis(200)), + ); } } @@ -345,7 +361,7 @@ fn boot_service_http_with_env( .stderr(Stdio::from(create_log_file(&stderr_path))); let mut service = ServiceProcess::spawn(command); - service.wait_until_listening(port, Duration::from_secs(90), &stdout_path, &stderr_path); + service.wait_until_openapi_ready(port, Duration::from_secs(90), &stdout_path, &stderr_path); (service, port, stdout_path, stderr_path) }