use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
command::{CommandEnvelope, CommandMetadata, ResourceBounds},
digest::ContentDigest,
error::StateError,
id::{Audience, CommandId, NamespaceId, Purpose},
revision::Revision,
usage::{
ConversationId, ConversationUsageRecord, PersonaId, ReplacedConversation, TurnId,
TurnUsage, UsageCommand, UsageOperation, UsageRollupRecord, usage_scope,
},
versioned::{EntryExpectation, MAX_MUTATIONS_PER_TRANSACTION, MAX_TRANSACTION_PAYLOAD_BYTES},
};
use crate::wire::{fixed_bytes, malformed, required};
pub(crate) fn expected_to_wire(value: EntryExpectation) -> pb::StateUsageExpectedEntry {
use pb::__buffa::oneof::state_usage_expected_entry::Expected;
let expected = match value {
EntryExpectation::Absent => Expected::from(pb::StateUsageExpectedAbsent {
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
EntryExpectation::Revision(revision) => Expected::from(pb::StateUsageExpectedRevision {
revision: revision.get(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
};
pb::StateUsageExpectedEntry {
expected: Some(expected),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn expected_from_wire(value: pb::StateUsageExpectedEntry) -> Result<EntryExpectation, StateError> {
use pb::__buffa::oneof::state_usage_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 usage operation declares its exact row premise",
)),
}
}
pub(crate) fn rollup_to_wire(value: &UsageRollupRecord) -> pb::StateUsageRollup {
pb::StateUsageRollup {
committed_turns: value.committed_turns(),
input_tokens: value.input_tokens(),
output_tokens: value.output_tokens(),
last_active_ms: value.last_active_ms(),
updated_at_ms: value.updated_at_ms(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) const fn rollup_from_wire(value: &pb::StateUsageRollup) -> UsageRollupRecord {
UsageRollupRecord::from_parts(
value.committed_turns,
value.input_tokens,
value.output_tokens,
value.last_active_ms,
value.updated_at_ms,
)
}
pub(crate) fn conversation_to_wire(value: &ConversationUsageRecord) -> pb::StateUsageConversation {
pb::StateUsageConversation {
committed_turns: value.committed_turns(),
input_tokens: value.input_tokens(),
output_tokens: value.output_tokens(),
last_active_ms: value.last_active_ms(),
last_turn_id: value.last_turn_id().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
pub(crate) fn conversation_from_wire(value: pb::StateUsageConversation) -> ConversationUsageRecord {
ConversationUsageRecord::from_parts(
value.committed_turns,
value.input_tokens,
value.output_tokens,
value.last_active_ms,
TurnId::new(value.last_turn_id),
)
}
fn turn_to_wire(value: TurnUsage) -> pb::StateUsageTurn {
pb::StateUsageTurn {
input_tokens: value.input_tokens(),
output_tokens: value.output_tokens(),
at_ms: value.at_ms(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
const fn turn_from_wire(value: &pb::StateUsageTurn) -> TurnUsage {
TurnUsage::new(value.input_tokens, value.output_tokens, value.at_ms)
}
fn replaced_to_wire(value: &ReplacedConversation) -> pb::StateUsageReplacedConversation {
pb::StateUsageReplacedConversation {
conversation_id: value.conversation().as_str().to_owned(),
record: value
.record()
.map_or_else(buffa::MessageField::none, |record| {
buffa::MessageField::some(conversation_to_wire(record))
}),
expected: buffa::MessageField::some(expected_to_wire(value.expected())),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn replaced_from_wire(
value: pb::StateUsageReplacedConversation,
) -> Result<ReplacedConversation, StateError> {
Ok(ReplacedConversation::new(
ConversationId::new(value.conversation_id),
value.record.into_option().map(conversation_from_wire),
expected_from_wire(required(
"expected",
"a re-derived usage row carries its premise",
value.expected,
)?)?,
))
}
pub(crate) fn operation_to_wire(value: &UsageOperation) -> pb::StateUsageOperation {
use pb::__buffa::oneof::state_usage_operation::Operation;
let operation = match value {
UsageOperation::Accrue {
now_ms,
conversation,
turn,
usage,
rollup,
record,
rollup_expected,
record_expected,
} => Operation::from(pb::StateUsageAccrueOperation {
now_ms: *now_ms,
conversation_id: conversation.as_str().to_owned(),
turn_id: turn.as_str().to_owned(),
usage: buffa::MessageField::some(turn_to_wire(*usage)),
rollup: buffa::MessageField::some(rollup_to_wire(rollup)),
record: buffa::MessageField::some(conversation_to_wire(record)),
rollup_expected: buffa::MessageField::some(expected_to_wire(*rollup_expected)),
record_expected: buffa::MessageField::some(expected_to_wire(*record_expected)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
UsageOperation::Fold {
now_ms,
absorbed,
rollup,
rollup_expected,
absorbed_expected,
} => Operation::from(pb::StateUsageFoldOperation {
now_ms: *now_ms,
absorbed_persona_id: absorbed.as_str().to_owned(),
rollup: buffa::MessageField::some(rollup_to_wire(rollup)),
rollup_expected: buffa::MessageField::some(expected_to_wire(*rollup_expected)),
absorbed_expected: buffa::MessageField::some(expected_to_wire(*absorbed_expected)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
UsageOperation::Replace {
now_ms,
rollup,
rollup_expected,
conversations,
} => Operation::from(pb::StateUsageReplaceOperation {
now_ms: *now_ms,
rollup: buffa::MessageField::some(rollup_to_wire(rollup)),
rollup_expected: buffa::MessageField::some(expected_to_wire(*rollup_expected)),
conversations: conversations.iter().map(replaced_to_wire).collect(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
};
pb::StateUsageOperation {
operation: Some(operation),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn operation_from_wire(value: pb::StateUsageOperation) -> Result<UsageOperation, StateError> {
use pb::__buffa::oneof::state_usage_operation::Operation;
match value.operation {
Some(Operation::Accrue(value)) => Ok(UsageOperation::Accrue {
now_ms: value.now_ms,
conversation: ConversationId::new(value.conversation_id),
turn: TurnId::new(value.turn_id),
usage: turn_from_wire(&required(
"usage",
"an accrual carries what its turn contributed",
value.usage,
)?),
rollup: rollup_from_wire(&required(
"rollup",
"an accrual carries its rollup result",
value.rollup,
)?),
record: conversation_from_wire(required(
"record",
"an accrual carries its conversation result",
value.record,
)?),
rollup_expected: expected_from_wire(required(
"rollup_expected",
"an accrual carries its rollup premise",
value.rollup_expected,
)?)?,
record_expected: expected_from_wire(required(
"record_expected",
"an accrual carries its conversation premise",
value.record_expected,
)?)?,
}),
Some(Operation::Fold(value)) => Ok(UsageOperation::Fold {
now_ms: value.now_ms,
absorbed: PersonaId::new(value.absorbed_persona_id),
rollup: rollup_from_wire(&required(
"rollup",
"a fold carries its survivor result",
value.rollup,
)?),
rollup_expected: expected_from_wire(required(
"rollup_expected",
"a fold carries its survivor premise",
value.rollup_expected,
)?)?,
absorbed_expected: expected_from_wire(required(
"absorbed_expected",
"a fold carries its absorbed premise",
value.absorbed_expected,
)?)?,
}),
Some(Operation::Replace(value)) => Ok(UsageOperation::Replace {
now_ms: value.now_ms,
rollup: rollup_from_wire(&required(
"rollup",
"a re-derivation carries its rollup result",
value.rollup,
)?),
rollup_expected: expected_from_wire(required(
"rollup_expected",
"a re-derivation carries its rollup premise",
value.rollup_expected,
)?)?,
conversations: value
.conversations
.into_iter()
.map(replaced_from_wire)
.collect::<Result<Vec<_>, _>>()?,
}),
None => Err(malformed(
"operation",
"a usage command names one operation",
)),
}
}
pub(crate) fn metadata_to_wire(command: &UsageCommand) -> pb::StateUsageCommandMetadata {
let value = command.metadata();
pb::StateUsageCommandMetadata {
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::StateUsageCommandMetadata,
persona_id: String,
operation: pb::StateUsageOperation,
) -> Result<UsageCommand, StateError> {
let namespace = NamespaceId::new(metadata.namespace);
let digest = ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
"digest",
&metadata.digest,
)?);
Ok(UsageCommand::new(
CommandMetadata::new(
CommandId::new(metadata.command_id),
polyc_state::usage::family(),
digest,
usage_scope(&namespace),
CommandEnvelope::new(
Purpose::new(metadata.purpose),
Audience::new(metadata.command_audience),
ResourceBounds::new(MAX_TRANSACTION_PAYLOAD_BYTES, MAX_MUTATIONS_PER_TRANSACTION),
),
),
PersonaId::new(persona_id),
operation_from_wire(operation)?,
))
}
#[cfg(test)]
mod tests {
use super::*;
const NOW: u64 = 1_700_000_000_000;
const TURN: &str = "0190a1b2-c3d4-7e5f-8a9b-000000000001";
fn rollup() -> UsageRollupRecord {
UsageRollupRecord::from_parts(1, 2, 3, NOW, NOW)
}
fn record() -> ConversationUsageRecord {
ConversationUsageRecord::from_parts(1, 2, 3, NOW, TurnId::new(TURN))
}
#[test]
fn every_operation_round_trips_through_the_wire() {
let operations = [
UsageOperation::Accrue {
now_ms: NOW,
conversation: ConversationId::new("conv-1"),
turn: TurnId::new(TURN),
usage: TurnUsage::new(2, 3, NOW),
rollup: rollup(),
record: record(),
rollup_expected: EntryExpectation::Absent,
record_expected: EntryExpectation::Revision(Revision::new(9)),
},
UsageOperation::Fold {
now_ms: NOW + 1,
absorbed: PersonaId::new("p-2"),
rollup: rollup(),
rollup_expected: EntryExpectation::Revision(Revision::new(11)),
absorbed_expected: EntryExpectation::Revision(Revision::new(12)),
},
UsageOperation::Replace {
now_ms: NOW + 2,
rollup: rollup(),
rollup_expected: EntryExpectation::Absent,
conversations: vec![
ReplacedConversation::new(
ConversationId::new("conv-1"),
Some(record()),
EntryExpectation::Revision(Revision::new(13)),
),
ReplacedConversation::new(
ConversationId::new("conv-2"),
None,
EntryExpectation::Revision(Revision::new(14)),
),
],
},
];
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 a_removed_row_stays_removed() {
let removed = ReplacedConversation::new(
ConversationId::new("conv-1"),
None,
EntryExpectation::Absent,
);
let restored = replaced_from_wire(replaced_to_wire(&removed)).expect("round trip");
assert!(
restored.record().is_none(),
"a removal came back as a write"
);
}
#[test]
fn the_rows_round_trip_through_the_wire() {
assert_eq!(rollup_from_wire(&rollup_to_wire(&rollup())), rollup());
assert_eq!(
conversation_from_wire(conversation_to_wire(&record())),
record()
);
}
#[test]
fn a_missing_oneof_and_a_missing_premise_fail_closed() {
assert!(
operation_from_wire(pb::StateUsageOperation::default()).is_err(),
"an operation with no variant was accepted"
);
assert!(
expected_from_wire(pb::StateUsageExpectedEntry::default()).is_err(),
"a premise with no variant was accepted"
);
let fold = pb::StateUsageFoldOperation {
now_ms: NOW,
absorbed_persona_id: "p-2".to_owned(),
rollup: buffa::MessageField::some(rollup_to_wire(&rollup())),
rollup_expected: buffa::MessageField::none(),
absorbed_expected: buffa::MessageField::none(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
assert!(
operation_from_wire(pb::StateUsageOperation {
operation: Some(pb::__buffa::oneof::state_usage_operation::Operation::from(
fold
)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
.is_err(),
"a fold with no premise was accepted"
);
}
}