use std::time::Duration;
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
claims::{
AcquireClaim, AttemptHistory, AttemptOutcome, AttemptRecord, ClaimLease, ClaimStatus,
CompleteClaim, Granted, ReleaseClaim, RenewClaim, WorkDisposition, WorkOutcome,
},
command::{
CommandEnvelope, CommandMetadata, CommandScope, FencingToken, Precondition, ResourceBounds,
},
deadline::{Deadline, MonotonicInstant},
digest::ContentDigest,
error::StateError,
id::{
AggregateId, AttemptId, Audience, CommandId, NamespaceId, OperationFamily, OwnerId,
PartitionId, ProtocolVersion, Purpose, WorkId,
},
journal::JournalRecordDraft,
revision::{CommitRoot, JournalHead, JournalPosition, Revision},
};
use crate::wire::{Kernel, fixed_bytes, known, malformed, required};
pub(crate) fn scope_to_wire(scope: &CommandScope) -> pb::ClaimScope {
pb::ClaimScope {
aggregate: scope.aggregate().as_str().to_owned(),
partition: scope.partition().as_str().to_owned(),
namespace: scope.namespace().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn scope_from_wire(value: pb::ClaimScope) -> CommandScope {
CommandScope::new(
AggregateId::new(value.aggregate),
PartitionId::new(value.partition),
NamespaceId::new(value.namespace),
)
}
fn metadata_to_wire(value: &CommandMetadata) -> pb::ClaimCommandMetadata {
pb::ClaimCommandMetadata {
command_id: value.command_id().as_str().to_owned(),
digest: value.digest().as_bytes().to_vec(),
scope: buffa::MessageField::some(scope_to_wire(value.scope())),
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(FencingToken::get),
protocol_version: value.protocol_version().get(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn metadata_from_wire(value: pb::ClaimCommandMetadata) -> Result<CommandMetadata, StateError> {
let mut metadata = CommandMetadata::new(
CommandId::new(value.command_id),
OperationFamily::new(polyc_state::claims::FAMILY),
ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
"digest",
&value.digest,
)?),
scope_from_wire(required(
"scope",
"a claim command names its work item",
value.scope,
)?),
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 claim command declares its precondition",
value.precondition,
)?)?
.into_inner(),
)
.with_protocol_version(ProtocolVersion::new(value.protocol_version));
if let Some(fence) = value.fence {
metadata = metadata.with_fence(FencingToken::new(fence));
}
Ok(metadata)
}
fn changes_to_wire(changes: &[JournalRecordDraft]) -> Vec<pb::JournalRecordDraft> {
changes
.iter()
.map(|change| pb::JournalRecordDraft::from(Kernel(change)))
.collect()
}
fn changes_from_wire(
changes: Vec<pb::JournalRecordDraft>,
) -> Result<Vec<JournalRecordDraft>, StateError> {
changes
.into_iter()
.map(|change| Kernel::<JournalRecordDraft>::try_from(change).map(Kernel::into_inner))
.collect()
}
pub(crate) fn acquire_to_wire(command: &AcquireClaim) -> pb::AcquireClaimCommand {
pb::AcquireClaimCommand {
metadata: buffa::MessageField::some(metadata_to_wire(command.metadata())),
work: command.work().as_str().to_owned(),
attempt: command.attempt().as_str().to_owned(),
owner: command.owner().as_str().to_owned(),
lease_nanos: u64::try_from(command.lease().as_nanos()).unwrap_or(u64::MAX),
changes: changes_to_wire(command.changes()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn acquire_from_wire(
value: pb::AcquireClaimCommand,
) -> Result<AcquireClaim, StateError> {
Ok(AcquireClaim::new(
metadata_from_wire(required(
"metadata",
"an acquire carries metadata",
value.metadata,
)?)?,
WorkId::new(value.work),
AttemptId::new(value.attempt),
OwnerId::new(value.owner),
Duration::from_nanos(value.lease_nanos),
changes_from_wire(value.changes)?,
))
}
pub(crate) fn renew_to_wire(command: &RenewClaim) -> pb::RenewClaimCommand {
pb::RenewClaimCommand {
metadata: buffa::MessageField::some(metadata_to_wire(command.metadata())),
work: command.work().as_str().to_owned(),
attempt: command.attempt().as_str().to_owned(),
owner: command.owner().as_str().to_owned(),
lease_nanos: u64::try_from(command.lease().as_nanos()).unwrap_or(u64::MAX),
changes: changes_to_wire(command.changes()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn renew_from_wire(value: pb::RenewClaimCommand) -> Result<RenewClaim, StateError> {
Ok(RenewClaim::new(
metadata_from_wire(required(
"metadata",
"a renewal carries metadata",
value.metadata,
)?)?,
WorkId::new(value.work),
AttemptId::new(value.attempt),
OwnerId::new(value.owner),
Duration::from_nanos(value.lease_nanos),
changes_from_wire(value.changes)?,
))
}
const fn outcome_to_wire(value: WorkOutcome) -> pb::ClaimWorkOutcome {
match value {
WorkOutcome::Succeeded => pb::ClaimWorkOutcome::Succeeded,
WorkOutcome::Failed => pb::ClaimWorkOutcome::Failed,
}
}
fn outcome_from_wire(
field: &str,
value: buffa::EnumValue<pb::ClaimWorkOutcome>,
) -> Result<WorkOutcome, StateError> {
match known(field, value)? {
pb::ClaimWorkOutcome::Succeeded => Ok(WorkOutcome::Succeeded),
pb::ClaimWorkOutcome::Failed => Ok(WorkOutcome::Failed),
pb::ClaimWorkOutcome::Unspecified => {
Err(malformed(field, "a completion names its outcome"))
}
}
}
pub(crate) fn complete_to_wire(command: &CompleteClaim) -> pb::CompleteClaimCommand {
pb::CompleteClaimCommand {
metadata: buffa::MessageField::some(metadata_to_wire(command.metadata())),
work: command.work().as_str().to_owned(),
attempt: command.attempt().as_str().to_owned(),
owner: command.owner().as_str().to_owned(),
outcome: buffa::EnumValue::from(outcome_to_wire(command.outcome())),
changes: changes_to_wire(command.changes()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn complete_from_wire(
value: pb::CompleteClaimCommand,
) -> Result<CompleteClaim, StateError> {
Ok(CompleteClaim::new(
metadata_from_wire(required(
"metadata",
"a completion carries metadata",
value.metadata,
)?)?,
WorkId::new(value.work),
AttemptId::new(value.attempt),
OwnerId::new(value.owner),
outcome_from_wire("outcome", value.outcome)?,
changes_from_wire(value.changes)?,
))
}
pub(crate) fn release_to_wire(command: &ReleaseClaim) -> pb::ReleaseClaimCommand {
pb::ReleaseClaimCommand {
metadata: buffa::MessageField::some(metadata_to_wire(command.metadata())),
work: command.work().as_str().to_owned(),
attempt: command.attempt().as_str().to_owned(),
owner: command.owner().as_str().to_owned(),
changes: changes_to_wire(command.changes()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn release_from_wire(
value: pb::ReleaseClaimCommand,
) -> Result<ReleaseClaim, StateError> {
Ok(ReleaseClaim::new(
metadata_from_wire(required(
"metadata",
"a release carries metadata",
value.metadata,
)?)?,
WorkId::new(value.work),
AttemptId::new(value.attempt),
OwnerId::new(value.owner),
changes_from_wire(value.changes)?,
))
}
fn lease_to_wire(value: &ClaimLease) -> pb::ClaimLease {
pb::ClaimLease {
work: value.work().as_str().to_owned(),
attempt: value.attempt().as_str().to_owned(),
owner: value.owner().as_str().to_owned(),
fence: value.fence().get(),
deadline_nanos: value.deadline().instant().as_nanos(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn lease_from_wire(value: pb::ClaimLease) -> ClaimLease {
ClaimLease::new(
WorkId::new(value.work),
AttemptId::new(value.attempt),
OwnerId::new(value.owner),
FencingToken::new(value.fence),
Deadline::at(MonotonicInstant::from_nanos(value.deadline_nanos)),
)
}
pub(crate) fn granted_to_wire(value: &Granted) -> pb::GrantedClaim {
pb::GrantedClaim {
lease: buffa::MessageField::some(lease_to_wire(value.lease())),
receipt: buffa::MessageField::some(pb::Receipt::from(Kernel(value.receipt()))),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn granted_from_wire(value: pb::GrantedClaim) -> Result<Granted, StateError> {
Ok(Granted::new(
lease_from_wire(required(
"lease",
"a granted claim carries its lease",
value.lease,
)?),
Kernel::<polyc_state::receipt::Receipt>::try_from(required(
"receipt",
"a granted claim carries its durable receipt",
value.receipt,
)?)?
.into_inner(),
))
}
const fn attempt_outcome_to_wire(value: AttemptOutcome) -> pb::ClaimAttemptOutcome {
match value {
AttemptOutcome::InProgress => pb::ClaimAttemptOutcome::InProgress,
AttemptOutcome::Completed(WorkOutcome::Succeeded) => pb::ClaimAttemptOutcome::Succeeded,
AttemptOutcome::Completed(WorkOutcome::Failed) => pb::ClaimAttemptOutcome::Failed,
AttemptOutcome::Released => pb::ClaimAttemptOutcome::Released,
AttemptOutcome::Expired => pb::ClaimAttemptOutcome::Expired,
}
}
fn attempt_outcome_from_wire(
value: buffa::EnumValue<pb::ClaimAttemptOutcome>,
) -> Result<AttemptOutcome, StateError> {
match known("attempt.outcome", value)? {
pb::ClaimAttemptOutcome::InProgress => Ok(AttemptOutcome::InProgress),
pb::ClaimAttemptOutcome::Succeeded => Ok(AttemptOutcome::Completed(WorkOutcome::Succeeded)),
pb::ClaimAttemptOutcome::Failed => Ok(AttemptOutcome::Completed(WorkOutcome::Failed)),
pb::ClaimAttemptOutcome::Released => Ok(AttemptOutcome::Released),
pb::ClaimAttemptOutcome::Expired => Ok(AttemptOutcome::Expired),
pb::ClaimAttemptOutcome::Unspecified => Err(malformed(
"attempt.outcome",
"an attempt record names its outcome",
)),
}
}
fn record_to_wire(value: &AttemptRecord) -> pb::ClaimAttemptRecord {
pb::ClaimAttemptRecord {
attempt: value.attempt().as_str().to_owned(),
owner: value.owner().as_str().to_owned(),
fence: value.fence().get(),
granted_at_nanos: value.granted_at().as_nanos(),
deadline_nanos: value.deadline().instant().as_nanos(),
outcome: buffa::EnumValue::from(attempt_outcome_to_wire(value.outcome())),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn record_from_wire(value: pb::ClaimAttemptRecord) -> Result<AttemptRecord, StateError> {
Ok(AttemptRecord::new(
AttemptId::new(value.attempt),
OwnerId::new(value.owner),
FencingToken::new(value.fence),
MonotonicInstant::from_nanos(value.granted_at_nanos),
Deadline::at(MonotonicInstant::from_nanos(value.deadline_nanos)),
attempt_outcome_from_wire(value.outcome)?,
))
}
pub(crate) fn history_to_wire(value: &AttemptHistory) -> pb::ClaimAttemptHistory {
pb::ClaimAttemptHistory {
work: value.work().as_str().to_owned(),
records: value.records().iter().map(record_to_wire).collect(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn history_from_wire(
value: pb::ClaimAttemptHistory,
) -> Result<AttemptHistory, StateError> {
Ok(AttemptHistory::new(
WorkId::new(value.work),
value
.records
.into_iter()
.map(record_from_wire)
.collect::<Result<_, _>>()?,
))
}
const fn disposition_to_wire(value: WorkDisposition) -> pb::ClaimDisposition {
match value {
WorkDisposition::Unclaimed => pb::ClaimDisposition::Unclaimed,
WorkDisposition::Claimed => pb::ClaimDisposition::Claimed,
WorkDisposition::Completed(WorkOutcome::Succeeded) => pb::ClaimDisposition::Succeeded,
WorkDisposition::Completed(WorkOutcome::Failed) => pb::ClaimDisposition::Failed,
}
}
pub(crate) fn status_to_wire(value: &ClaimStatus) -> pb::ClaimStatus {
pb::ClaimStatus {
work: value.work().as_str().to_owned(),
revision: value.revision().get(),
fence: value.fence().get(),
holder: value
.holder()
.map_or_else(buffa::MessageField::none, |holder| {
buffa::MessageField::some(lease_to_wire(holder))
}),
disposition: buffa::EnumValue::from(disposition_to_wire(value.disposition())),
attempts_granted: value.attempts_granted(),
changes_position: value.changes().position().get(),
changes_root: value.changes().root().map(|root| root.as_bytes().to_vec()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn status_from_wire(value: pb::ClaimStatus) -> Result<ClaimStatus, StateError> {
let changes = JournalHead::new(
JournalPosition::new(value.changes_position),
value
.changes_root
.map(|root| {
fixed_bytes::<{ CommitRoot::LEN }>("changes_root", &root)
.map(CommitRoot::from_bytes)
})
.transpose()?,
);
let holder_present = value.holder.is_set();
let work = WorkId::new(value.work);
let mut status = ClaimStatus::unclaimed(work, changes)
.at(
Revision::new(value.revision),
FencingToken::new(value.fence),
)
.with_attempts(value.attempts_granted);
status = match known("disposition", value.disposition)? {
pb::ClaimDisposition::Unclaimed => status,
pb::ClaimDisposition::Claimed => status.held_by(lease_from_wire(required(
"holder",
"a claimed status names its holder",
value.holder,
)?)),
pb::ClaimDisposition::Succeeded => status.settled(WorkOutcome::Succeeded),
pb::ClaimDisposition::Failed => status.settled(WorkOutcome::Failed),
pb::ClaimDisposition::Unspecified => {
return Err(malformed(
"disposition",
"a claim status names its disposition",
));
}
};
if status.disposition() != WorkDisposition::Claimed && holder_present {
return Err(malformed(
"holder",
"only a claimed status may carry a live holder",
));
}
Ok(status)
}
#[cfg(test)]
mod tests {
use super::*;
use polyc_state::{
claims::{ClaimWrite as _, adapter::ClaimCaseAdapter as _, cases},
context::CallContext,
};
#[test]
fn command_round_trip_preserves_the_presented_fence() {
let command = cases::renew(
"renew-1",
"work-1",
"attempt-1",
"owner-1",
FencingToken::new(42),
cases::LEASE,
);
let decoded = renew_from_wire(renew_to_wire(&command)).expect("decode renewal");
assert_eq!(decoded, command);
}
#[test]
fn unknown_attempt_outcome_is_refused() {
let error = attempt_outcome_from_wire(buffa::EnumValue::from(99))
.expect_err("unknown outcome must fail closed");
assert!(matches!(error, StateError::Malformed { .. }));
}
#[test]
fn granted_round_trip_preserves_state_issued_fence() {
let state = polyc_state::claims::memory::MemoryClaims::new();
let context = CallContext::new(
polyc_state::deadline::Deadline::after(state.now(), cases::BUDGET),
polyc_state::cancel::CancellationToken::new(),
);
let granted = state
.acquire(
cases::acquire("acquire-1", "work-1", "attempt-1", "owner-1", cases::LEASE),
&context,
)
.expect("grant");
let decoded = granted_from_wire(granted_to_wire(&granted)).expect("decode grant");
assert_eq!(decoded, granted);
}
}