use std::time::Duration;
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
cancel::CancellationToken,
command::{
CommandEnvelope, CommandMetadata, CommandScope, FencingToken, ObservedState, Precondition,
ResourceBounds,
},
conformance::{SyntheticCommand, SyntheticOp, SyntheticRecord, family},
consistency::{Consistency, Watermark},
context::CallContext,
digest::ContentDigest,
error::StateError,
id::{
AggregateId, Audience, CommandId, NamespaceId, OperationFamily, PartitionId, Purpose,
SnapshotId,
},
page::{Cursor, Page, PageCompleteness, PageRequest, Positioned, ReadStart},
receipt::{CommitDisposition, CommitEvidence, Receipt},
revision::{CommitRoot, JournalPosition, Revision},
stream::{
Backpressure, CompactedRange, Delivery, StreamChunk, StreamContract, StreamEnd,
StreamRequest,
},
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Kernel<T>(pub T);
impl<T> Kernel<T> {
#[must_use]
pub fn into_inner(self) -> T {
self.0
}
}
fn cancelled_token(withdrawn: bool) -> CancellationToken {
let token = CancellationToken::new();
if withdrawn {
token.cancel();
}
token
}
pub(crate) fn malformed(field: &str, reason: impl Into<String>) -> StateError {
StateError::Malformed {
field: field.to_owned(),
reason: reason.into(),
}
}
pub(crate) fn fixed_bytes<const N: usize>(
field: &str,
bytes: &[u8],
) -> Result<[u8; N], StateError> {
<[u8; N]>::try_from(bytes).map_err(|_| {
malformed(
field,
format!("expected {N} bytes, received {}", bytes.len()),
)
})
}
pub(crate) fn required<T: Default, P: buffa::ProtoBox<T>>(
field: &str,
reason: &str,
value: buffa::MessageField<T, P>,
) -> Result<T, StateError> {
value
.into_option()
.ok_or_else(|| malformed(field, reason.to_owned()))
}
pub(crate) fn known<E: buffa::Enumeration>(
field: &str,
value: buffa::EnumValue<E>,
) -> Result<E, StateError> {
value.as_known().ok_or_else(|| {
malformed(
field,
format!("this build does not know value {}", value.to_i32()),
)
})
}
impl From<Kernel<Consistency>> for pb::Consistency {
fn from(value: Kernel<Consistency>) -> Self {
match value.0 {
Consistency::LinearizableCurrent => Self::CONSISTENCY_LINEARIZABLE_CURRENT,
Consistency::RevisionBound => Self::CONSISTENCY_REVISION_BOUND,
Consistency::SnapshotConsistent => Self::CONSISTENCY_SNAPSHOT_CONSISTENT,
Consistency::OrderedPerAggregate => Self::CONSISTENCY_ORDERED_PER_AGGREGATE,
Consistency::EventualWithWatermark => Self::CONSISTENCY_EVENTUAL_WITH_WATERMARK,
Consistency::ExplicitlyEphemeral => Self::CONSISTENCY_EXPLICITLY_EPHEMERAL,
}
}
}
pub(crate) fn consistency(
field: &str,
value: buffa::EnumValue<pb::Consistency>,
) -> Result<Consistency, StateError> {
match known(field, value)? {
pb::Consistency::CONSISTENCY_LINEARIZABLE_CURRENT => Ok(Consistency::LinearizableCurrent),
pb::Consistency::CONSISTENCY_REVISION_BOUND => Ok(Consistency::RevisionBound),
pb::Consistency::CONSISTENCY_SNAPSHOT_CONSISTENT => Ok(Consistency::SnapshotConsistent),
pb::Consistency::CONSISTENCY_ORDERED_PER_AGGREGATE => Ok(Consistency::OrderedPerAggregate),
pb::Consistency::CONSISTENCY_EVENTUAL_WITH_WATERMARK => {
Ok(Consistency::EventualWithWatermark)
}
pb::Consistency::CONSISTENCY_EXPLICITLY_EPHEMERAL => Ok(Consistency::ExplicitlyEphemeral),
pb::Consistency::CONSISTENCY_UNSPECIFIED => Err(malformed(
field,
"an operation declares its consistency class",
)),
}
}
impl From<Kernel<CommitDisposition>> for pb::CommitDisposition {
fn from(value: Kernel<CommitDisposition>) -> Self {
match value.0 {
CommitDisposition::New => Self::COMMIT_DISPOSITION_NEW,
CommitDisposition::Deduplicated => Self::COMMIT_DISPOSITION_DEDUPLICATED,
}
}
}
fn disposition(
field: &str,
value: buffa::EnumValue<pb::CommitDisposition>,
) -> Result<CommitDisposition, StateError> {
match known(field, value)? {
pb::CommitDisposition::COMMIT_DISPOSITION_NEW => Ok(CommitDisposition::New),
pb::CommitDisposition::COMMIT_DISPOSITION_DEDUPLICATED => {
Ok(CommitDisposition::Deduplicated)
}
pb::CommitDisposition::COMMIT_DISPOSITION_UNSPECIFIED => Err(malformed(
field,
"a receipt reports whether it committed or replayed",
)),
}
}
impl From<Kernel<PageCompleteness>> for pb::PageCompleteness {
fn from(value: Kernel<PageCompleteness>) -> Self {
match value.0 {
PageCompleteness::Complete => Self::PAGE_COMPLETENESS_COMPLETE,
PageCompleteness::Truncated => Self::PAGE_COMPLETENESS_TRUNCATED,
}
}
}
pub(crate) fn completeness(
field: &str,
value: buffa::EnumValue<pb::PageCompleteness>,
) -> Result<PageCompleteness, StateError> {
match known(field, value)? {
pb::PageCompleteness::PAGE_COMPLETENESS_COMPLETE => Ok(PageCompleteness::Complete),
pb::PageCompleteness::PAGE_COMPLETENESS_TRUNCATED => Ok(PageCompleteness::Truncated),
pb::PageCompleteness::PAGE_COMPLETENESS_UNSPECIFIED => Err(malformed(
field,
"a page reports whether records remain after it",
)),
}
}
impl From<Kernel<StreamEnd>> for pb::StreamEnd {
fn from(value: Kernel<StreamEnd>) -> Self {
match value.0 {
StreamEnd::More => Self::STREAM_END_MORE,
StreamEnd::Exhausted => Self::STREAM_END_EXHAUSTED,
StreamEnd::Drained => Self::STREAM_END_DRAINED,
}
}
}
pub(crate) fn stream_end(
field: &str,
value: buffa::EnumValue<pb::StreamEnd>,
) -> Result<StreamEnd, StateError> {
match known(field, value)? {
pb::StreamEnd::STREAM_END_MORE => Ok(StreamEnd::More),
pb::StreamEnd::STREAM_END_EXHAUSTED => Ok(StreamEnd::Exhausted),
pb::StreamEnd::STREAM_END_DRAINED => Ok(StreamEnd::Drained),
pb::StreamEnd::STREAM_END_UNSPECIFIED => {
Err(malformed(field, "a chunk reports why it stopped"))
}
}
}
impl From<Kernel<Delivery>> for pb::Delivery {
fn from(value: Kernel<Delivery>) -> Self {
match value.0 {
Delivery::AtLeastOnce => Self::DELIVERY_AT_LEAST_ONCE,
Delivery::AtMostOnce => Self::DELIVERY_AT_MOST_ONCE,
}
}
}
fn delivery(field: &str, value: buffa::EnumValue<pb::Delivery>) -> Result<Delivery, StateError> {
match known(field, value)? {
pb::Delivery::DELIVERY_AT_LEAST_ONCE => Ok(Delivery::AtLeastOnce),
pb::Delivery::DELIVERY_AT_MOST_ONCE => Ok(Delivery::AtMostOnce),
pb::Delivery::DELIVERY_UNSPECIFIED => {
Err(malformed(field, "a stream declares its delivery guarantee"))
}
}
}
impl From<Kernel<CompactedRange>> for pb::CompactedRange {
fn from(value: Kernel<CompactedRange>) -> Self {
match value.0 {
CompactedRange::NeverCompacted => Self::COMPACTED_RANGE_NEVER_COMPACTED,
CompactedRange::FailsOnCompactedCursor => {
Self::COMPACTED_RANGE_FAILS_ON_COMPACTED_CURSOR
}
}
}
}
fn compacted_range(
field: &str,
value: buffa::EnumValue<pb::CompactedRange>,
) -> Result<CompactedRange, StateError> {
match known(field, value)? {
pb::CompactedRange::COMPACTED_RANGE_NEVER_COMPACTED => Ok(CompactedRange::NeverCompacted),
pb::CompactedRange::COMPACTED_RANGE_FAILS_ON_COMPACTED_CURSOR => {
Ok(CompactedRange::FailsOnCompactedCursor)
}
pb::CompactedRange::COMPACTED_RANGE_UNSPECIFIED => Err(malformed(
field,
"a stream declares what a compacted cursor does",
)),
}
}
impl From<Kernel<Backpressure>> for pb::Backpressure {
fn from(value: Kernel<Backpressure>) -> Self {
match value.0 {
Backpressure::BlockProducer => Self::BACKPRESSURE_BLOCK_PRODUCER,
Backpressure::ShedSlowConsumer => Self::BACKPRESSURE_SHED_SLOW_CONSUMER,
}
}
}
fn backpressure(
field: &str,
value: buffa::EnumValue<pb::Backpressure>,
) -> Result<Backpressure, StateError> {
match known(field, value)? {
pb::Backpressure::BACKPRESSURE_BLOCK_PRODUCER => Ok(Backpressure::BlockProducer),
pb::Backpressure::BACKPRESSURE_SHED_SLOW_CONSUMER => Ok(Backpressure::ShedSlowConsumer),
pb::Backpressure::BACKPRESSURE_UNSPECIFIED => Err(malformed(
field,
"a stream declares what it does with a slow consumer",
)),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct CallContextVersion(u32);
impl CallContextVersion {
pub const CURRENT: Self = Self(2);
#[must_use]
pub const fn new(version: u32) -> Self {
Self(version)
}
#[must_use]
pub const fn get(self) -> u32 {
self.0
}
}
impl std::fmt::Display for CallContextVersion {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(formatter, "v{}", self.0)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DeclaredCall {
pub version: CallContextVersion,
pub audience: Audience,
pub budget: Duration,
pub cancelled: bool,
}
impl DeclaredCall {
#[must_use]
pub const fn bounded(audience: Audience, remaining: Duration) -> Self {
Self::live(audience, remaining)
}
#[must_use]
pub const fn live(audience: Audience, budget: Duration) -> Self {
Self {
version: CallContextVersion::CURRENT,
audience,
budget,
cancelled: false,
}
}
#[must_use]
pub const fn withdrawn(mut self) -> Self {
self.cancelled = true;
self
}
#[must_use]
pub fn origin_relative_context(&self) -> CallContext {
CallContext::from_origin_relative_budget(self.budget, cancelled_token(self.cancelled))
}
#[must_use]
pub const fn remaining_budget(&self) -> Duration {
self.budget
}
}
impl From<Kernel<&DeclaredCall>> for pb::CallContext {
fn from(value: Kernel<&DeclaredCall>) -> Self {
let call = value.0;
Self {
protocol_version: call.version.get(),
audience: call.audience.as_str().to_owned(),
deadline_nanos: nanos_from_duration(call.budget),
cancelled: call.cancelled,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl From<pb::CallContext> for Kernel<DeclaredCall> {
fn from(value: pb::CallContext) -> Self {
Self(DeclaredCall {
version: CallContextVersion::new(value.protocol_version),
audience: Audience::new(value.audience),
budget: duration_from_nanos(value.deadline_nanos),
cancelled: value.cancelled,
})
}
}
pub fn declared_call(
context: impl Into<Option<pb::CallContext>>,
) -> Result<DeclaredCall, StateError> {
let message = context.into().ok_or_else(|| {
malformed(
"context",
"every call declares its version, audience, and budget",
)
})?;
Ok(Kernel::<DeclaredCall>::from(message).into_inner())
}
impl From<Kernel<Precondition>> for pb::Precondition {
fn from(value: Kernel<Precondition>) -> Self {
use pb::__buffa::oneof::precondition::Kind;
let kind = match value.0 {
Precondition::Unconditional => Kind::from(pb::Unconditional::default()),
Precondition::NoExistingState => Kind::from(pb::NoExistingState::default()),
Precondition::Revision(revision) => Kind::Revision(revision.get()),
Precondition::JournalHead(position) => Kind::JournalHead(position.get()),
};
Self {
kind: Some(kind),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::Precondition> for Kernel<Precondition> {
type Error = StateError;
fn try_from(value: pb::Precondition) -> Result<Self, Self::Error> {
use pb::__buffa::oneof::precondition::Kind;
let precondition = match value.kind {
Some(Kind::Unconditional(_)) => Precondition::Unconditional,
Some(Kind::NoExistingState(_)) => Precondition::NoExistingState,
Some(Kind::Revision(revision)) => Precondition::Revision(Revision::new(revision)),
Some(Kind::JournalHead(position)) => {
Precondition::JournalHead(JournalPosition::new(position))
}
None => {
return Err(malformed(
"precondition",
"a precondition names what durable state the command requires",
));
}
};
Ok(Self(precondition))
}
}
impl From<Kernel<ObservedState>> for pb::ObservedState {
fn from(value: Kernel<ObservedState>) -> Self {
use pb::__buffa::oneof::observed_state::Kind;
let kind = match value.0 {
ObservedState::Absent => Kind::from(pb::AbsentState::default()),
ObservedState::Revision(revision) => Kind::Revision(revision.get()),
ObservedState::JournalHead(position) => Kind::JournalHead(position.get()),
};
Self {
kind: Some(kind),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ObservedState> for Kernel<ObservedState> {
type Error = StateError;
fn try_from(value: pb::ObservedState) -> Result<Self, Self::Error> {
use pb::__buffa::oneof::observed_state::Kind;
let observed = match value.kind {
Some(Kind::Absent(_)) => ObservedState::Absent,
Some(Kind::Revision(revision)) => ObservedState::Revision(Revision::new(revision)),
Some(Kind::JournalHead(position)) => {
ObservedState::JournalHead(JournalPosition::new(position))
}
None => {
return Err(malformed(
"observed",
"a conflict names the durable state it was checked against",
));
}
};
Ok(Self(observed))
}
}
impl From<Kernel<&Cursor>> for pb::Cursor {
fn from(value: Kernel<&Cursor>) -> Self {
Self {
position: value.0.position().get(),
snapshot: value
.0
.snapshot()
.map(|snapshot| snapshot.as_str().to_owned()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl From<pb::Cursor> for Kernel<Cursor> {
fn from(value: pb::Cursor) -> Self {
let position = JournalPosition::new(value.position);
Self(value.snapshot.map_or_else(
|| Cursor::at(position),
|snapshot| Cursor::in_snapshot(SnapshotId::new(snapshot), position),
))
}
}
impl From<Kernel<&CommitEvidence>> for pb::CommitEvidence {
fn from(value: Kernel<&CommitEvidence>) -> Self {
let evidence = value.0;
Self {
revision: evidence.revision().map(Revision::get),
position: evidence.position().map(JournalPosition::get),
root: evidence.root().map(|root| root.as_bytes().to_vec()),
snapshot: evidence
.snapshot()
.map(|snapshot| snapshot.as_str().to_owned()),
cursor: evidence
.cursor()
.map_or_else(buffa::MessageField::default, |cursor| {
buffa::MessageField::some(pb::Cursor::from(Kernel(cursor)))
}),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::CommitEvidence> for Kernel<CommitEvidence> {
type Error = StateError;
fn try_from(value: pb::CommitEvidence) -> Result<Self, Self::Error> {
let mut evidence = CommitEvidence::new();
if let Some(revision) = value.revision {
evidence = evidence.with_revision(Revision::new(revision));
}
if let Some(position) = value.position {
evidence = evidence.with_position(JournalPosition::new(position));
}
if let Some(root) = value.root {
let bytes = fixed_bytes::<{ CommitRoot::LEN }>("root", &root)?;
evidence = evidence.with_root(CommitRoot::from_bytes(bytes));
}
if let Some(snapshot) = value.snapshot {
evidence = evidence.with_snapshot(SnapshotId::new(snapshot));
}
if let Some(cursor) = value.cursor.into_option() {
evidence = evidence.with_cursor(Kernel::<Cursor>::from(cursor).into_inner());
}
Ok(Self(evidence))
}
}
impl From<Kernel<&Receipt>> for pb::Receipt {
fn from(value: Kernel<&Receipt>) -> Self {
let receipt = value.0;
Self {
command_id: receipt.command_id().as_str().to_owned(),
family: receipt.family().as_str().to_owned(),
digest: receipt.digest().as_bytes().to_vec(),
disposition: pb::CommitDisposition::from(Kernel(receipt.disposition())).into(),
evidence: buffa::MessageField::some(pb::CommitEvidence::from(Kernel(
receipt.evidence(),
))),
consistency: pb::Consistency::from(Kernel(receipt.consistency())).into(),
fence: receipt.fence().map(FencingToken::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::Receipt> for Kernel<Receipt> {
type Error = StateError;
fn try_from(value: pb::Receipt) -> Result<Self, Self::Error> {
let disposition = disposition("disposition", value.disposition)?;
let consistency = consistency("consistency", value.consistency)?;
let evidence = Kernel::<CommitEvidence>::try_from(required(
"evidence",
"a receipt carries what its commit produced",
value.evidence,
)?)?
.into_inner();
let digest = ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
"digest",
&value.digest,
)?);
let mut metadata = CommandMetadata::new(
CommandId::new(value.command_id),
OperationFamily::new(value.family),
digest,
CommandScope::new(
AggregateId::new(family::AGGREGATE),
PartitionId::new(family::AGGREGATE),
NamespaceId::new(family::NAMESPACE),
),
CommandEnvelope::new(
Purpose::new(family::PURPOSE),
Audience::new(family::AUDIENCE),
ResourceBounds::new(family::MAX_PAYLOAD_BYTES, family::MAX_RECORDS_PER_COMMAND),
),
);
if let Some(fence) = value.fence {
metadata = metadata.with_fence(FencingToken::new(fence));
}
let receipt = Receipt::committed(&metadata, evidence, consistency);
Ok(Self(match disposition {
CommitDisposition::New => receipt,
CommitDisposition::Deduplicated => receipt.as_replay(),
}))
}
}
impl From<Kernel<&ReadStart>> for pb::ReadStart {
fn from(value: Kernel<&ReadStart>) -> Self {
use pb::__buffa::oneof::read_start::Start;
let start = match value.0 {
ReadStart::Snapshot(snapshot) => Start::Snapshot(snapshot.as_str().to_owned()),
ReadStart::Resume(cursor) => Start::from(pb::Cursor::from(Kernel(cursor))),
};
Self {
start: Some(start),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ReadStart> for Kernel<ReadStart> {
type Error = StateError;
fn try_from(value: pb::ReadStart) -> Result<Self, Self::Error> {
use pb::__buffa::oneof::read_start::Start;
let start = match value.start {
Some(Start::Snapshot(snapshot)) => ReadStart::Snapshot(SnapshotId::new(snapshot)),
Some(Start::Resume(cursor)) => {
ReadStart::Resume(Kernel::<Cursor>::from(*cursor).into_inner())
}
None => return Err(malformed("start", "a bounded read names where it begins")),
};
Ok(Self(start))
}
}
impl From<Kernel<&PageRequest>> for pb::PageRequest {
fn from(value: Kernel<&PageRequest>) -> Self {
Self {
start: buffa::MessageField::some(pb::ReadStart::from(Kernel(value.0.start()))),
limit: value.0.limit(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::PageRequest> for Kernel<PageRequest> {
type Error = StateError;
fn try_from(value: pb::PageRequest) -> Result<Self, Self::Error> {
let start = Kernel::<ReadStart>::try_from(required(
"start",
"a bounded read names where it begins",
value.start,
)?)?
.into_inner();
Ok(Self(PageRequest::new(start, value.limit)))
}
}
impl From<Kernel<&StreamRequest>> for pb::StreamRequest {
fn from(value: Kernel<&StreamRequest>) -> Self {
Self {
start: buffa::MessageField::some(pb::ReadStart::from(Kernel(value.0.start()))),
max_chunk_records: value.0.max_chunk_records(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::StreamRequest> for Kernel<StreamRequest> {
type Error = StateError;
fn try_from(value: pb::StreamRequest) -> Result<Self, Self::Error> {
let start = Kernel::<ReadStart>::try_from(required(
"start",
"a bounded read names where it begins",
value.start,
)?)?
.into_inner();
Ok(Self(StreamRequest::new(start, value.max_chunk_records)))
}
}
impl From<Kernel<&SyntheticRecord>> for pb::SyntheticRecord {
fn from(value: Kernel<&SyntheticRecord>) -> Self {
Self {
position: value.0.position().get(),
command_id: value.0.command_id().as_str().to_owned(),
amount: value.0.amount(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl From<pb::SyntheticRecord> for Kernel<SyntheticRecord> {
fn from(value: pb::SyntheticRecord) -> Self {
Self(SyntheticRecord::new(
JournalPosition::new(value.position),
CommandId::new(value.command_id),
value.amount,
))
}
}
impl From<Kernel<&Page<SyntheticRecord>>> for pb::Page {
fn from(value: Kernel<&Page<SyntheticRecord>>) -> Self {
let page = value.0;
Self {
records: page
.records()
.iter()
.map(|record| pb::SyntheticRecord::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(),
watermark: page.watermark().map(Watermark::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::Page> for Kernel<Page<SyntheticRecord>> {
type Error = StateError;
fn try_from(value: pb::Page) -> Result<Self, Self::Error> {
let completeness = completeness("completeness", value.completeness)?;
let consistency = consistency("consistency", value.consistency)?;
let records = value
.records
.into_iter()
.map(|record| Kernel::<SyntheticRecord>::from(record).into_inner())
.collect();
let next = value
.next
.into_option()
.map(|cursor| Kernel::<Cursor>::from(cursor).into_inner());
let page = Page::new(records, next, completeness, consistency);
Ok(Self(match value.watermark {
Some(watermark) => page.with_watermark(Watermark::new(watermark)),
None => page,
}))
}
}
impl From<Kernel<&StreamChunk<SyntheticRecord>>> for pb::StreamChunk {
fn from(value: Kernel<&StreamChunk<SyntheticRecord>>) -> Self {
let chunk = value.0;
Self {
records: chunk
.records()
.iter()
.map(|record| pb::SyntheticRecord::from(Kernel(record)))
.collect(),
next: chunk
.next_cursor()
.map_or_else(buffa::MessageField::default, |cursor| {
buffa::MessageField::some(pb::Cursor::from(Kernel(cursor)))
}),
end: pb::StreamEnd::from(Kernel(chunk.end())).into(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::StreamChunk> for Kernel<StreamChunk<SyntheticRecord>> {
type Error = StateError;
fn try_from(value: pb::StreamChunk) -> Result<Self, Self::Error> {
let end = stream_end("end", value.end)?;
let records = value
.records
.into_iter()
.map(|record| Kernel::<SyntheticRecord>::from(record).into_inner())
.collect();
let next = value
.next
.into_option()
.map(|cursor| Kernel::<Cursor>::from(cursor).into_inner());
Ok(Self(StreamChunk::new(records, next, end)))
}
}
impl From<Kernel<StreamContract>> for pb::StreamContract {
fn from(value: Kernel<StreamContract>) -> Self {
let contract = value.0;
Self {
consistency: pb::Consistency::from(Kernel(contract.consistency())).into(),
delivery: pb::Delivery::from(Kernel(contract.delivery())).into(),
max_chunk_records: contract.max_chunk_records(),
compacted_range: pb::CompactedRange::from(Kernel(contract.compacted_range())).into(),
backpressure: pb::Backpressure::from(Kernel(contract.backpressure())).into(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::StreamContract> for Kernel<StreamContract> {
type Error = StateError;
fn try_from(value: pb::StreamContract) -> Result<Self, Self::Error> {
Ok(Self(StreamContract::new(
consistency("consistency", value.consistency)?,
delivery("delivery", value.delivery)?,
value.max_chunk_records,
compacted_range("compacted_range", value.compacted_range)?,
backpressure("backpressure", value.backpressure)?,
)))
}
}
impl From<Kernel<SyntheticOp>> for pb::SyntheticOp {
fn from(value: Kernel<SyntheticOp>) -> Self {
use pb::__buffa::oneof::synthetic_op::Kind;
let unknown = buffa::UnknownFields::default;
let kind = match value.0 {
SyntheticOp::Append { amount } => Kind::from(pb::SyntheticAppend {
amount,
__buffa_unknown_fields: unknown(),
}),
SyntheticOp::Oversized {
amount,
payload_bytes,
} => Kind::from(pb::SyntheticOversized {
amount,
payload_bytes,
__buffa_unknown_fields: unknown(),
}),
SyntheticOp::AppendThenLoseResponse { amount } => {
Kind::from(pb::SyntheticAppendThenLoseResponse {
amount,
__buffa_unknown_fields: unknown(),
})
}
SyntheticOp::AppendThenObserveCancellation { amount } => {
Kind::from(pb::SyntheticAppendThenObserveCancellation {
amount,
__buffa_unknown_fields: unknown(),
})
}
};
Self {
kind: Some(kind),
__buffa_unknown_fields: unknown(),
}
}
}
impl TryFrom<pb::SyntheticOp> for Kernel<SyntheticOp> {
type Error = StateError;
fn try_from(value: pb::SyntheticOp) -> Result<Self, Self::Error> {
use pb::__buffa::oneof::synthetic_op::Kind;
let op = match value.kind {
Some(Kind::Append(op)) => SyntheticOp::Append { amount: op.amount },
Some(Kind::Oversized(op)) => SyntheticOp::Oversized {
amount: op.amount,
payload_bytes: op.payload_bytes,
},
Some(Kind::AppendThenLoseResponse(op)) => {
SyntheticOp::AppendThenLoseResponse { amount: op.amount }
}
Some(Kind::AppendThenObserveCancellation(op)) => {
SyntheticOp::AppendThenObserveCancellation { amount: op.amount }
}
None => {
return Err(malformed(
"op",
"a command names what it asks the module to do",
));
}
};
Ok(Self(op))
}
}
impl From<Kernel<&SyntheticCommand>> for pb::SyntheticCommand {
fn from(value: Kernel<&SyntheticCommand>) -> Self {
let metadata = value.0.metadata();
Self {
command_id: metadata.command_id().as_str().to_owned(),
op: buffa::MessageField::some(pb::SyntheticOp::from(Kernel(value.0.op()))),
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::SyntheticCommand> for Kernel<SyntheticCommand> {
type Error = StateError;
fn try_from(value: pb::SyntheticCommand) -> Result<Self, Self::Error> {
let op = Kernel::<SyntheticOp>::try_from(required(
"op",
"a command names what it asks the module to do",
value.op,
)?)?
.into_inner();
let precondition = Kernel::<Precondition>::try_from(required(
"precondition",
"a precondition names what durable state the command requires",
value.precondition,
)?)?
.into_inner();
let mut command =
SyntheticCommand::new(value.command_id.as_str(), op).with_precondition(precondition);
if let Some(fence) = value.fence {
command = command.with_fence(FencingToken::new(fence));
}
Ok(Self(command))
}
}
#[must_use]
pub const fn duration_from_nanos(nanos: u64) -> Duration {
Duration::from_nanos(nanos)
}
#[must_use]
pub fn nanos_from_duration(budget: Duration) -> u64 {
u64::try_from(budget.as_nanos()).unwrap_or(u64::MAX)
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs, clippy::unwrap_used)]
use super::*;
use polyc_state::{conformance::ConformanceAdapter, memory::MemoryState};
fn live() -> DeclaredCall {
DeclaredCall::live(Audience::new(family::AUDIENCE), Duration::from_secs(30))
}
#[test]
fn a_call_context_round_trips() {
let declared = live();
let wire = pb::CallContext::from(Kernel(&declared));
assert_eq!(
wire.protocol_version, 2,
"CallContext field 1 versions wire semantics independently of persisted commands"
);
assert_eq!(Kernel::<DeclaredCall>::from(wire).into_inner(), declared);
assert!(!declared.origin_relative_context().is_cancelled());
}
#[test]
fn a_withdrawn_call_arrives_withdrawn() {
let declared = live().withdrawn();
let wire = pb::CallContext::from(Kernel(&declared));
assert!(wire.cancelled);
assert!(declared.origin_relative_context().is_cancelled());
}
#[test]
fn a_command_round_trips_with_its_derived_digest() {
let command = SyntheticCommand::new(
"cmd-1",
SyntheticOp::Oversized {
amount: 3,
payload_bytes: 99,
},
)
.with_precondition(Precondition::Revision(Revision::new(4)))
.with_fence(FencingToken::new(7));
let back =
Kernel::<SyntheticCommand>::try_from(pb::SyntheticCommand::from(Kernel(&command)))
.unwrap()
.into_inner();
assert_eq!(back, command);
assert_eq!(back.metadata().digest(), command.metadata().digest());
}
#[test]
fn a_receipt_round_trips_field_for_field() {
let mut state = MemoryState::new();
let declared = live();
let receipt = state
.submit(
SyntheticCommand::new("cmd-1", SyntheticOp::Append { amount: 2 })
.with_fence(FencingToken::new(3)),
&declared.origin_relative_context(),
)
.unwrap();
let back = Kernel::<Receipt>::try_from(pb::Receipt::from(Kernel(&receipt)))
.unwrap()
.into_inner();
assert_eq!(back, receipt);
let replay = receipt.as_replay();
let back_replay = Kernel::<Receipt>::try_from(pb::Receipt::from(Kernel(&replay)))
.unwrap()
.into_inner();
assert_eq!(back_replay, replay);
assert!(back_replay.is_deduplicated());
}
#[test]
fn a_page_and_a_chunk_round_trip() {
let mut state = MemoryState::new();
let declared = live();
for amount in 1..=3_u64 {
state
.submit(
SyntheticCommand::new(&format!("cmd-{amount}"), SyntheticOp::Append { amount }),
&declared.origin_relative_context(),
)
.unwrap();
}
let snapshot = state
.create_snapshot(&declared.origin_relative_context())
.unwrap();
let page = state
.read_page(
PageRequest::new(ReadStart::Snapshot(snapshot.clone()), 2),
&declared.origin_relative_context(),
)
.unwrap();
assert_eq!(
Kernel::<Page<SyntheticRecord>>::try_from(pb::Page::from(Kernel(&page)))
.unwrap()
.into_inner(),
page
);
let chunk = state
.read_chunk(
StreamRequest::new(ReadStart::Snapshot(snapshot), 2),
&declared.origin_relative_context(),
)
.unwrap();
assert_eq!(
Kernel::<StreamChunk<SyntheticRecord>>::try_from(pb::StreamChunk::from(Kernel(&chunk)))
.unwrap()
.into_inner(),
chunk
);
let contract = state.stream_contract();
assert_eq!(
Kernel::<StreamContract>::try_from(pb::StreamContract::from(Kernel(contract)))
.unwrap()
.into_inner(),
contract
);
}
#[test]
fn read_requests_round_trip() {
let page = PageRequest::new(ReadStart::Snapshot(SnapshotId::new("snap-1")), 3);
assert_eq!(
Kernel::<PageRequest>::try_from(pb::PageRequest::from(Kernel(&page)))
.unwrap()
.into_inner(),
page
);
let chunk = StreamRequest::new(
ReadStart::Resume(Cursor::in_snapshot(
SnapshotId::new("snap-1"),
JournalPosition::new(2),
)),
2,
);
assert_eq!(
Kernel::<StreamRequest>::try_from(pb::StreamRequest::from(Kernel(&chunk)))
.unwrap()
.into_inner(),
chunk
);
}
#[test]
fn a_short_digest_is_malformed() {
let wire = pb::Receipt {
command_id: "cmd-1".to_owned(),
family: family::FAMILY.to_owned(),
digest: vec![1, 2, 3],
disposition: pb::CommitDisposition::COMMIT_DISPOSITION_NEW.into(),
evidence: buffa::MessageField::some(pb::CommitEvidence::default()),
consistency: pb::Consistency::CONSISTENCY_ORDERED_PER_AGGREGATE.into(),
fence: None,
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let error = Kernel::<Receipt>::try_from(wire).unwrap_err();
assert!(matches!(error, StateError::Malformed { ref field, .. } if field == "digest"));
}
#[test]
fn an_unset_oneof_is_malformed() {
assert!(matches!(
Kernel::<Precondition>::try_from(pb::Precondition::default()).unwrap_err(),
StateError::Malformed { .. }
));
assert!(matches!(
Kernel::<ObservedState>::try_from(pb::ObservedState::default()).unwrap_err(),
StateError::Malformed { .. }
));
assert!(matches!(
Kernel::<ReadStart>::try_from(pb::ReadStart::default()).unwrap_err(),
StateError::Malformed { .. }
));
assert!(matches!(
Kernel::<SyntheticOp>::try_from(pb::SyntheticOp::default()).unwrap_err(),
StateError::Malformed { .. }
));
}
#[test]
fn an_unspecified_enum_is_malformed() {
assert!(matches!(
consistency(
"consistency",
pb::Consistency::CONSISTENCY_UNSPECIFIED.into()
)
.unwrap_err(),
StateError::Malformed { .. }
));
assert!(matches!(
stream_end("end", pb::StreamEnd::STREAM_END_UNSPECIFIED.into()).unwrap_err(),
StateError::Malformed { .. }
));
assert!(matches!(
delivery("delivery", pb::Delivery::DELIVERY_UNSPECIFIED.into()).unwrap_err(),
StateError::Malformed { .. }
));
assert!(matches!(
backpressure(
"backpressure",
pb::Backpressure::BACKPRESSURE_UNSPECIFIED.into()
)
.unwrap_err(),
StateError::Malformed { .. }
));
assert!(matches!(
compacted_range(
"compacted_range",
pb::CompactedRange::COMPACTED_RANGE_UNSPECIFIED.into()
)
.unwrap_err(),
StateError::Malformed { .. }
));
assert!(matches!(
completeness(
"completeness",
pb::PageCompleteness::PAGE_COMPLETENESS_UNSPECIFIED.into()
)
.unwrap_err(),
StateError::Malformed { .. }
));
assert!(matches!(
disposition(
"disposition",
pb::CommitDisposition::COMMIT_DISPOSITION_UNSPECIFIED.into()
)
.unwrap_err(),
StateError::Malformed { .. }
));
}
#[test]
fn an_unknown_enum_value_is_malformed() {
assert!(matches!(
consistency("consistency", buffa::EnumValue::Unknown(99)).unwrap_err(),
StateError::Malformed { .. }
));
}
#[test]
fn a_missing_call_context_is_malformed() {
assert!(matches!(
declared_call(buffa::MessageField::<_, buffa::Inline<_>>::default()).unwrap_err(),
StateError::Malformed { ref field, .. } if field == "context"
));
}
#[test]
fn preconditions_and_observed_state_round_trip() {
for precondition in [
Precondition::Unconditional,
Precondition::NoExistingState,
Precondition::Revision(Revision::new(9)),
Precondition::JournalHead(JournalPosition::new(4)),
] {
assert_eq!(
Kernel::<Precondition>::try_from(pb::Precondition::from(Kernel(precondition)))
.unwrap()
.into_inner(),
precondition
);
}
for observed in [
ObservedState::Absent,
ObservedState::Revision(Revision::new(2)),
ObservedState::JournalHead(JournalPosition::new(6)),
] {
assert_eq!(
Kernel::<ObservedState>::try_from(pb::ObservedState::from(Kernel(observed)))
.unwrap()
.into_inner(),
observed
);
}
}
#[test]
fn durations_convert_both_ways() {
assert_eq!(nanos_from_duration(Duration::from_millis(3)), 3_000_000);
assert_eq!(duration_from_nanos(3_000_000), Duration::from_millis(3));
}
}