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},
immutable::ContentReference,
projection::{
self, FamilyId, ProjectionGeneration, ProjectionKey, ProjectionManifest, PublisherFence,
PublisherId,
artifact::{ExactObjectRef, ObjectNamespace},
},
revision::PartitionIncarnation,
};
use crate::wire::{Kernel, fixed_bytes, required};
impl From<Kernel<&ProjectionKey>> for pb::ProjectionKey {
fn from(value: Kernel<&ProjectionKey>) -> Self {
Self {
family: value.0.family().as_str().to_owned(),
source_partition: value.0.source().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ProjectionKey> for Kernel<ProjectionKey> {
type Error = StateError;
fn try_from(value: pb::ProjectionKey) -> Result<Self, Self::Error> {
let pb::ProjectionKey {
family,
source_partition,
__buffa_unknown_fields: _,
} = value;
if family.is_empty() || source_partition.is_empty() {
return Err(StateError::Malformed {
field: "projection_key".to_owned(),
reason: "a projection key names a family and source partition".to_owned(),
});
}
Ok(Self(ProjectionKey::new(
FamilyId::new(family),
PartitionId::new(source_partition),
)))
}
}
impl From<Kernel<&PublisherFence>> for pb::ProjectionPublisherFence {
fn from(value: Kernel<&PublisherFence>) -> Self {
Self {
key: buffa::MessageField::some(pb::ProjectionKey::from(Kernel(value.0.key()))),
source_incarnation: value.0.incarnation().as_bytes().to_vec(),
term: value.0.term(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ProjectionPublisherFence> for Kernel<PublisherFence> {
type Error = StateError;
fn try_from(value: pb::ProjectionPublisherFence) -> Result<Self, Self::Error> {
let pb::ProjectionPublisherFence {
key,
source_incarnation,
term,
__buffa_unknown_fields: _,
} = value;
if term == 0 {
return Err(StateError::Malformed {
field: "term".to_owned(),
reason: "a publisher fence has a granted nonzero term".to_owned(),
});
}
Ok(Self(PublisherFence::new(
Kernel::<ProjectionKey>::try_from(required(
"key",
"a publisher fence names one projection key",
key,
)?)?
.into_inner(),
PartitionIncarnation::from_bytes(fixed_bytes::<{ PartitionIncarnation::LEN }>(
"source_incarnation",
&source_incarnation,
)?),
term,
)))
}
}
impl From<Kernel<&ProjectionManifest>> for pb::ProjectionManifest {
fn from(value: Kernel<&ProjectionManifest>) -> Self {
Self {
key: buffa::MessageField::some(pb::ProjectionKey::from(Kernel(value.0.key()))),
generation: value.0.generation().get(),
checkpoint: buffa::MessageField::some(pb::SourceCheckpoint::from(Kernel(
value.0.checkpoint(),
))),
schema_version: value.0.schema_version(),
fact_version: value.0.fact_version(),
object: buffa::MessageField::some(pb::ObjectDescriptor::from(Kernel(
value.0.object_descriptor(),
))),
publisher: value.0.publisher().as_str().to_owned(),
fence: buffa::MessageField::some(pb::ProjectionPublisherFence::from(Kernel(
value.0.fence(),
))),
artifact_object: buffa::MessageField::some(pb::ExactObjectRef::from(Kernel(
value.0.artifact_object(),
))),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl From<Kernel<&ExactObjectRef>> for pb::ExactObjectRef {
fn from(value: Kernel<&ExactObjectRef>) -> Self {
Self {
namespace: value.0.namespace().as_str().to_owned(),
key: value.0.key().as_str().to_owned(),
backend_generation: value.0.generation(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ExactObjectRef> for Kernel<ExactObjectRef> {
type Error = StateError;
fn try_from(value: pb::ExactObjectRef) -> Result<Self, Self::Error> {
let pb::ExactObjectRef {
namespace,
key,
backend_generation,
__buffa_unknown_fields: _,
} = value;
Ok(Self(ExactObjectRef::try_new(
ObjectNamespace::try_new(namespace)?,
ContentReference::try_new(key)?,
backend_generation,
)?))
}
}
impl TryFrom<pb::ProjectionManifest> for Kernel<ProjectionManifest> {
type Error = StateError;
fn try_from(value: pb::ProjectionManifest) -> Result<Self, Self::Error> {
let pb::ProjectionManifest {
key,
generation,
checkpoint,
schema_version,
fact_version,
object,
publisher,
fence,
artifact_object,
__buffa_unknown_fields: _,
} = value;
if generation == 0 || schema_version == 0 || fact_version == 0 || publisher.is_empty() {
return Err(StateError::Malformed {
field: "manifest".to_owned(),
reason: "a manifest names nonzero generation and versions and one publisher"
.to_owned(),
});
}
let manifest = ProjectionManifest::new(
Kernel::try_from(required("key", "a manifest names its projection key", key)?)?
.into_inner(),
ProjectionGeneration::new(generation),
Kernel::try_from(required(
"checkpoint",
"a manifest carries its exact source checkpoint",
checkpoint,
)?)?
.into_inner(),
schema_version,
fact_version,
Kernel::try_from(required(
"object",
"a manifest carries its immutable descriptor",
object,
)?)?
.into_inner(),
Kernel::try_from(required(
"artifact_object",
"a manifest names where its signed artifact is stored",
artifact_object,
)?)?
.into_inner(),
PublisherId::new(publisher),
Kernel::try_from(required(
"fence",
"a manifest carries its authority",
fence,
)?)?
.into_inner(),
);
manifest
.validate_structure()
.map_err(|_| StateError::Malformed {
field: "manifest".to_owned(),
reason: "a manifest binds its key, source, generation, descriptor, and fence"
.to_owned(),
})?;
Ok(Self(manifest))
}
}
impl From<Kernel<&CommandMetadata>> for pb::ProjectionCommandMetadata {
fn from(value: Kernel<&CommandMetadata>) -> Self {
let metadata = value.0;
Self {
command_id: metadata.command_id().as_str().to_owned(),
aggregate: metadata.scope().aggregate().as_str().to_owned(),
partition: metadata.scope().partition().as_str().to_owned(),
namespace: metadata.scope().namespace().as_str().to_owned(),
purpose: metadata.envelope().purpose().as_str().to_owned(),
command_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),
max_payload_bytes: metadata.envelope().bounds().max_payload_bytes(),
max_records: metadata.envelope().bounds().max_records(),
protocol_version: metadata.protocol_version().get(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ProjectionCommandMetadata> for Kernel<CommandMetadata> {
type Error = StateError;
fn try_from(value: pb::ProjectionCommandMetadata) -> Result<Self, Self::Error> {
let pb::ProjectionCommandMetadata {
command_id,
aggregate,
partition,
namespace,
purpose,
command_audience,
digest,
precondition,
fence,
max_payload_bytes,
max_records,
protocol_version,
__buffa_unknown_fields: _,
} = value;
let mut metadata = CommandMetadata::new(
CommandId::new(command_id),
projection::family(),
ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>("digest", &digest)?),
CommandScope::new(
AggregateId::new(aggregate),
PartitionId::new(partition),
NamespaceId::new(namespace),
),
CommandEnvelope::new(
Purpose::new(purpose),
Audience::new(command_audience),
ResourceBounds::new(max_payload_bytes, max_records),
),
)
.with_precondition(
Kernel::<Precondition>::try_from(required(
"precondition",
"a projection command names its precondition",
precondition,
)?)?
.into_inner(),
)
.with_protocol_version(ProtocolVersion::new(protocol_version));
if let Some(fence) = fence {
metadata = metadata.with_fence(FencingToken::new(fence));
}
Ok(Self(metadata))
}
}
#[cfg(test)]
mod tests {
use super::*;
use polyc_state::projection::cases;
#[test]
fn command_metadata_round_trips_every_field() {
let command = cases::register("wire-register", &cases::source(1));
let wire = pb::ProjectionCommandMetadata::from(Kernel(command.metadata()));
let decoded = Kernel::<CommandMetadata>::try_from(wire)
.expect("decode metadata")
.into_inner();
assert_eq!(decoded, *command.metadata());
}
#[test]
fn malformed_incarnation_is_refused() {
let fence = PublisherFence::new(
cases::key(),
PartitionIncarnation::from_bytes([1; PartitionIncarnation::LEN]),
1,
);
let mut wire = pb::ProjectionPublisherFence::from(Kernel(&fence));
wire.source_incarnation.pop();
assert!(Kernel::<PublisherFence>::try_from(wire).is_err());
}
#[test]
fn manifest_round_trips_every_field() {
let source = cases::source(1);
let fence = PublisherFence::new(cases::key(), source.incarnation(), 1);
let manifest = cases::transport_record(&source, 1, 7, &fence);
let wire = pb::ProjectionManifest::from(Kernel(&manifest));
let decoded = Kernel::<ProjectionManifest>::try_from(wire)
.expect("decode manifest")
.into_inner();
assert_eq!(decoded, manifest);
}
#[test]
fn contradictory_manifest_bindings_are_refused() {
let source = cases::source(1);
let fence = PublisherFence::new(cases::key(), source.incarnation(), 1);
let manifest = cases::transport_record(&source, 1, 7, &fence);
let mut wrong_partition = pb::ProjectionManifest::from(Kernel(&manifest));
wrong_partition
.checkpoint
.as_option_mut()
.expect("checkpoint")
.source
.as_option_mut()
.expect("source")
.partition = "other".to_owned();
assert!(Kernel::<ProjectionManifest>::try_from(wrong_partition).is_err());
let mut wrong_generation = pb::ProjectionManifest::from(Kernel(&manifest));
wrong_generation
.object
.as_option_mut()
.expect("object")
.generation = 2;
assert!(Kernel::<ProjectionManifest>::try_from(wrong_generation).is_err());
let mut wrong_fence = pb::ProjectionManifest::from(Kernel(&manifest));
wrong_fence
.fence
.as_option_mut()
.expect("fence")
.key
.as_option_mut()
.expect("key")
.source_partition = "other".to_owned();
assert!(Kernel::<ProjectionManifest>::try_from(wrong_fence).is_err());
}
}