use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
command::{
CommandEnvelope, CommandMetadata, CommandScope, FencingToken, Precondition, ResourceBounds,
},
digest::ContentDigest,
error::StateError,
id::{AggregateId, Audience, CommandId, NamespaceId, PartitionId, Purpose},
journal::{
self, CommitJournalBatch, DestroyPartition, ExcisePartitionRecords, GetJournalHead,
GetJournalProof, GetJournalRoot, JournalAnchor, JournalAttestation, JournalInclusionProof,
JournalProof, JournalRange, JournalRecord, JournalRecordDraft, ListJournalPartitions,
PartitionDescriptor, PartitionPage, QuarantinedRecord, ReadJournalRange, RecordDecision,
RecordDisposition, RecordKind, RecordTrust, RepairPartition,
},
revision::{CommitRoot, JournalHead, JournalPosition},
};
use crate::{
journal::verify::{PartitionVerification, ViolationCategory},
wire::{Kernel, completeness, fixed_bytes, known, malformed, required},
};
impl From<Kernel<RecordTrust>> for pb::RecordTrust {
fn from(value: Kernel<RecordTrust>) -> Self {
match value.0 {
RecordTrust::Unspecified => Self::RECORD_TRUST_UNSPECIFIED,
RecordTrust::TrustedUser => Self::RECORD_TRUST_TRUSTED_USER,
RecordTrust::QuarantinedContent => Self::RECORD_TRUST_QUARANTINED_CONTENT,
}
}
}
fn record_trust(
field: &str,
value: buffa::EnumValue<pb::RecordTrust>,
) -> Result<RecordTrust, StateError> {
match known(field, value)? {
pb::RecordTrust::RECORD_TRUST_UNSPECIFIED => Ok(RecordTrust::Unspecified),
pb::RecordTrust::RECORD_TRUST_TRUSTED_USER => Ok(RecordTrust::TrustedUser),
pb::RecordTrust::RECORD_TRUST_QUARANTINED_CONTENT => Ok(RecordTrust::QuarantinedContent),
}
}
impl From<Kernel<&JournalRecordDraft>> for pb::JournalRecordDraft {
fn from(value: Kernel<&JournalRecordDraft>) -> Self {
Self {
kind: value.0.kind().as_str().to_owned(),
trust: pb::RecordTrust::from(Kernel(value.0.trust())).into(),
payload: value.0.payload().to_vec(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::JournalRecordDraft> for Kernel<JournalRecordDraft> {
type Error = StateError;
fn try_from(value: pb::JournalRecordDraft) -> Result<Self, Self::Error> {
let kind = RecordKind::new(value.kind);
if kind.is_empty() {
return Err(malformed(
"kind",
"a record declares the schema its payload follows",
));
}
let trust = record_trust("trust", value.trust)?;
Ok(Self(
JournalRecordDraft::new(kind, value.payload).with_trust(trust),
))
}
}
impl From<Kernel<&JournalRecord>> for pb::JournalRecord {
fn from(value: Kernel<&JournalRecord>) -> Self {
use polyc_state::page::Positioned as _;
Self {
position: value.0.position().get(),
kind: value.0.kind().as_str().to_owned(),
trust: pb::RecordTrust::from(Kernel(value.0.trust())).into(),
payload: value.0.payload().to_vec(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::JournalRecord> for Kernel<JournalRecord> {
type Error = StateError;
fn try_from(value: pb::JournalRecord) -> Result<Self, Self::Error> {
let trust = record_trust("trust", value.trust)?;
Ok(Self(JournalRecord::new(
JournalPosition::new(value.position),
RecordKind::new(value.kind),
trust,
value.payload,
)))
}
}
const fn journal_bounds() -> ResourceBounds {
ResourceBounds::new(
journal::MAX_BATCH_PAYLOAD_BYTES,
journal::MAX_RECORDS_PER_BATCH,
)
}
struct WireMetadata {
command_id: String,
partition: String,
aggregate: String,
namespace: String,
purpose: String,
audience: String,
digest: Vec<u8>,
precondition: Option<pb::Precondition>,
fence: Option<u64>,
}
impl TryFrom<WireMetadata> for Kernel<CommandMetadata> {
type Error = StateError;
fn try_from(wire: WireMetadata) -> Result<Self, Self::Error> {
let digest = ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
"digest",
&wire.digest,
)?);
let precondition = wire.precondition.ok_or_else(|| {
malformed(
"precondition",
"a precondition names what durable state the command requires",
)
})?;
let precondition = Kernel::<Precondition>::try_from(precondition)?.into_inner();
let mut metadata = CommandMetadata::new(
CommandId::new(wire.command_id),
journal::family(),
digest,
CommandScope::new(
AggregateId::new(wire.aggregate),
PartitionId::new(wire.partition),
NamespaceId::new(wire.namespace),
),
CommandEnvelope::new(
Purpose::new(wire.purpose),
Audience::new(wire.audience),
journal_bounds(),
),
)
.with_precondition(precondition);
if let Some(fence) = wire.fence {
metadata = metadata.with_fence(FencingToken::new(fence));
}
Ok(Self(metadata))
}
}
impl From<Kernel<&CommitJournalBatch>> for pb::CommitJournalBatch {
fn from(value: Kernel<&CommitJournalBatch>) -> Self {
let metadata = value.0.metadata();
Self {
command_id: metadata.command_id().as_str().to_owned(),
partition: metadata.scope().partition().as_str().to_owned(),
aggregate: metadata.scope().aggregate().as_str().to_owned(),
namespace: metadata.scope().namespace().as_str().to_owned(),
purpose: metadata.envelope().purpose().as_str().to_owned(),
audience: metadata.envelope().audience().as_str().to_owned(),
digest: metadata.digest().as_bytes().to_vec(),
precondition: buffa::MessageField::some(pb::Precondition::from(Kernel(
metadata.precondition(),
))),
fence: metadata.fence().map(FencingToken::get),
records: value
.0
.records()
.iter()
.map(|record| pb::JournalRecordDraft::from(Kernel(record)))
.collect(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::CommitJournalBatch> for Kernel<CommitJournalBatch> {
type Error = StateError;
fn try_from(value: pb::CommitJournalBatch) -> Result<Self, Self::Error> {
let metadata = Kernel::<CommandMetadata>::try_from(WireMetadata {
command_id: value.command_id,
partition: value.partition,
aggregate: value.aggregate,
namespace: value.namespace,
purpose: value.purpose,
audience: value.audience,
digest: value.digest,
precondition: value.precondition.into_option(),
fence: value.fence,
})?
.into_inner();
let records = value
.records
.into_iter()
.map(|record| Kernel::<JournalRecordDraft>::try_from(record).map(Kernel::into_inner))
.collect::<Result<Vec<_>, _>>()?;
Ok(Self(CommitJournalBatch::new(metadata, records)))
}
}
impl From<Kernel<&JournalAnchor>> for pb::JournalAnchor {
fn from(value: Kernel<&JournalAnchor>) -> Self {
Self {
partition: value.0.partition().as_str().to_owned(),
head: value.0.head().get(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl From<pb::JournalAnchor> for Kernel<JournalAnchor> {
fn from(value: pb::JournalAnchor) -> Self {
Self(JournalAnchor::new(
PartitionId::new(value.partition),
JournalPosition::new(value.head),
))
}
}
impl From<Kernel<&ReadJournalRange>> for pb::ReadJournalRange {
fn from(value: Kernel<&ReadJournalRange>) -> Self {
Self {
partition: value.0.partition().as_str().to_owned(),
start: value.0.start().get(),
limit: value.0.limit(),
anchor: value
.0
.anchor()
.map_or_else(buffa::MessageField::default, |anchor| {
buffa::MessageField::some(pb::JournalAnchor::from(Kernel(anchor)))
}),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl From<pb::ReadJournalRange> for Kernel<ReadJournalRange> {
fn from(value: pb::ReadJournalRange) -> Self {
let request = ReadJournalRange::new(
PartitionId::new(value.partition),
JournalPosition::new(value.start),
value.limit,
);
Self(match value.anchor.into_option() {
Some(anchor) => request.within(Kernel::<JournalAnchor>::from(anchor).into_inner()),
None => request,
})
}
}
impl From<Kernel<&JournalRange>> for pb::JournalRange {
fn from(value: Kernel<&JournalRange>) -> Self {
Self {
anchor: buffa::MessageField::some(pb::JournalAnchor::from(Kernel(value.0.anchor()))),
records: value
.0
.records()
.iter()
.map(|record| pb::JournalRecord::from(Kernel(record)))
.collect(),
next_position: value.0.next_position().map(JournalPosition::get),
completeness: pb::PageCompleteness::from(Kernel(value.0.completeness())).into(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::JournalRange> for Kernel<JournalRange> {
type Error = StateError;
fn try_from(value: pb::JournalRange) -> Result<Self, Self::Error> {
let anchor = Kernel::<JournalAnchor>::from(required(
"anchor",
"a bounded read names the immutable prefix it observed",
value.anchor,
)?)
.into_inner();
let completeness = completeness("completeness", value.completeness)?;
let records = value
.records
.into_iter()
.map(|record| Kernel::<JournalRecord>::try_from(record).map(Kernel::into_inner))
.collect::<Result<Vec<_>, _>>()?;
Ok(Self(JournalRange::new(
anchor,
records,
value.next_position.map(JournalPosition::new),
completeness,
)))
}
}
impl From<Kernel<JournalHead>> for pb::JournalHead {
fn from(value: Kernel<JournalHead>) -> Self {
Self {
position: value.0.position().get(),
root: value.0.root().map(|root| root.as_bytes().to_vec()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::JournalHead> for Kernel<JournalHead> {
type Error = StateError;
fn try_from(value: pb::JournalHead) -> Result<Self, Self::Error> {
let root = value
.root
.map(|root| {
fixed_bytes::<{ CommitRoot::LEN }>("root", &root).map(CommitRoot::from_bytes)
})
.transpose()?;
Ok(Self(JournalHead::new(
JournalPosition::new(value.position),
root,
)))
}
}
impl From<Kernel<&PartitionDescriptor>> for pb::PartitionDescriptor {
fn from(value: Kernel<&PartitionDescriptor>) -> Self {
Self {
partition: value.0.partition().as_str().to_owned(),
head: value.0.head().get(),
record_count: value.0.record_count(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl From<pb::PartitionDescriptor> for Kernel<PartitionDescriptor> {
fn from(value: pb::PartitionDescriptor) -> Self {
Self(PartitionDescriptor::new(
PartitionId::new(value.partition),
JournalPosition::new(value.head),
value.record_count,
))
}
}
impl From<Kernel<&PartitionPage>> for pb::PartitionPage {
fn from(value: Kernel<&PartitionPage>) -> Self {
Self {
partitions: value
.0
.partitions()
.iter()
.map(|descriptor| pb::PartitionDescriptor::from(Kernel(descriptor)))
.collect(),
next_after: value
.0
.next_after()
.map(|partition| partition.as_str().to_owned()),
completeness: pb::PageCompleteness::from(Kernel(value.0.completeness())).into(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::PartitionPage> for Kernel<PartitionPage> {
type Error = StateError;
fn try_from(value: pb::PartitionPage) -> Result<Self, Self::Error> {
let completeness = completeness("completeness", value.completeness)?;
let partitions = value
.partitions
.into_iter()
.map(|descriptor| Kernel::<PartitionDescriptor>::from(descriptor).into_inner())
.collect();
Ok(Self(PartitionPage::new(
partitions,
value.next_after.map(PartitionId::new),
completeness,
)))
}
}
impl From<Kernel<&JournalAttestation>> for pb::JournalAttestation {
fn from(value: Kernel<&JournalAttestation>) -> Self {
Self {
root: value.0.root().as_bytes().to_vec(),
leaf_count: value.0.leaf_count(),
signature: value.0.signature().to_vec(),
signer: value.0.signer().to_vec(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::JournalAttestation> for Kernel<JournalAttestation> {
type Error = StateError;
fn try_from(value: pb::JournalAttestation) -> Result<Self, Self::Error> {
let root = CommitRoot::from_bytes(fixed_bytes::<{ CommitRoot::LEN }>("root", &value.root)?);
Ok(Self(JournalAttestation::new(
root,
value.leaf_count,
value.signature,
value.signer,
)))
}
}
impl From<Kernel<&JournalInclusionProof>> for pb::JournalInclusionProof {
fn from(value: Kernel<&JournalInclusionProof>) -> Self {
Self {
position: value.0.position().get(),
leaf_count: value.0.leaf_count(),
inactive_peaks: value.0.inactive_peaks(),
digests: value
.0
.digests()
.iter()
.map(|digest| digest.to_vec())
.collect(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::JournalInclusionProof> for Kernel<JournalInclusionProof> {
type Error = StateError;
fn try_from(value: pb::JournalInclusionProof) -> Result<Self, Self::Error> {
let digests = value
.digests
.into_iter()
.map(|digest| fixed_bytes::<{ CommitRoot::LEN }>("digests", &digest))
.collect::<Result<Vec<_>, _>>()?;
Ok(Self(JournalInclusionProof::new(
JournalPosition::new(value.position),
value.leaf_count,
value.inactive_peaks,
digests,
)))
}
}
impl From<Kernel<&JournalProof>> for pb::JournalProof {
fn from(value: Kernel<&JournalProof>) -> Self {
Self {
attestation: buffa::MessageField::some(pb::JournalAttestation::from(Kernel(
value.0.attestation(),
))),
inclusion: buffa::MessageField::some(pb::JournalInclusionProof::from(Kernel(
value.0.inclusion(),
))),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::JournalProof> for Kernel<JournalProof> {
type Error = StateError;
fn try_from(value: pb::JournalProof) -> Result<Self, Self::Error> {
let attestation = Kernel::<JournalAttestation>::try_from(required(
"attestation",
"a proof names the signed root it resolves to",
value.attestation,
)?)?
.into_inner();
let inclusion = Kernel::<JournalInclusionProof>::try_from(required(
"inclusion",
"a proof names the record's path to its root",
value.inclusion,
)?)?
.into_inner();
Ok(Self(JournalProof::new(attestation, inclusion)))
}
}
impl From<Kernel<&RecordDecision>> for pb::RecordDecision {
fn from(value: Kernel<&RecordDecision>) -> Self {
use pb::__buffa::oneof::record_decision::Disposition;
let disposition = match value.0.disposition() {
RecordDisposition::Drop => Disposition::from(pb::DropRecord::default()),
RecordDisposition::Replace(payload) => Disposition::Replace(payload.clone()),
};
Self {
position: value.0.position().get(),
disposition: Some(disposition),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::RecordDecision> for Kernel<RecordDecision> {
type Error = StateError;
fn try_from(value: pb::RecordDecision) -> Result<Self, Self::Error> {
use pb::__buffa::oneof::record_decision::Disposition;
let disposition = match value.disposition {
Some(Disposition::Drop(_)) => RecordDisposition::Drop,
Some(Disposition::Replace(payload)) => RecordDisposition::Replace(payload),
None => {
return Err(malformed(
"disposition",
"an excision decision names what happens to the record",
));
}
};
Ok(Self(RecordDecision::new(
JournalPosition::new(value.position),
disposition,
)))
}
}
impl From<Kernel<&CommandMetadata>> for pb::JournalMutationCommand {
fn from(value: Kernel<&CommandMetadata>) -> Self {
let metadata = value.0;
Self {
command_id: metadata.command_id().as_str().to_owned(),
partition: metadata.scope().partition().as_str().to_owned(),
aggregate: metadata.scope().aggregate().as_str().to_owned(),
namespace: metadata.scope().namespace().as_str().to_owned(),
purpose: metadata.envelope().purpose().as_str().to_owned(),
audience: metadata.envelope().audience().as_str().to_owned(),
digest: metadata.digest().as_bytes().to_vec(),
precondition: buffa::MessageField::some(pb::Precondition::from(Kernel(
metadata.precondition(),
))),
fence: metadata.fence().map(FencingToken::get),
expected_root: None,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::JournalMutationCommand> for Kernel<CommandMetadata> {
type Error = StateError;
fn try_from(value: pb::JournalMutationCommand) -> Result<Self, Self::Error> {
Self::try_from(WireMetadata {
command_id: value.command_id,
partition: value.partition,
aggregate: value.aggregate,
namespace: value.namespace,
purpose: value.purpose,
audience: value.audience,
digest: value.digest,
precondition: value.precondition.into_option(),
fence: value.fence,
})
}
}
impl From<Kernel<&DestroyPartition>> for pb::JournalMutationCommand {
fn from(value: Kernel<&DestroyPartition>) -> Self {
let mut command = Self::from(Kernel(value.0.metadata()));
command.expected_root = value.0.expected_root().map(|root| root.as_bytes().to_vec());
command
}
}
impl TryFrom<pb::JournalMutationCommand> for Kernel<DestroyPartition> {
type Error = StateError;
fn try_from(value: pb::JournalMutationCommand) -> Result<Self, Self::Error> {
let expected_root = value
.expected_root
.as_deref()
.map(|bytes| fixed_bytes::<{ CommitRoot::LEN }>("expected_root", bytes))
.transpose()?
.map(CommitRoot::from_bytes);
let metadata = Kernel::<CommandMetadata>::try_from(value)?.into_inner();
let command = DestroyPartition::new(metadata);
Ok(Self(match expected_root {
Some(root) => command.at_root(root),
None => command,
}))
}
}
impl From<Kernel<&RepairPartition>> for pb::JournalMutationCommand {
fn from(value: Kernel<&RepairPartition>) -> Self {
Self::from(Kernel(value.0.metadata()))
}
}
impl From<Kernel<&ExcisePartitionRecords>> for pb::JournalMutationCommand {
fn from(value: Kernel<&ExcisePartitionRecords>) -> Self {
Self::from(Kernel(value.0.metadata()))
}
}
impl From<Kernel<&QuarantinedRecord>> for pb::QuarantinedRecord {
fn from(value: Kernel<&QuarantinedRecord>) -> Self {
Self {
position: value.0.position().get(),
reason: value.0.reason().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl From<pb::QuarantinedRecord> for Kernel<QuarantinedRecord> {
fn from(value: pb::QuarantinedRecord) -> Self {
Self(QuarantinedRecord::new(
JournalPosition::new(value.position),
value.reason,
))
}
}
impl From<Kernel<&GetJournalHead>> for String {
fn from(value: Kernel<&GetJournalHead>) -> Self {
value.0.partition().as_str().to_owned()
}
}
impl From<Kernel<&GetJournalRoot>> for String {
fn from(value: Kernel<&GetJournalRoot>) -> Self {
value.0.partition().as_str().to_owned()
}
}
#[must_use]
pub fn listing(start_after: Option<String>, limit: u32) -> ListJournalPartitions {
let request = ListJournalPartitions::new(limit);
match start_after {
Some(after) => request.after(PartitionId::new(after)),
None => request,
}
}
#[must_use]
pub fn proof_request(partition: String, position: u64) -> GetJournalProof {
GetJournalProof::new(PartitionId::new(partition), JournalPosition::new(position))
}
impl From<Kernel<&PartitionVerification>> for pb::VerifyPartitionReplayReply {
fn from(value: Kernel<&PartitionVerification>) -> Self {
use pb::__buffa::oneof::verify_partition_replay_reply::Verdict;
let verdict = match value.0 {
PartitionVerification::Verified {
event_count,
signed_root_count,
} => Verdict::from(pb::ReplayVerified {
event_count: *event_count,
signed_root_count: *signed_root_count,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
PartitionVerification::Violation { category, reason } => {
Verdict::from(pb::ReplayViolation {
reason: reason.clone(),
category: category.as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
};
Self {
verdict: Some(verdict),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::VerifyPartitionReplayReply> for Kernel<PartitionVerification> {
type Error = StateError;
fn try_from(value: pb::VerifyPartitionReplayReply) -> Result<Self, Self::Error> {
use pb::__buffa::oneof::verify_partition_replay_reply::Verdict;
let verification = match value.verdict {
Some(Verdict::Verified(verified)) => {
PartitionVerification::verified(verified.event_count, verified.signed_root_count)
}
Some(Verdict::Violation(violation)) => PartitionVerification::violation(
ViolationCategory::from_identifier(&violation.category),
violation.reason.clone(),
),
None => {
return Err(malformed(
"verdict",
"a verified replay answers with a verdict, and an absent one is not a pass",
));
}
};
Ok(Self(verification))
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs, clippy::unwrap_used)]
use super::*;
use polyc_state::{
consistency::Consistency,
id::SnapshotId,
page::PageCompleteness,
receipt::{CommitEvidence, Receipt},
};
fn drafts() -> Vec<JournalRecordDraft> {
vec![
JournalRecordDraft::new(RecordKind::new("turn_input"), b"one".to_vec()),
JournalRecordDraft::new(RecordKind::new("turn_output"), b"two".to_vec())
.with_trust(RecordTrust::QuarantinedContent),
]
}
fn batch() -> CommitJournalBatch {
let records = drafts();
let metadata = CommandMetadata::new(
CommandId::new("cmd-1"),
journal::family(),
ContentDigest::from_bytes([5; ContentDigest::LEN]),
CommandScope::new(
AggregateId::new("conv-1"),
PartitionId::new("conv-1"),
NamespaceId::new("default"),
),
CommandEnvelope::new(
Purpose::new("turn"),
crate::state_audience(),
journal_bounds(),
),
)
.with_precondition(Precondition::JournalHead(JournalPosition::new(4)))
.with_fence(FencingToken::new(3));
CommitJournalBatch::new(metadata, records)
}
#[test]
fn a_batch_round_trips_field_for_field() {
let original = batch();
let back =
Kernel::<CommitJournalBatch>::try_from(pb::CommitJournalBatch::from(Kernel(&original)))
.unwrap()
.into_inner();
assert_eq!(back, original);
assert_eq!(back.metadata().fence(), Some(FencingToken::new(3)));
assert_eq!(
back.metadata().precondition(),
Precondition::JournalHead(JournalPosition::new(4))
);
assert_eq!(back.canonical_bytes(), original.canonical_bytes());
}
#[test]
fn a_batch_with_a_wrong_width_digest_is_malformed() {
let mut wire = pb::CommitJournalBatch::from(Kernel(&batch()));
wire.digest = vec![1, 2, 3];
let error = Kernel::<CommitJournalBatch>::try_from(wire).unwrap_err();
assert!(matches!(error, StateError::Malformed { ref field, .. } if field == "digest"));
}
#[test]
fn a_record_with_no_kind_is_malformed() {
let mut wire = pb::CommitJournalBatch::from(Kernel(&batch()));
wire.records[0].kind = String::new();
let error = Kernel::<CommitJournalBatch>::try_from(wire).unwrap_err();
assert!(matches!(error, StateError::Malformed { ref field, .. } if field == "kind"));
}
#[test]
fn an_unknown_trust_class_is_refused() {
let mut wire = pb::CommitJournalBatch::from(Kernel(&batch()));
wire.records[0].trust = buffa::EnumValue::from(99);
let error = Kernel::<CommitJournalBatch>::try_from(wire).unwrap_err();
assert!(matches!(error, StateError::Malformed { ref field, .. } if field == "trust"));
}
#[test]
fn a_range_round_trips_with_its_anchor() {
let anchor = JournalAnchor::new(PartitionId::new("conv-1"), JournalPosition::new(7));
let original = JournalRange::new(
anchor.clone(),
vec![JournalRecord::new(
JournalPosition::new(1),
RecordKind::new("k"),
RecordTrust::TrustedUser,
b"body".to_vec(),
)],
Some(JournalPosition::new(2)),
PageCompleteness::Truncated,
);
let back = Kernel::<JournalRange>::try_from(pb::JournalRange::from(Kernel(&original)))
.unwrap()
.into_inner();
assert_eq!(back, original);
assert_eq!(back.anchor().snapshot(), anchor.snapshot());
}
#[test]
fn a_range_that_names_no_anchor_is_malformed() {
let mut wire = pb::JournalRange::from(Kernel(&JournalRange::new(
JournalAnchor::new(PartitionId::new("p"), JournalPosition::new(1)),
Vec::new(),
None,
PageCompleteness::Complete,
)));
wire.anchor = buffa::MessageField::default();
let error = Kernel::<JournalRange>::try_from(wire).unwrap_err();
assert!(matches!(error, StateError::Malformed { ref field, .. } if field == "anchor"));
}
#[test]
fn a_read_request_round_trips_with_and_without_an_anchor() {
let bare = ReadJournalRange::new(PartitionId::new("conv-1"), JournalPosition::new(3), 8);
assert_eq!(
Kernel::<ReadJournalRange>::from(pb::ReadJournalRange::from(Kernel(&bare)))
.into_inner(),
bare
);
let anchored = bare.clone().within(JournalAnchor::new(
PartitionId::new("conv-1"),
JournalPosition::new(9),
));
assert_eq!(
Kernel::<ReadJournalRange>::from(pb::ReadJournalRange::from(Kernel(&anchored)))
.into_inner(),
anchored
);
}
#[test]
fn a_head_round_trips_with_and_without_a_root() {
for head in [
JournalHead::new(JournalPosition::new(0), None),
JournalHead::new(
JournalPosition::new(5),
Some(CommitRoot::from_bytes([2; CommitRoot::LEN])),
),
] {
let back = Kernel::<JournalHead>::try_from(pb::JournalHead::from(Kernel(head)))
.unwrap()
.into_inner();
assert_eq!(back.position(), head.position());
assert_eq!(back.root(), head.root());
}
let mut wrong = pb::JournalHead::from(Kernel(JournalHead::new(
JournalPosition::new(1),
Some(CommitRoot::from_bytes([0; CommitRoot::LEN])),
)));
wrong.root = Some(vec![1]);
assert!(matches!(
Kernel::<JournalHead>::try_from(wrong).unwrap_err(),
StateError::Malformed { ref field, .. } if field == "root"
));
}
#[test]
fn a_partition_page_round_trips() {
let original = PartitionPage::new(
vec![PartitionDescriptor::new(
PartitionId::new("conv-1"),
JournalPosition::new(4),
4,
)],
Some(PartitionId::new("conv-1")),
PageCompleteness::Truncated,
);
let back = Kernel::<PartitionPage>::try_from(pb::PartitionPage::from(Kernel(&original)))
.unwrap()
.into_inner();
assert_eq!(back, original);
assert!(back.honors_its_bound());
}
#[test]
fn a_proof_round_trips_and_a_wrong_width_digest_is_refused() {
let original = JournalProof::new(
JournalAttestation::new(
CommitRoot::from_bytes([8; CommitRoot::LEN]),
6,
vec![1, 2, 3],
vec![4, 5, 6],
),
JournalInclusionProof::new(
JournalPosition::new(2),
6,
1,
vec![[9; CommitRoot::LEN], [10; CommitRoot::LEN]],
),
);
let wire = pb::JournalProof::from(Kernel(&original));
assert_eq!(
Kernel::<JournalProof>::try_from(wire.clone())
.unwrap()
.into_inner(),
original
);
let wire = pb::JournalProof::from(Kernel(&original));
let mut broken_inclusion = wire.inclusion.into_option().unwrap();
broken_inclusion.digests[0] = vec![1];
let broken = pb::JournalProof {
attestation: wire.attestation,
inclusion: buffa::MessageField::some(broken_inclusion),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
assert!(matches!(
Kernel::<JournalProof>::try_from(broken).unwrap_err(),
StateError::Malformed { ref field, .. } if field == "digests"
));
}
#[test]
fn an_excision_decision_round_trips_both_dispositions() {
for decision in [
RecordDecision::new(JournalPosition::new(1), RecordDisposition::Drop),
RecordDecision::new(
JournalPosition::new(2),
RecordDisposition::Replace(b"redacted".to_vec()),
),
] {
let back =
Kernel::<RecordDecision>::try_from(pb::RecordDecision::from(Kernel(&decision)))
.unwrap()
.into_inner();
assert_eq!(back, decision);
}
}
#[test]
fn a_decision_with_no_disposition_is_malformed() {
let mut wire = pb::RecordDecision::from(Kernel(&RecordDecision::new(
JournalPosition::new(1),
RecordDisposition::Drop,
)));
wire.disposition = None;
assert!(matches!(
Kernel::<RecordDecision>::try_from(wire).unwrap_err(),
StateError::Malformed { ref field, .. } if field == "disposition"
));
}
#[test]
fn a_mutation_commands_identity_and_preconditions_round_trip() {
let original = ExcisePartitionRecords::new(
batch().metadata().clone(),
vec![RecordDecision::new(
JournalPosition::new(1),
RecordDisposition::Drop,
)],
);
let back = Kernel::<CommandMetadata>::try_from(pb::JournalMutationCommand::from(Kernel(
&original,
)))
.unwrap()
.into_inner();
assert_eq!(&back, original.metadata());
let destroy = DestroyPartition::new(batch().metadata().clone())
.at_root(CommitRoot::from_bytes([7; CommitRoot::LEN]));
assert_eq!(
Kernel::<DestroyPartition>::try_from(pb::JournalMutationCommand::from(Kernel(
&destroy
)))
.unwrap()
.into_inner(),
destroy
);
let repair = RepairPartition::new(batch().metadata().clone());
assert_eq!(
Kernel::<CommandMetadata>::try_from(pb::JournalMutationCommand::from(Kernel(&repair)))
.unwrap()
.into_inner(),
*repair.metadata()
);
}
#[test]
fn a_quarantined_record_round_trips() {
let original = QuarantinedRecord::new(JournalPosition::new(3), "decode failed");
let back =
Kernel::<QuarantinedRecord>::from(pb::QuarantinedRecord::from(Kernel(&original)))
.into_inner();
assert_eq!(back, original);
}
#[test]
fn the_small_request_shapes_carry_what_they_name() {
assert_eq!(
String::from(Kernel(&GetJournalHead::new(PartitionId::new("conv-1")))),
"conv-1"
);
assert_eq!(
String::from(Kernel(&GetJournalRoot::new(PartitionId::new("conv-1")))),
"conv-1"
);
assert_eq!(
listing(Some("conv-1".to_owned()), 4).start_after(),
Some(&PartitionId::new("conv-1"))
);
assert_eq!(listing(None, 4).start_after(), None);
assert_eq!(
proof_request("conv-1".to_owned(), 2).position(),
JournalPosition::new(2)
);
}
#[test]
fn a_verdict_round_trips_and_an_absent_one_is_never_a_pass() {
for verdict in [
PartitionVerification::verified(0, 0),
PartitionVerification::verified(9, 2),
PartitionVerification::violation(
ViolationCategory::RootMismatch,
"a root does not hold",
),
PartitionVerification::violation(
ViolationCategory::TruncatedReplay,
"the replay is shorter than its durable floor",
),
PartitionVerification::violation(
ViolationCategory::Unclassified,
"something this build has no name for",
),
] {
let back = Kernel::<PartitionVerification>::try_from(
pb::VerifyPartitionReplayReply::from(Kernel(&verdict)),
)
.unwrap()
.into_inner();
assert_eq!(back, verdict);
}
let error =
Kernel::<PartitionVerification>::try_from(pb::VerifyPartitionReplayReply::default())
.unwrap_err();
assert!(
matches!(error, StateError::Malformed { ref field, .. } if field == "verdict"),
"got {error}"
);
}
#[test]
fn a_journal_receipt_keeps_its_position_root_and_anchor() {
let anchor = JournalAnchor::new(PartitionId::new("conv-1"), JournalPosition::new(2));
let evidence = CommitEvidence::new()
.with_position(JournalPosition::new(2))
.with_root(CommitRoot::from_bytes([6; CommitRoot::LEN]))
.with_snapshot(anchor.snapshot());
let receipt = Receipt::committed(
batch().metadata(),
evidence,
Consistency::OrderedPerAggregate,
);
let back = Kernel::<Receipt>::try_from(pb::Receipt::from(Kernel(&receipt)))
.unwrap()
.into_inner();
assert_eq!(back.evidence(), receipt.evidence());
assert_eq!(
back.evidence().snapshot().map(SnapshotId::as_str),
Some("journal:conv-1@2")
);
}
}