use connectrpc::client::{ClientConfig, ClientTransport};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
error::StateError,
id::{CommandId, NamespaceId},
receipt::Receipt,
revision::Revision,
usage::{
AccrualView, ConversationId, ConversationUsageRecord, PersonaId, UsageCommand, UsageFact,
UsageRollupRecord,
},
};
use crate::{
MAX_WIRE_MESSAGE_BYTES,
error::{TransportFallback, from_connect_error},
trace::bounded_traced_options,
wire::{DeclaredCall, Kernel},
};
use super::wire::{conversation_from_wire, metadata_to_wire, operation_to_wire, rollup_from_wire};
pub struct UsageClient<T> {
inner: pb::StateUsageServiceClient<T>,
namespace: NamespaceId,
}
impl<T> UsageClient<T>
where
T: ClientTransport,
<T::ResponseBody as connectrpc::http_body::Body>::Error: std::fmt::Display,
{
#[must_use]
pub fn new(transport: T, config: ClientConfig, namespace: NamespaceId) -> Self {
Self {
inner: pb::StateUsageServiceClient::new(
transport,
config.with_default_max_message_size(MAX_WIRE_MESSAGE_BYTES),
),
namespace,
}
}
fn fallback(declared: &DeclaredCall, attempted: usize) -> TransportFallback {
TransportFallback::new(
polyc_state::usage::family(),
MAX_WIRE_MESSAGE_BYTES as u64,
attempted as u64,
declared.budget,
)
}
pub async fn transact(
&self,
declared: &DeclaredCall,
command: &UsageCommand,
) -> Result<Receipt, StateError> {
if command.metadata().scope().namespace() != &self.namespace {
return Err(StateError::Denied {
family: polyc_state::usage::family(),
});
}
let request = pb::TransactStateUsageRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
metadata: buffa::MessageField::some(metadata_to_wire(command)),
persona_id: command.persona().as_str().to_owned(),
operation: buffa::MessageField::some(operation_to_wire(command.operation())),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.transact_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
.into_owned();
Kernel::<Receipt>::try_from(reply.receipt.into_option().ok_or_else(|| {
StateError::Malformed {
field: "receipt".into(),
reason: "a successful usage mutation returns a receipt".into(),
}
})?)
.map(Kernel::into_inner)
}
pub async fn rollup(
&self,
declared: &DeclaredCall,
persona: &PersonaId,
) -> Result<UsageFact<Option<UsageRollupRecord>>, StateError> {
let request = pb::GetStateUsageRollupRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
namespace: self.namespace.as_str().to_owned(),
persona_id: persona.as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.get_rollup_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
.into_owned();
Ok(UsageFact::observed(
reply.rollup.into_option().as_ref().map(rollup_from_wire),
Revision::new(reply.snapshot_revision),
reply.entry_revision.map(Revision::new),
))
}
pub async fn conversation(
&self,
declared: &DeclaredCall,
persona: &PersonaId,
conversation: &ConversationId,
) -> Result<UsageFact<ConversationUsageRecord>, StateError> {
let request = pb::GetStateUsageConversationRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
namespace: self.namespace.as_str().to_owned(),
persona_id: persona.as_str().to_owned(),
conversation_id: conversation.as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.get_conversation_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
.into_owned();
Ok(UsageFact::observed(
conversation_from_wire(reply.record.into_option().ok_or_else(|| {
StateError::Malformed {
field: "record".into(),
reason: "a usage reply carries its accrual row".into(),
}
})?),
Revision::new(reply.snapshot_revision),
reply.entry_revision.map(Revision::new),
))
}
pub async fn accrual_view(
&self,
declared: &DeclaredCall,
persona: &PersonaId,
conversation: &ConversationId,
) -> Result<AccrualView, StateError> {
let request = pb::GetStateUsageAccrualViewRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
namespace: self.namespace.as_str().to_owned(),
persona_id: persona.as_str().to_owned(),
conversation_id: conversation.as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.get_accrual_view_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
.into_owned();
let rollup = reply.rollup.into_option();
Ok(AccrualView::observed(
rollup.as_ref().map(rollup_from_wire),
conversation_from_wire(reply.record.into_option().ok_or_else(|| {
StateError::Malformed {
field: "record".into(),
reason: "an accrual view carries its accrual row".into(),
}
})?),
Revision::new(reply.snapshot_revision),
reply.rollup_revision.map(Revision::new),
reply.record_revision.map(Revision::new),
))
}
pub async fn committed_receipt(
&self,
declared: &DeclaredCall,
command_id: &CommandId,
) -> Result<Option<Receipt>, StateError> {
let request = pb::GetStateUsageReceiptRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
namespace: self.namespace.as_str().to_owned(),
command_id: command_id.as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.get_receipt_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
.into_owned();
reply
.receipt
.into_option()
.map(|value| Kernel::<Receipt>::try_from(value).map(Kernel::into_inner))
.transpose()
}
}