use buffa::EnumValue;
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
command::{CommandEnvelope, CommandMetadata, FencingToken, ResourceBounds},
digest::ContentDigest,
error::StateError,
id::{Audience, CommandId, NamespaceId, Purpose},
revision::Revision,
tasks::{
ContextId, ContextIndex, EdgeId, TaskCommand, TaskId, TaskOperation, TaskRecord, TaskState,
task_scope,
},
versioned::{EntryExpectation, MAX_MUTATIONS_PER_TRANSACTION, MAX_TRANSACTION_PAYLOAD_BYTES},
};
use crate::wire::{fixed_bytes, malformed, required};
fn expected_to_wire(value: EntryExpectation) -> pb::StateTaskExpectedEntry {
use pb::__buffa::oneof::state_task_expected_entry::Expected;
let expected = match value {
EntryExpectation::Absent => Expected::from(pb::StateTaskExpectedAbsent {
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
EntryExpectation::Revision(revision) => Expected::from(pb::StateTaskExpectedRevision {
revision: revision.get(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
};
pb::StateTaskExpectedEntry {
expected: Some(expected),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn expected_from_wire(value: pb::StateTaskExpectedEntry) -> Result<EntryExpectation, StateError> {
use pb::__buffa::oneof::state_task_expected_entry::Expected;
match value.expected {
Some(Expected::Absent(_)) => Ok(EntryExpectation::Absent),
Some(Expected::Revision(value)) => {
Ok(EntryExpectation::Revision(Revision::new(value.revision)))
}
None => Err(malformed(
"expected",
"a task operation declares its exact row premise",
)),
}
}
pub(crate) const fn state_to_wire(value: TaskState) -> pb::StateTaskState {
match value {
TaskState::Submitted => pb::StateTaskState::Submitted,
TaskState::Working => pb::StateTaskState::Working,
TaskState::InputRequired => pb::StateTaskState::InputRequired,
TaskState::AuthRequired => pb::StateTaskState::AuthRequired,
TaskState::Completed => pb::StateTaskState::Completed,
TaskState::Failed => pb::StateTaskState::Failed,
TaskState::Canceled => pb::StateTaskState::Canceled,
TaskState::Rejected => pb::StateTaskState::Rejected,
}
}
fn state_from_wire(value: EnumValue<pb::StateTaskState>) -> Result<TaskState, StateError> {
match value {
EnumValue::Known(pb::StateTaskState::Submitted) => Ok(TaskState::Submitted),
EnumValue::Known(pb::StateTaskState::Working) => Ok(TaskState::Working),
EnumValue::Known(pb::StateTaskState::InputRequired) => Ok(TaskState::InputRequired),
EnumValue::Known(pb::StateTaskState::AuthRequired) => Ok(TaskState::AuthRequired),
EnumValue::Known(pb::StateTaskState::Completed) => Ok(TaskState::Completed),
EnumValue::Known(pb::StateTaskState::Failed) => Ok(TaskState::Failed),
EnumValue::Known(pb::StateTaskState::Canceled) => Ok(TaskState::Canceled),
EnumValue::Known(pb::StateTaskState::Rejected) => Ok(TaskState::Rejected),
EnumValue::Known(pb::StateTaskState::Unspecified) | EnumValue::Unknown(_) => {
Err(malformed("state", "a task names one known lifecycle state"))
}
}
}
pub(crate) fn record_to_wire(value: &TaskRecord) -> pb::StateTaskRecord {
pb::StateTaskRecord {
task_id: value.id().as_str().to_owned(),
context_id: value.context_id().as_str().to_owned(),
owner_edge_id: value.owner().as_str().to_owned(),
state: EnumValue::Known(state_to_wire(value.state())),
status_detail: value.status_detail().to_vec(),
artifacts: value.artifacts().to_vec(),
history: value.history().to_vec(),
metadata: value.metadata().to_vec(),
created_at_ms: value.created_at_ms(),
updated_at_ms: value.updated_at_ms(),
ownership_fence: value.ownership_fence().map(FencingToken::get),
ownership_binding: value.ownership_binding().unwrap_or_default().to_vec(),
last_owned_operation: value.last_owned_operation().unwrap_or_default().to_vec(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn record_from_wire(value: pb::StateTaskRecord) -> Result<TaskRecord, StateError> {
let state = state_from_wire(value.state)?;
let ownership_fence = value.ownership_fence.map(FencingToken::new);
let ownership_binding =
(!value.ownership_binding.is_empty()).then_some(value.ownership_binding);
let last_owned_operation =
(!value.last_owned_operation.is_empty()).then_some(value.last_owned_operation);
if ownership_fence == Some(FencingToken::UNCLAIMED) {
return Err(malformed(
"ownership_fence",
"a claimed task carries a State-issued nonzero ownership fence",
));
}
if value.owner_edge_id.is_empty() {
return Err(malformed(
"owner_edge_id",
"a task record names the edge that created it",
));
}
Ok(TaskRecord::from_parts(
TaskId::new(value.task_id),
ContextId::new(value.context_id),
EdgeId::new(value.owner_edge_id),
state,
value.status_detail,
value.artifacts,
value.history,
value.metadata,
ownership_fence,
ownership_binding,
last_owned_operation,
value.created_at_ms,
value.updated_at_ms,
))
}
pub(crate) fn index_to_wire(value: &ContextIndex) -> pb::StateTaskContextIndex {
pb::StateTaskContextIndex {
task_ids: value
.entries()
.iter()
.map(|id| id.as_str().to_owned())
.collect(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn index_from_wire(value: pb::StateTaskContextIndex) -> ContextIndex {
ContextIndex::new(value.task_ids.into_iter().map(TaskId::new).collect())
}
pub(crate) fn operation_to_wire(value: &TaskOperation) -> pb::StateTaskOperation {
use pb::__buffa::oneof::state_task_operation::Operation;
let operation = match value {
TaskOperation::Create {
task,
task_expected,
index,
index_expected,
} => Operation::from(pb::StateTaskCreate {
task: buffa::MessageField::some(record_to_wire(task)),
task_expected: buffa::MessageField::some(expected_to_wire(*task_expected)),
index: buffa::MessageField::some(index_to_wire(index)),
index_expected: buffa::MessageField::some(expected_to_wire(*index_expected)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
TaskOperation::Transition {
now_ms,
task,
task_expected,
ownership_fence,
} => Operation::from(pb::StateTaskTransition {
now_ms: *now_ms,
task: buffa::MessageField::some(record_to_wire(task)),
task_expected: buffa::MessageField::some(expected_to_wire(*task_expected)),
ownership_fence: ownership_fence.map(FencingToken::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
TaskOperation::Cancel {
now_ms,
task,
task_expected,
} => Operation::from(pb::StateTaskCancel {
now_ms: *now_ms,
task: buffa::MessageField::some(record_to_wire(task)),
task_expected: buffa::MessageField::some(expected_to_wire(*task_expected)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
};
pb::StateTaskOperation {
operation: Some(operation),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn operation_from_wire(value: pb::StateTaskOperation) -> Result<TaskOperation, StateError> {
use pb::__buffa::oneof::state_task_operation::Operation;
let task = |value| {
required("task", "a task operation carries its record", value).and_then(record_from_wire)
};
let expected = |field, value| {
expected_from_wire(required(
field,
"a task operation carries its premise",
value,
)?)
};
match value.operation {
Some(Operation::Create(value)) => Ok(TaskOperation::Create {
task: task(value.task)?,
task_expected: expected("task_expected", value.task_expected)?,
index: index_from_wire(required(
"index",
"a task creation carries its context index",
value.index,
)?),
index_expected: expected("index_expected", value.index_expected)?,
}),
Some(Operation::Transition(value)) => Ok(TaskOperation::Transition {
now_ms: value.now_ms,
task: task(value.task)?,
task_expected: expected("task_expected", value.task_expected)?,
ownership_fence: value.ownership_fence.map(FencingToken::new),
}),
Some(Operation::Cancel(value)) => Ok(TaskOperation::Cancel {
now_ms: value.now_ms,
task: task(value.task)?,
task_expected: expected("task_expected", value.task_expected)?,
}),
None => Err(malformed("operation", "a task command names one operation")),
}
}
pub(crate) fn metadata_to_wire(command: &TaskCommand) -> pb::StateTaskCommandMetadata {
let value = command.metadata();
pb::StateTaskCommandMetadata {
command_id: value.command_id().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(),
digest: value.digest().as_bytes().to_vec(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn command_from_wire(
metadata: pb::StateTaskCommandMetadata,
operation: pb::StateTaskOperation,
) -> Result<TaskCommand, StateError> {
let namespace = NamespaceId::new(metadata.namespace);
let digest = ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
"digest",
&metadata.digest,
)?);
Ok(TaskCommand::new(
CommandMetadata::new(
CommandId::new(metadata.command_id),
polyc_state::tasks::family(),
digest,
task_scope(&namespace),
CommandEnvelope::new(
Purpose::new(metadata.purpose),
Audience::new(metadata.command_audience),
ResourceBounds::new(MAX_TRANSACTION_PAYLOAD_BYTES, MAX_MUTATIONS_PER_TRANSACTION),
),
),
operation_from_wire(operation)?,
))
}
#[cfg(test)]
mod tests {
use super::*;
fn record() -> TaskRecord {
TaskRecord::submitted(
TaskId::new("task-1"),
ContextId::new("context-1"),
EdgeId::new("edge-1"),
vec![b"open".to_vec()],
10,
)
.transitioned_under(
TaskState::Working,
b"running".to_vec(),
vec![b"artifact".to_vec()],
b"{}".to_vec(),
&[b"reply".to_vec()],
11,
FencingToken::new(7),
vec![7; 32],
vec![8; 32],
)
}
#[test]
fn every_operation_round_trips_through_the_wire() {
let operations = [
TaskOperation::Create {
task: TaskRecord::submitted(
TaskId::new("task-1"),
ContextId::new("context-1"),
EdgeId::new("edge-1"),
vec![b"open".to_vec()],
10,
),
task_expected: EntryExpectation::Absent,
index: ContextIndex::default().with(TaskId::new("task-1")),
index_expected: EntryExpectation::Absent,
},
TaskOperation::Transition {
now_ms: 11,
task: record(),
task_expected: EntryExpectation::Revision(Revision::new(7)),
ownership_fence: Some(FencingToken::new(7)),
},
TaskOperation::Cancel {
now_ms: 12,
task: record(),
task_expected: EntryExpectation::Revision(Revision::new(9)),
},
];
for operation in operations {
let restored = operation_from_wire(operation_to_wire(&operation))
.expect("an operation survives its own encoding");
assert_eq!(restored, operation);
}
}
#[test]
fn the_record_and_index_round_trip_through_the_wire() {
assert_eq!(
record_from_wire(record_to_wire(&record())).expect("record round trip"),
record()
);
let index = ContextIndex::default()
.with(TaskId::new("a"))
.with(TaskId::new("b"));
assert_eq!(index_from_wire(index_to_wire(&index)), index);
}
#[test]
fn a_record_without_an_owner_is_refused_on_the_wire() {
let mut ownerless = record_to_wire(&record());
ownerless.owner_edge_id = String::new();
assert!(
record_from_wire(ownerless).is_err(),
"a task record with no owner decoded"
);
let mut absent = record_to_wire(&record());
absent.owner_edge_id = pb::StateTaskRecord::default().owner_edge_id;
assert!(
record_from_wire(absent).is_err(),
"an omitted owner field decoded"
);
}
#[test]
fn a_missing_oneof_and_an_unknown_enum_fail_closed() {
assert!(
operation_from_wire(pb::StateTaskOperation::default()).is_err(),
"an operation with no variant was accepted"
);
assert!(
expected_from_wire(pb::StateTaskExpectedEntry::default()).is_err(),
"a premise with no variant was accepted"
);
let mut unknown = record_to_wire(&record());
unknown.state = EnumValue::Unknown(999);
assert!(
record_from_wire(unknown).is_err(),
"an unknown lifecycle state was accepted"
);
let mut unspecified = record_to_wire(&record());
unspecified.state = EnumValue::Known(pb::StateTaskState::Unspecified);
assert!(
record_from_wire(unspecified).is_err(),
"the unspecified lifecycle state was accepted"
);
let mut unclaimed = record_to_wire(&record());
unclaimed.ownership_fence = Some(0);
assert!(
record_from_wire(unclaimed).is_err(),
"a zero ownership fence was accepted as a State grant"
);
let create = pb::StateTaskCreate {
task: buffa::MessageField::some(record_to_wire(&record())),
task_expected: buffa::MessageField::none(),
index: buffa::MessageField::none(),
index_expected: buffa::MessageField::none(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
assert!(
operation_from_wire(pb::StateTaskOperation {
operation: Some(pb::__buffa::oneof::state_task_operation::Operation::from(
create
)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
.is_err(),
"a creation with no premise was accepted"
);
}
}