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, ProtocolVersion, Purpose},
journal::{JournalRecord, JournalRecordDraft},
persona_memory::journal::{
AppendMemoryRecords, DestroyMemoryRecords, FencedMemoryWrite, MemoryAppendOutcome,
MemoryJournalPartition, MemoryMigrationOutcome, MemoryPartitionPage, MemoryReplayPage,
MemoryRewriteEdit, MemoryRewriteOutcome, MigrateMemoryRecords, RewriteMemoryRecords,
SignedMemoryHead,
},
receipt::Receipt,
revision::{CommitRoot, JournalHead, JournalPosition},
};
use crate::wire::{Kernel, fixed_bytes, required};
fn metadata_to_wire(value: &CommandMetadata) -> pb::PersonaMemoryJournalCommandMetadata {
pb::PersonaMemoryJournalCommandMetadata {
command_id: value.command_id().as_str().to_owned(),
digest: value.digest().as_bytes().to_vec(),
aggregate: value.scope().aggregate().as_str().to_owned(),
partition: value.scope().partition().as_str().to_owned(),
namespace: value.scope().namespace().as_str().to_owned(),
purpose: value.envelope().purpose().as_str().to_owned(),
command_audience: value.envelope().audience().as_str().to_owned(),
max_payload_bytes: value.envelope().bounds().max_payload_bytes(),
max_records: value.envelope().bounds().max_records(),
precondition: buffa::MessageField::some(pb::Precondition::from(Kernel(
value.precondition(),
))),
fence: value.fence().map_or(0, FencingToken::get),
protocol_version: value.protocol_version().get(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn metadata_from_wire(
value: pb::PersonaMemoryJournalCommandMetadata,
) -> Result<CommandMetadata, StateError> {
Ok(CommandMetadata::new(
CommandId::new(value.command_id),
polyc_state::persona_memory::journal::family(),
ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
"digest",
&value.digest,
)?),
CommandScope::new(
AggregateId::new(value.aggregate),
PartitionId::new(value.partition),
NamespaceId::new(value.namespace),
),
CommandEnvelope::new(
Purpose::new(value.purpose),
Audience::new(value.command_audience),
ResourceBounds::new(value.max_payload_bytes, value.max_records),
),
)
.with_precondition(
Kernel::<Precondition>::try_from(required(
"precondition",
"a memory-journal command declares its precondition",
value.precondition,
)?)?
.into_inner(),
)
.with_fence(FencingToken::new(value.fence))
.with_protocol_version(ProtocolVersion::new(value.protocol_version)))
}
fn records_to_wire(records: &[JournalRecordDraft]) -> Vec<pb::JournalRecordDraft> {
records
.iter()
.map(|record| pb::JournalRecordDraft::from(Kernel(record)))
.collect()
}
fn records_from_wire(
records: Vec<pb::JournalRecordDraft>,
) -> Result<Vec<JournalRecordDraft>, StateError> {
records
.into_iter()
.map(|record| Kernel::<JournalRecordDraft>::try_from(record).map(Kernel::into_inner))
.collect()
}
pub(super) fn append_to_wire(value: &AppendMemoryRecords) -> pb::AppendPersonaMemoryRecordsCommand {
pb::AppendPersonaMemoryRecordsCommand {
metadata: buffa::MessageField::some(metadata_to_wire(value.write().metadata())),
records: records_to_wire(value.records()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(super) fn append_from_wire(
value: pb::AppendPersonaMemoryRecordsCommand,
) -> Result<AppendMemoryRecords, StateError> {
let metadata = metadata_from_wire(required(
"metadata",
"an append carries metadata",
value.metadata,
)?)?;
let partition = MemoryJournalPartition::parse(metadata.scope().partition().as_str())?;
AppendMemoryRecords::new(
FencedMemoryWrite::new(metadata)?,
partition,
records_from_wire(value.records)?,
)
}
fn edit_to_wire(value: &MemoryRewriteEdit) -> pb::PersonaMemoryRewriteEdit {
use pb::__buffa::oneof::persona_memory_rewrite_edit::Edit;
let edit = match value {
MemoryRewriteEdit::Drop { position } => Edit::from(pb::PersonaMemoryRewriteDrop {
position: *position,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
MemoryRewriteEdit::Replace { position, payload } => {
Edit::from(pb::PersonaMemoryRewriteReplace {
position: *position,
payload: payload.clone(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
};
pb::PersonaMemoryRewriteEdit {
edit: Some(edit),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn edit_from_wire(value: pb::PersonaMemoryRewriteEdit) -> Result<MemoryRewriteEdit, StateError> {
use pb::__buffa::oneof::persona_memory_rewrite_edit::Edit;
match value.edit {
Some(Edit::Drop(value)) => Ok(MemoryRewriteEdit::Drop {
position: value.position,
}),
Some(Edit::Replace(value)) => Ok(MemoryRewriteEdit::Replace {
position: value.position,
payload: value.payload,
}),
None => Err(StateError::Malformed {
field: "edit".into(),
reason: "a rewrite edit names drop or replace".into(),
}),
}
}
pub(super) fn rewrite_to_wire(
value: &RewriteMemoryRecords,
) -> pb::RewritePersonaMemoryRecordsCommand {
pb::RewritePersonaMemoryRecordsCommand {
metadata: buffa::MessageField::some(metadata_to_wire(value.write().metadata())),
edits: value.edits().iter().map(edit_to_wire).collect(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(super) fn rewrite_from_wire(
value: pb::RewritePersonaMemoryRecordsCommand,
) -> Result<RewriteMemoryRecords, StateError> {
let metadata = metadata_from_wire(required(
"metadata",
"a rewrite carries metadata",
value.metadata,
)?)?;
let partition = MemoryJournalPartition::parse(metadata.scope().partition().as_str())?;
RewriteMemoryRecords::new(
FencedMemoryWrite::new(metadata)?,
partition,
value
.edits
.into_iter()
.map(edit_from_wire)
.collect::<Result<_, _>>()?,
)
}
pub(super) fn migrate_to_wire(
value: &MigrateMemoryRecords,
) -> pb::MigratePersonaMemoryRecordsCommand {
pb::MigratePersonaMemoryRecordsCommand {
metadata: buffa::MessageField::some(metadata_to_wire(value.write().metadata())),
to_partition: value.to().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(super) fn migrate_from_wire(
value: pb::MigratePersonaMemoryRecordsCommand,
) -> Result<MigrateMemoryRecords, StateError> {
let metadata = metadata_from_wire(required(
"metadata",
"a migration carries metadata",
value.metadata,
)?)?;
let from = MemoryJournalPartition::parse(metadata.scope().partition().as_str())?;
MigrateMemoryRecords::new(
FencedMemoryWrite::new(metadata)?,
from,
MemoryJournalPartition::parse(value.to_partition)?,
)
}
pub(super) fn destroy_to_wire(
value: &DestroyMemoryRecords,
) -> pb::DestroyPersonaMemoryRecordsCommand {
pb::DestroyPersonaMemoryRecordsCommand {
metadata: buffa::MessageField::some(metadata_to_wire(value.write().metadata())),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(super) fn destroy_from_wire(
value: pb::DestroyPersonaMemoryRecordsCommand,
) -> Result<DestroyMemoryRecords, StateError> {
let metadata = metadata_from_wire(required(
"metadata",
"a destruction carries metadata",
value.metadata,
)?)?;
let partition = MemoryJournalPartition::parse(metadata.scope().partition().as_str())?;
DestroyMemoryRecords::new(FencedMemoryWrite::new(metadata)?, partition)
}
pub(super) fn head_to_wire(value: JournalHead) -> pb::StatePersonaMemoryJournalHead {
pb::StatePersonaMemoryJournalHead {
position: value.position().get(),
root: value.root().map(|root| root.as_bytes().to_vec()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(super) fn head_from_wire(
value: pb::StatePersonaMemoryJournalHead,
) -> Result<JournalHead, StateError> {
Ok(JournalHead::new(
JournalPosition::new(value.position),
value
.root
.map(|root| {
fixed_bytes::<{ CommitRoot::LEN }>("root", &root).map(CommitRoot::from_bytes)
})
.transpose()?,
))
}
fn signed_head_to_wire(value: &SignedMemoryHead) -> pb::SignedPersonaMemoryHead {
pb::SignedPersonaMemoryHead {
namespace: value.namespace().as_str().to_owned(),
partition: value.partition().as_str().to_owned(),
head: buffa::MessageField::some(head_to_wire(value.head())),
mmr_root: value.mmr_root().as_bytes().to_vec(),
signer_key_id: value.signer_key_id().to_owned(),
signature: value.signature().to_vec(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn signed_head_from_wire(
value: pb::SignedPersonaMemoryHead,
) -> Result<SignedMemoryHead, StateError> {
SignedMemoryHead::new(
NamespaceId::new(value.namespace),
MemoryJournalPartition::parse(value.partition)?,
head_from_wire(required(
"signed_head.head",
"a signed head names the exact journal head",
value.head,
)?)?,
CommitRoot::from_bytes(fixed_bytes::<{ CommitRoot::LEN }>(
"signed_head.mmr_root",
&value.mmr_root,
)?),
value.signer_key_id,
value.signature,
)
}
pub(super) fn replay_to_wire(value: &MemoryReplayPage) -> pb::ReplayPersonaMemoryRecordsReply {
pb::ReplayPersonaMemoryRecordsReply {
records: value
.records()
.iter()
.map(|record| pb::JournalRecord::from(Kernel(record)))
.collect(),
next: value.next(),
head: buffa::MessageField::some(head_to_wire(value.head())),
signed_head: value
.signed_head()
.map(|head| buffa::MessageField::some(signed_head_to_wire(head)))
.unwrap_or_default(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(super) fn replay_from_wire(
value: pb::ReplayPersonaMemoryRecordsReply,
) -> Result<MemoryReplayPage, StateError> {
let page = MemoryReplayPage::new(
value
.records
.into_iter()
.map(|record| Kernel::<JournalRecord>::try_from(record).map(Kernel::into_inner))
.collect::<Result<_, _>>()?,
value.next,
head_from_wire(required(
"head",
"a replay returns its observed head",
value.head,
)?)?,
);
let signed_head = signed_head_from_wire(required(
"signed_head",
"a served replay carries State's signed full head",
value.signed_head,
)?)?;
if signed_head.head() != page.head() {
return Err(StateError::Malformed {
field: "signed_head.head".into(),
reason: "the signed and replay heads agree".into(),
});
}
Ok(page.with_signed_head(signed_head))
}
pub(super) fn partition_page_to_wire(
value: &MemoryPartitionPage,
) -> pb::ListPersonaMemoryPartitionsReply {
pb::ListPersonaMemoryPartitionsReply {
partitions: value
.partitions()
.iter()
.map(|value| value.as_str().to_owned())
.collect(),
next: value.next().map(str::to_owned),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(super) fn partition_page_from_wire(
value: pb::ListPersonaMemoryPartitionsReply,
) -> Result<MemoryPartitionPage, StateError> {
Ok(MemoryPartitionPage::new(
value
.partitions
.into_iter()
.map(MemoryJournalPartition::parse)
.collect::<Result<_, _>>()?,
value.next,
))
}
pub(super) fn append_outcome_from_wire(
receipt: pb::Receipt,
positions: Vec<u64>,
) -> Result<MemoryAppendOutcome, StateError> {
Ok(MemoryAppendOutcome::new(
Kernel::<Receipt>::try_from(receipt)?.into_inner(),
positions,
))
}
pub(super) fn rewrite_outcome_from_wire(
receipt: pb::Receipt,
dropped: u64,
) -> Result<MemoryRewriteOutcome, StateError> {
Ok(MemoryRewriteOutcome::new(
Kernel::<Receipt>::try_from(receipt)?.into_inner(),
usize::try_from(dropped).map_err(|_| StateError::Malformed {
field: "dropped".into(),
reason: "the count fits this architecture".into(),
})?,
))
}
pub(super) fn migration_outcome_from_wire(
receipt: pb::Receipt,
copied: u64,
) -> Result<MemoryMigrationOutcome, StateError> {
Ok(MemoryMigrationOutcome::new(
Kernel::<Receipt>::try_from(receipt)?.into_inner(),
usize::try_from(copied).map_err(|_| StateError::Malformed {
field: "copied".into(),
reason: "the count fits this architecture".into(),
})?,
))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_served_replay_without_a_signed_head_fails_closed() {
let reply = pb::ReplayPersonaMemoryRecordsReply {
records: Vec::new(),
next: None,
head: buffa::MessageField::some(head_to_wire(JournalHead::new(
JournalPosition::new(0),
None,
))),
signed_head: buffa::MessageField::none(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
assert!(replay_from_wire(reply).is_err());
}
#[test]
fn a_signed_head_round_trips_every_bound_field() {
let partition = MemoryJournalPartition::parse("persona-wire-mem").unwrap();
let head = JournalHead::new(JournalPosition::new(0), None);
let signed = SignedMemoryHead::new(
NamespaceId::new("polychrome"),
partition,
head,
CommitRoot::from_bytes([1; CommitRoot::LEN]),
"key".into(),
vec![2; polyc_state::persona_memory::journal::MEMORY_HEAD_SIGNATURE_BYTES],
)
.unwrap();
let page = MemoryReplayPage::new(Vec::new(), None, head).with_signed_head(signed);
assert_eq!(replay_from_wire(replay_to_wire(&page)).unwrap(), page);
}
}