Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion src/jupiter/storage/native_metadata_install.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1056,7 +1056,7 @@ fn unavailable(message: &str) -> SnapshotError {

#[path = "native_metadata_install_capability.rs"]
mod capability;
pub(crate) use capability::ValidatedLegacyInstallCapability;
pub(crate) use capability::{LegacyPayloadInstallWork, ValidatedLegacyInstallCapability};

#[cfg(test)]
#[path = "native_metadata_install_capability_tests.rs"]
Expand Down
200 changes: 197 additions & 3 deletions src/jupiter/storage/native_metadata_install_capability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,50 @@ pub(crate) struct ValidatedLegacyInstallCapability {
install_seal: [u8; 32],
}

/// Payload batch work only, excluding barriers, intent/capability setup,
/// finalize and the complete COMMITTED replay oracle. Classification reuses
/// the requested-member query; its batch count is not an extra SQL trip.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub(crate) struct LegacyPayloadInstallWork {
pub requested_pages: u64,
pub payload_pages_omitted: u64,
pub payload_pages_encoded: u64,
pub requested_payload_bytes_validated: u64,
pub payload_bytes_omitted: u64,
pub payload_bytes_encoded: u64,
pub metadata_parameter_bytes: u64,
pub payload_parameter_bytes: u64,
pub transactions: u64,
pub registration_queries: u64,
pub requested_member_queries: u64,
pub classification_batches: u64,
pub insert_statements: u64,
pub byte_comparison_queries: u64,
pub committed_replay_pages: u64,
pub elapsed_micros: u64,
}

impl LegacyPayloadInstallWork {
pub(crate) fn record(&mut self, batch: Self) {
self.requested_pages += batch.requested_pages;
self.payload_pages_omitted += batch.payload_pages_omitted;
self.payload_pages_encoded += batch.payload_pages_encoded;
self.requested_payload_bytes_validated += batch.requested_payload_bytes_validated;
self.payload_bytes_omitted += batch.payload_bytes_omitted;
self.payload_bytes_encoded += batch.payload_bytes_encoded;
self.metadata_parameter_bytes += batch.metadata_parameter_bytes;
self.payload_parameter_bytes += batch.payload_parameter_bytes;
self.transactions += batch.transactions;
self.registration_queries += batch.registration_queries;
self.requested_member_queries += batch.requested_member_queries;
self.classification_batches += batch.classification_batches;
self.insert_statements += batch.insert_statements;
self.byte_comparison_queries += batch.byte_comparison_queries;
self.committed_replay_pages += batch.committed_replay_pages;
self.elapsed_micros += batch.elapsed_micros;
}
}

