use std::time::Duration;
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
command::{CommandEnvelope, CommandMetadata, CommandScope, FencingToken, Precondition},
deadline::MonotonicInstant,
digest::ContentDigest,
error::StateError,
id::{AggregateId, Audience, CommandId, NamespaceId, OwnerId, PartitionId, Purpose},
immutable::{
self, AtRestProtection, Classification, ContentReference, Generation, ObjectDescriptor,
ObjectHead, ObjectId, ObjectMetadata, Retention,
},
page::Page,
};
use crate::wire::{
Kernel, completeness, consistency, fixed_bytes, known, malformed, nanos_from_duration, required,
};
impl From<Kernel<Classification>> for pb::ObjectClassification {
fn from(value: Kernel<Classification>) -> Self {
match value.0 {
Classification::Public => Self::OBJECT_CLASSIFICATION_PUBLIC,
Classification::Internal => Self::OBJECT_CLASSIFICATION_INTERNAL,
Classification::Confidential => Self::OBJECT_CLASSIFICATION_CONFIDENTIAL,
Classification::Restricted => Self::OBJECT_CLASSIFICATION_RESTRICTED,
}
}
}
pub(crate) fn classification(
field: &str,
value: buffa::EnumValue<pb::ObjectClassification>,
) -> Result<Classification, StateError> {
match known(field, value)? {
pb::ObjectClassification::OBJECT_CLASSIFICATION_PUBLIC => Ok(Classification::Public),
pb::ObjectClassification::OBJECT_CLASSIFICATION_INTERNAL => Ok(Classification::Internal),
pb::ObjectClassification::OBJECT_CLASSIFICATION_CONFIDENTIAL => {
Ok(Classification::Confidential)
}
pb::ObjectClassification::OBJECT_CLASSIFICATION_RESTRICTED => {
Ok(Classification::Restricted)
}
pb::ObjectClassification::OBJECT_CLASSIFICATION_UNSPECIFIED => Err(malformed(
field,
"object metadata declares one known classification",
)),
}
}
impl From<Kernel<AtRestProtection>> for pb::ObjectAtRestProtection {
fn from(value: Kernel<AtRestProtection>) -> Self {
match value.0 {
AtRestProtection::None => Self::OBJECT_AT_REST_PROTECTION_NONE,
AtRestProtection::ManagedKey => Self::OBJECT_AT_REST_PROTECTION_MANAGED_KEY,
AtRestProtection::TenantKey => Self::OBJECT_AT_REST_PROTECTION_TENANT_KEY,
}
}
}
pub(crate) fn protection(
field: &str,
value: buffa::EnumValue<pb::ObjectAtRestProtection>,
) -> Result<AtRestProtection, StateError> {
match known(field, value)? {
pb::ObjectAtRestProtection::OBJECT_AT_REST_PROTECTION_NONE => Ok(AtRestProtection::None),
pb::ObjectAtRestProtection::OBJECT_AT_REST_PROTECTION_MANAGED_KEY => {
Ok(AtRestProtection::ManagedKey)
}
pb::ObjectAtRestProtection::OBJECT_AT_REST_PROTECTION_TENANT_KEY => {
Ok(AtRestProtection::TenantKey)
}
pb::ObjectAtRestProtection::OBJECT_AT_REST_PROTECTION_UNSPECIFIED => Err(malformed(
field,
"the module declares one known protection tier",
)),
}
}
impl From<Kernel<Retention>> for pb::ObjectRetention {
fn from(value: Kernel<Retention>) -> Self {
use pb::__buffa::oneof::object_retention::Policy;
let policy = match value.0 {
Retention::For(duration) => Policy::from(pb::RetainFor {
duration_nanos: nanos_from_duration(duration),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
Retention::UntilReleased => Policy::from(pb::RetainUntilReleased {
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
};
Self {
policy: Some(policy),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ObjectRetention> for Kernel<Retention> {
type Error = StateError;
fn try_from(value: pb::ObjectRetention) -> Result<Self, Self::Error> {
use pb::__buffa::oneof::object_retention::Policy;
Ok(Self(match value.policy {
Some(Policy::ForDuration(value)) => {
Retention::For(Duration::from_nanos(value.duration_nanos))
}
Some(Policy::UntilReleased(_)) => Retention::UntilReleased,
None => return Err(malformed("retention", "metadata declares its retention")),
}))
}
}
impl From<Kernel<&ObjectDescriptor>> for pb::ObjectDescriptor {
fn from(value: Kernel<&ObjectDescriptor>) -> Self {
Self {
object: value.0.object().as_str().to_owned(),
generation: value.0.generation().get(),
digest: value.0.digest().as_bytes().to_vec(),
owner: value.0.owner().as_str().to_owned(),
classification: pb::ObjectClassification::from(Kernel(value.0.classification())).into(),
retention: buffa::MessageField::some(pb::ObjectRetention::from(Kernel(
value.0.retention(),
))),
byte_len: value.0.byte_len(),
content_reference: value.0.content_reference().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ObjectDescriptor> for Kernel<ObjectDescriptor> {
type Error = StateError;
fn try_from(value: pb::ObjectDescriptor) -> Result<Self, Self::Error> {
let pb::ObjectDescriptor {
object,
generation,
digest,
owner,
classification: encoded_classification,
retention,
byte_len,
content_reference,
__buffa_unknown_fields: _,
} = value;
let digest =
ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>("digest", &digest)?);
let retention = Kernel::<Retention>::try_from(required(
"retention",
"metadata declares its retention",
retention,
)?)?
.into_inner();
Ok(Self(ObjectDescriptor::new(
ObjectId::new(object),
Generation::new(generation),
digest,
OwnerId::new(owner),
classification("classification", encoded_classification)?,
retention,
byte_len,
ContentReference::try_new(content_reference)?,
)))
}
}
impl From<Kernel<&ObjectMetadata>> for pb::ObjectMetadata {
fn from(value: Kernel<&ObjectMetadata>) -> Self {
Self {
descriptor: buffa::MessageField::some(pb::ObjectDescriptor::from(Kernel(
value.0.descriptor(),
))),
recorded_at_nanos: value.0.recorded_at().as_nanos(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ObjectMetadata> for Kernel<ObjectMetadata> {
type Error = StateError;
fn try_from(value: pb::ObjectMetadata) -> Result<Self, Self::Error> {
let descriptor = Kernel::<ObjectDescriptor>::try_from(required(
"descriptor",
"object metadata carries its descriptor",
value.descriptor,
)?)?
.into_inner();
Ok(Self(ObjectMetadata::new(
descriptor,
MonotonicInstant::from_nanos(value.recorded_at_nanos),
)))
}
}
impl From<Kernel<&ObjectHead>> for pb::ObjectHead {
fn from(value: Kernel<&ObjectHead>) -> Self {
Self {
object: value.0.object().as_str().to_owned(),
generation: value.0.generation().get(),
digest: value.0.digest().as_bytes().to_vec(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ObjectHead> for Kernel<ObjectHead> {
type Error = StateError;
fn try_from(value: pb::ObjectHead) -> Result<Self, Self::Error> {
Ok(Self(ObjectHead::new(
ObjectId::new(value.object),
Generation::new(value.generation),
ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
"digest",
&value.digest,
)?),
)))
}
}
impl From<Kernel<&Page<ObjectMetadata>>> for pb::ObjectMetadataPage {
fn from(value: Kernel<&Page<ObjectMetadata>>) -> Self {
let page = value.0;
Self {
records: page
.records()
.iter()
.map(|record| pb::ObjectMetadata::from(Kernel(record)))
.collect(),
next: page
.next_cursor()
.map_or_else(buffa::MessageField::default, |cursor| {
buffa::MessageField::some(pb::Cursor::from(Kernel(cursor)))
}),
completeness: pb::PageCompleteness::from(Kernel(page.completeness())).into(),
consistency: pb::Consistency::from(Kernel(page.consistency())).into(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ObjectMetadataPage> for Kernel<Page<ObjectMetadata>> {
type Error = StateError;
fn try_from(value: pb::ObjectMetadataPage) -> Result<Self, Self::Error> {
let records = value
.records
.into_iter()
.map(|record| Kernel::<ObjectMetadata>::try_from(record).map(Kernel::into_inner))
.collect::<Result<Vec<_>, _>>()?;
let next = value
.next
.into_option()
.map(|cursor| Kernel::from(cursor).into_inner());
Ok(Self(Page::new(
records,
next,
completeness("completeness", value.completeness)?,
consistency("consistency", value.consistency)?,
)))
}
}
impl From<Kernel<&CommandMetadata>> for pb::ObjectCommandMetadata {
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),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ObjectCommandMetadata> for Kernel<CommandMetadata> {
type Error = StateError;
fn try_from(value: pb::ObjectCommandMetadata) -> Result<Self, Self::Error> {
let precondition = Kernel::<Precondition>::try_from(required(
"precondition",
"an object command names its precondition",
value.precondition,
)?)?
.into_inner();
let mut metadata = CommandMetadata::new(
CommandId::new(value.command_id),
immutable::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),
immutable::command_bounds(),
),
)
.with_precondition(precondition);
if let Some(fence) = value.fence {
metadata = metadata.with_fence(FencingToken::new(fence));
}
Ok(Self(metadata))
}
}
#[cfg(test)]
mod tests {
use super::*;
use polyc_state::immutable::{Classification, Retention, cases};
#[test]
fn descriptors_metadata_heads_and_commands_round_trip() {
let descriptor = cases::descriptor(
7,
b"content",
Classification::Confidential,
Retention::UntilReleased,
);
let wire = pb::ObjectDescriptor::from(Kernel(&descriptor));
assert_eq!(
Kernel::<ObjectDescriptor>::try_from(wire)
.expect("decode")
.into_inner(),
descriptor
);
let metadata = cases::metadata("object-wire", &descriptor.canonical_bytes());
let wire = pb::ObjectCommandMetadata::from(Kernel(&metadata));
let back = Kernel::<CommandMetadata>::try_from(wire)
.expect("decode")
.into_inner();
assert_eq!(back.command_id(), metadata.command_id());
assert_eq!(back.digest(), metadata.digest());
assert_eq!(back.scope(), metadata.scope());
}
#[test]
fn unset_retention_and_unknown_classification_fail_closed() {
assert!(Kernel::<Retention>::try_from(pb::ObjectRetention::default()).is_err());
let mut descriptor = pb::ObjectDescriptor::from(Kernel(&cases::plain(1, b"one")));
descriptor.classification =
pb::ObjectClassification::OBJECT_CLASSIFICATION_UNSPECIFIED.into();
assert!(Kernel::<ObjectDescriptor>::try_from(descriptor).is_err());
}
}