use std::collections::{BTreeMap, BTreeSet};
use prikk_error::{PrikkError, Result};
use prikk_object::{
BlockKind, BlockPayload, CanonicalEncode, ObjectEnvelope, ObjectId, ObjectType,
RecognitionClaimPayload, RefKind, RefStatePayload, RefUpdatePayload,
};
use crate::block_state::{CandidateStateDerivationError, derive_next_state_root_for_candidate};
use crate::container::decode_container_records;
use crate::fsutil::read_file_if_exists;
use crate::layout::{ContainerSlot, RepositoryLayout, persisted_object_types};
use crate::lifecycle_cache::replay::LifecycleReplayError;
use crate::lock::ActiveLock;
use crate::maintainer_signing::{MaintainerSigner, maintainer_signature};
use crate::object_store::{ObjectReadSnapshot, ObjectReader, ObjectWriteSession, ObjectWriter};
use crate::patch_exchange::accepted_but_unsealed_patch_ids;
use crate::recognition_claim::maintainer_trust_policy_or_empty;
use crate::recognition_claim::{ClaimSignatureVerification, verify_claim_signature};
use crate::refs::{RefPublication, RefStore, validate_local_branch_ref};
use crate::trust::{GatedOperation, verify_signer_trusted};
use crate::wal::Wal;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SealFromAcceptedOutcome {
Sealed {
ref_name: String,
block_id: ObjectId,
ref_state_id: ObjectId,
patch_count: usize,
claim_signature_outcome: ClaimSignatureVerification,
},
AlreadySealed {
ref_name: String,
claim_id: ObjectId,
},
}
pub fn seal_from_accepted_claim(
layout: &RepositoryLayout,
ref_name: &str,
claim_id: ObjectId,
signer: &impl MaintainerSigner,
) -> Result<SealFromAcceptedOutcome> {
layout.require_current_format()?;
let canonical_ref = validate_local_branch_ref(ref_name)?;
crate::refs::ensure_no_incomplete_publication(layout)?;
let read_snapshot = ObjectReadSnapshot::open(layout)?;
let claim_envelope = read_snapshot
.read_typed(claim_id, ObjectType::RecognitionClaim)?
.ok_or_else(|| {
PrikkError::Integrity(format!("recognition claim {claim_id} does not exist"))
})?;
let claim = RecognitionClaimPayload::decode_canonical(&claim_envelope.canonical_payload)?;
for &patch_id in &claim.patch_ids {
if !read_snapshot.contains_object(ObjectType::Patch, patch_id) {
return Err(PrikkError::Integrity(format!(
"recognition claim {claim_id} names patch {patch_id}, which does not exist -- \
refusing the whole seal, no partial application"
)));
}
}
let unsealed: BTreeSet<ObjectId> = accepted_but_unsealed_patch_ids(layout)?
.into_iter()
.collect();
let selected_patch_ids: Vec<ObjectId> = claim
.patch_ids
.iter()
.copied()
.filter(|patch_id| unsealed.contains(patch_id))
.collect();
if selected_patch_ids.is_empty() {
return Ok(SealFromAcceptedOutcome::AlreadySealed {
ref_name: canonical_ref,
claim_id,
});
}
refuse_if_order_ambiguous(layout, claim_id, &claim)?;
verify_signer_trusted(layout, signer, GatedOperation::SyncSeal)?;
let trust_policy = maintainer_trust_policy_or_empty(layout)?;
let claim_signature_outcome = verify_claim_signature(&claim_envelope, &trust_policy)?;
let active_lock = ActiveLock::acquire(layout)?;
crate::refs::ensure_no_incomplete_publication(layout)?;
let wal = Wal::for_layout(layout);
let replay = wal.replay()?;
if replay.trailing_partial_bytes != 0 {
return Err(PrikkError::Integrity(format!(
"active WAL has {} trailing partial bytes; run verify/doctor before sealing from \
an accepted claim",
replay.trailing_partial_bytes
)));
}
if replay.has_item_failure() {
return Err(PrikkError::Integrity(
"active WAL has a damaged record; run verify/doctor before sealing from an accepted \
claim"
.to_string(),
));
}
if !replay.records.is_empty() {
return Err(PrikkError::LockConflict(
"sealing from an accepted claim requires an empty active WAL -- seal or discard \
local work first"
.to_string(),
));
}
let ref_store = RefStore::new(layout.clone());
let current = read_current_tip(&read_snapshot, &ref_store, &canonical_ref)?;
let parent = current.as_ref().map(|tip| tip.target_block_id);
let state_merkle_root =
match derive_next_state_root_for_candidate(&read_snapshot, parent, &selected_patch_ids) {
Ok(root) => root,
Err(CandidateStateDerivationError::Lineage(err)) => return Err(err),
Err(CandidateStateDerivationError::Patch(err)) => {
return Err(classify_patch_application_failure(err));
}
};
let block_payload = BlockPayload {
parent_block_ids: parent.into_iter().collect(),
kind: if current.is_some() {
BlockKind::Normal
} else {
BlockKind::Root
},
patch_ids: selected_patch_ids.clone(),
state_merkle_root,
snapshot_blob_ref: None,
mainline_parent_id: None,
merge_baseline_block_id: None,
};
let block_envelope = signed_envelope(
ObjectType::Block,
2,
block_payload.to_canonical_bytes()?,
signer,
)?;
let block_id = block_envelope.object_id();
let mut object_store = ObjectWriteSession::open(layout)?;
object_store.write_object(&block_envelope)?;
let update_seq = current.as_ref().map_or(1, |tip| tip.update_seq + 1);
let previous_ref_state_id = current.as_ref().map(|tip| tip.ref_state_id);
let ref_state_payload = RefStatePayload {
ref_name: canonical_ref.clone(),
kind: RefKind::Branch,
target_object_id: block_id,
update_seq,
previous_ref_state_id,
required_attestation_ids: Vec::new(),
closed: false,
};
let ref_state_envelope = signed_envelope(
ObjectType::RefState,
1,
ref_state_payload.to_canonical_bytes()?,
signer,
)?;
let ref_state_id = ref_state_envelope.object_id();
let ref_update_payload = RefUpdatePayload {
ref_name: canonical_ref.clone(),
old_ref_state_id: previous_ref_state_id,
new_ref_state_id: ref_state_id,
new_target_object_id: block_id,
update_seq,
created_at: 0,
author_key_id: signer.key_id().to_string(),
};
let ref_update_envelope = signed_envelope(
ObjectType::RefUpdate,
1,
ref_update_payload.to_canonical_bytes()?,
signer,
)?;
let published_ref_state_id = ref_store.publish_with_object_store(
&mut object_store,
&RefPublication {
ref_name: canonical_ref.clone(),
expected_previous_ref_state_id: previous_ref_state_id,
ref_state: ref_state_envelope,
ref_update: ref_update_envelope,
},
)?;
drop(active_lock);
Ok(SealFromAcceptedOutcome::Sealed {
ref_name: canonical_ref,
block_id,
ref_state_id: published_ref_state_id,
patch_count: selected_patch_ids.len(),
claim_signature_outcome,
})
}
fn classify_patch_application_failure(error: LifecycleReplayError) -> PrikkError {
match error {
LifecycleReplayError::InconsistentLifecycleEffect { .. }
| LifecycleReplayError::TextSpanResolutionFailed { .. } => PrikkError::Integrity(format!(
"seal refused: divergence -- an accepted patch did not apply cleanly to this \
repository's own tip ({error}); the two histories moved differently, nothing is \
corrupt, and resolving this is merge's job, not this operation's"
)),
LifecycleReplayError::MissingBlockInLineage { .. }
| LifecycleReplayError::UnreadableBlockInLineage { .. }
| LifecycleReplayError::MergeLineageUnsupported { .. }
| LifecycleReplayError::LineageCycle { .. }
| LifecycleReplayError::HorizonNotInLineage { .. }
| LifecycleReplayError::MalformedPatchInLineage { .. }
| LifecycleReplayError::MissingBlobForLifecycleEffect { .. } => PrikkError::Integrity(
format!("seal refused: integrity -- this repository's own state is broken ({error})"),
),
}
}
struct CurrentTip {
ref_state_id: ObjectId,
target_block_id: ObjectId,
update_seq: u64,
}
fn read_current_tip(
object_store: &impl ObjectReader,
ref_store: &RefStore,
ref_name: &str,
) -> Result<Option<CurrentTip>> {
let Some(ref_state_id) = ref_store.read_current_ref_state_id(ref_name)? else {
return Ok(None);
};
let envelope = object_store
.read_typed(ref_state_id, ObjectType::RefState)?
.ok_or_else(|| {
PrikkError::Integrity(format!(
"ref {ref_name} points to missing RefState {ref_state_id}"
))
})?;
let payload =
RefStatePayload::decode_canonical(&envelope.canonical_payload, envelope.schema_version)?;
if payload.ref_name != ref_name {
return Err(PrikkError::Integrity(format!(
"RefState name mismatch for {ref_name}: got {}",
payload.ref_name
)));
}
if object_store
.read_typed(payload.target_object_id, ObjectType::Block)?
.is_none()
{
return Err(PrikkError::Integrity(format!(
"ref {ref_name} targets missing block {}",
payload.target_object_id
)));
}
Ok(Some(CurrentTip {
ref_state_id,
target_block_id: payload.target_object_id,
update_seq: payload.update_seq,
}))
}
fn signed_envelope(
object_type: ObjectType,
schema_version: u32,
canonical_payload: Vec<u8>,
signer: &impl MaintainerSigner,
) -> Result<ObjectEnvelope> {
let mut envelope = ObjectEnvelope::unsigned(object_type, schema_version, canonical_payload);
let object_id = envelope.object_id();
envelope.add_signature(maintainer_signature(signer, object_type, object_id)?)?;
Ok(envelope)
}
fn refuse_if_order_ambiguous(
layout: &RepositoryLayout,
selected_claim_id: ObjectId,
selected: &RecognitionClaimPayload,
) -> Result<()> {
let selected_set: BTreeSet<ObjectId> = selected.patch_ids.iter().copied().collect();
for (other_id, other) in enumerate_stored_claims(layout)? {
if other_id == selected_claim_id {
continue;
}
let overlaps = other.patch_ids.iter().any(|id| selected_set.contains(id));
if !overlaps {
continue;
}
if !orders_agree_on_overlap(&selected.patch_ids, &other.patch_ids) {
return Err(PrikkError::Integrity(format!(
"recognition claims {selected_claim_id} and {other_id} name overlapping patches \
but disagree on their relative order -- refusing rather than guessing which is \
right"
)));
}
}
Ok(())
}
fn orders_agree_on_overlap(a: &[ObjectId], b: &[ObjectId]) -> bool {
let mut b_index: BTreeMap<ObjectId, usize> = BTreeMap::new();
for (index, id) in b.iter().enumerate() {
b_index.entry(*id).or_insert(index);
}
let mut last_seen: Option<usize> = None;
for id in a {
if let Some(&index) = b_index.get(id) {
if let Some(last) = last_seen {
if index < last {
return false;
}
}
last_seen = Some(index);
}
}
true
}
fn enumerate_stored_claims(
layout: &RepositoryLayout,
) -> Result<Vec<(ObjectId, RecognitionClaimPayload)>> {
debug_assert!(
persisted_object_types().contains(&ObjectType::RecognitionClaim),
"RecognitionClaim must remain a persisted, containerized object type"
);
let container_path = layout.container_slot_path(ObjectType::RecognitionClaim, ContainerSlot::A);
let relative = layout.repository_relative(&container_path)?;
let mut claims = Vec::new();
if let Some(bytes) = read_file_if_exists(layout.repository_mutation_root(), &relative)? {
let replay = decode_container_records(ObjectType::RecognitionClaim, &bytes)?;
for record in replay.records {
let payload =
RecognitionClaimPayload::decode_canonical(&record.envelope.canonical_payload)?;
claims.push((record.envelope.object_id(), payload));
}
}
Ok(claims)
}
#[cfg(test)]
mod tests;