impl PostgresMetadataInstallRepository {
pub(crate) async fn mint_legacy_install_capability(
&self,
Expand Down Expand Up @@ -173,6 +217,144 @@ impl PostgresMetadataInstallRepository {
.await
}

/// Classify and install in one existing payload transaction. Presence is
/// only a write hint; finalize still reads and validates every stored byte.
pub(crate) async fn install_missing_pages_validated(
&self,
capability: &ValidatedLegacyInstallCapability,
payloads: &[MetadataPagePayload],
) -> Result<LegacyPayloadInstallWork, MetadataInstallError> {
let started = std::time::Instant::now();
if capability.scope != self.storage_scope {
return Err(
integrity("install capability belongs to another captured primary scope").into(),
);
}
if payloads.is_empty() || payloads.len() > 64 {
return Err(SnapshotError::new(
SnapshotErrorCode::LimitExceeded,
"metadata installation batch must contain 1..=64 pages",
)
.into());
}
let mut ids = BTreeSet::new();
for payload in payloads {
validate_payload(payload)?;
if !ids.insert(payload.id) {
return Err(integrity("duplicate page in metadata installation batch").into());
}
if capability
.members
.get(&payload.id)
.map(|member| member.0 as u64)
!= Some(payload.size)
{
return Err(integrity(
"metadata payload is not a member of its validated installation",
)
.into());
}
}
let metadata = serde_json::to_string(
&payloads
.iter()
.map(|page| json!({"page_id":hex::encode(page.id),"size":page.size}))
.collect::<Vec<_>>(),
)
.map_err(internal)?;
let txn = self.transaction().await?;
let result = async {
self.capability_barrier(&txn).await?;
let state = read_registered_prepare(&txn, capability).await?;
let rows = read_requested_members(&txn, capability, &metadata, payloads.len()).await?;
let mut present = BTreeSet::new();
for row in rows {
if !row.try_get::<bool>("", "payload_present").map_err(internal)? {
continue;
}
if row.try_get::<Option<i16>>("", "payload_codec").map_err(internal)?
!= Some(capability.identity.metadata_codec as i16)
|| row.try_get::<Option<i32>>("", "payload_size").map_err(internal)?
!= row.try_get::<Option<i32>>("", "expected_size").map_err(internal)?
{
return Err(integrity("immutable metadata payload profile conflicts with its fixed member"));
}
let page: Vec<u8> = row.try_get("", "page_id").map_err(internal)?;
present.insert(<[u8; 32]>::try_from(page.as_slice()).map_err(internal)?);
}
let committed = state == "COMMITTED";
if committed {
// A committed replay never repairs missing or corrupt bytes.
let stored = require_plan(&txn, &capability.intent).await?;
load_installed_dag(&txn, &stored).await?;
check_payload_coverage(&txn, &stored).await?;
verify_graph(&txn, &stored).await?;
}
let submitted: Vec<_> = payloads
.iter()
.filter(|page| committed || !present.contains(&page.id))
.collect();
let requested_bytes = payloads.iter().map(|page| page.size).sum::<u64>();
let submitted_bytes = submitted.iter().map(|page| page.size).sum::<u64>();
let mut work = LegacyPayloadInstallWork {
requested_pages: payloads.len() as u64,
payload_pages_omitted: (payloads.len() - submitted.len()) as u64,
payload_pages_encoded: submitted.len() as u64,
requested_payload_bytes_validated: requested_bytes,
payload_bytes_omitted: requested_bytes - submitted_bytes,
payload_bytes_encoded: submitted_bytes,
metadata_parameter_bytes: metadata.len() as u64,
transactions: 1,
registration_queries: 1,
requested_member_queries: 1,
classification_batches: u64::from(!committed),
committed_replay_pages: if committed { payloads.len() as u64 } else { 0 },
..LegacyPayloadInstallWork::default()
};
if !submitted.is_empty() {
let pages: Vec<_> = submitted.iter().map(|page| json!({
"page_id":hex::encode(page.id),"size":page.size,"payload":hex::encode(&page.bytes)
})).collect();
let encoded = serde_json::to_string(&pages).map_err(internal)?;
let encoded_len = encoded.len() as u64;
if !committed {
txn.execute_raw(statement(
"INSERT INTO mst2_metadata_payload(page_id,metadata_codec,byte_size,payload)
SELECT decode(p.page_id,'hex'),$1,p.size,decode(p.payload,'hex')
FROM jsonb_to_recordset($2::jsonb) AS p(page_id text,size integer,payload text)
ON CONFLICT(page_id) DO NOTHING",
[(capability.identity.metadata_codec as i16).into(),encoded.clone().into()],
)).await.map_err(internal)?;
work.insert_statements = 1;
work.payload_parameter_bytes += encoded_len;
}
let bad = txn.query_one_raw(statement(
"SELECT p.page_id FROM jsonb_to_recordset($2::jsonb) AS p(page_id text,size integer,payload text)
LEFT JOIN mst2_metadata_payload b ON b.page_id=decode(p.page_id,'hex')
WHERE b.page_id IS NULL OR b.metadata_codec<>$1 OR b.byte_size<>p.size
OR b.payload<>decode(p.payload,'hex') LIMIT 1",
[(capability.identity.metadata_codec as i16).into(),encoded.into()],
)).await.map_err(internal)?;
work.byte_comparison_queries = 1;
work.payload_parameter_bytes += encoded_len;
if bad.is_some() {
return Err(integrity("immutable metadata payload identity conflicts with stored bytes"));
}
}
Ok(work)
}.await;
let mut work = commit(
txn,
result,
&capability.intent.operation_id,
capability.intent.manifest_digest,
MetadataCommitPhase::Payload,
)
.await?;
work.elapsed_micros = started.elapsed().as_micros().min(u64::MAX as u128) as u64;
Ok(work)
}

pub(super) async fn capability_barrier(
&self,
txn: &DatabaseTransaction,
Expand Down Expand Up @@ -302,9 +484,21 @@ async fn check_requested_members(
encoded: &str,
expected_count: usize,
) -> Result<(), SnapshotError> {
read_requested_members(txn, capability, encoded, expected_count)
.await
.map(|_| ())
}

async fn read_requested_members(
txn: &DatabaseTransaction,
capability: &ValidatedLegacyInstallCapability,
encoded: &str,
expected_count: usize,
) -> Result<Vec<QueryResult>, SnapshotError> {
let rows=txn.query_all_raw(statement(
"SELECT decode(p.page_id,'hex') AS page_id,m.expected_size,m.generation AS member_generation,
b.page_id IS NOT NULL AS payload_present,b.generation AS payload_generation,
b.metadata_codec AS payload_codec,b.byte_size AS payload_size,
c.generation AS current_generation,l.generation AS lifetime_generation,l.graph_domain,
l.state AS lifetime_state,l.metadata_codec AS lifetime_codec,l.expected_size AS lifetime_size,
n.state AS graph_state,n.kind AS graph_kind,n.bytes AS graph_bytes,
Expand All @@ -321,7 +515,7 @@ async fn check_requested_members(
if rows.len() != expected_count {
return Err(integrity("requested install membership is incomplete"));
}
for row in rows {
for row in &rows {
let page: Vec<u8> = row.try_get("", "page_id").map_err(internal)?;
let page: [u8; 32] = page.as_slice().try_into().map_err(internal)?;
let size = row
Expand All @@ -336,12 +530,12 @@ async fn check_requested_members(
));
}
check_physical_member(
&row,
row,
capability.identity.metadata_codec as i16,
size.ok_or_else(|| integrity("requested member is missing"))?,
)?;
}
Ok(())
Ok(rows)
}

fn check_physical_member(row: &QueryResult, codec: i16, size: i32) -> Result<(), SnapshotError> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -961,3 +961,6 @@ async fn install_capability_batch_fault_rolls_back_and_raw_writer_waits_before_r
);
install(&repository, &cap, &pages).await;
}

#[path = "native_metadata_install_missing_tests.rs"]
mod missing;
Loading
Loading