use connectrpc::client::{ClientConfig, ClientTransport};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
claims::{
AcquireClaim, AttemptHistory, ClaimStatus, CompleteClaim, GetAttemptHistory,
GetClaimStatus, Granted, ReleaseClaim, RenewClaim,
},
command::CommandScope,
error::StateError,
id::CommandId,
receipt::Receipt,
};
use crate::{
MAX_CLAIMS_WIRE_MESSAGE_BYTES,
error::{TransportFallback, from_connect_error},
trace::bounded_traced_options,
wire::{DeclaredCall, Kernel},
};
use super::wire::{
acquire_to_wire, complete_to_wire, granted_from_wire, history_from_wire, release_to_wire,
renew_to_wire, scope_to_wire, status_from_wire,
};
pub struct ClaimsClient<T> {
inner: pb::StateClaimsServiceClient<T>,
}
impl<T> ClaimsClient<T>
where
T: ClientTransport,
<T::ResponseBody as connectrpc::http_body::Body>::Error: std::fmt::Display,
{
#[must_use]
pub fn new(transport: T, config: ClientConfig) -> Self {
Self {
inner: pb::StateClaimsServiceClient::new(
transport,
config.with_default_max_message_size(MAX_CLAIMS_WIRE_MESSAGE_BYTES),
),
}
}
fn fallback(attempted: usize) -> TransportFallback {
TransportFallback::new(
polyc_state::claims::family(),
MAX_CLAIMS_WIRE_MESSAGE_BYTES as u64,
attempted as u64,
)
}
pub async fn acquire(
&self,
declared: &DeclaredCall,
command: &AcquireClaim,
) -> Result<Granted, StateError> {
let request = pb::AcquireClaimRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
command: buffa::MessageField::some(acquire_to_wire(command)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.acquire_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
granted_from_wire(required(
reply.granted,
"a successful acquire returns its grant",
)?)
}
pub async fn renew(
&self,
declared: &DeclaredCall,
command: &RenewClaim,
) -> Result<Granted, StateError> {
let request = pb::RenewClaimRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
command: buffa::MessageField::some(renew_to_wire(command)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.renew_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
granted_from_wire(required(
reply.granted,
"a successful renewal returns its grant",
)?)
}
pub async fn complete(
&self,
declared: &DeclaredCall,
command: &CompleteClaim,
) -> Result<Receipt, StateError> {
let request = pb::CompleteClaimRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
command: buffa::MessageField::some(complete_to_wire(command)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.complete_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
receipt(reply.receipt, "a successful completion returns its receipt")
}
pub async fn release(
&self,
declared: &DeclaredCall,
command: &ReleaseClaim,
) -> Result<Receipt, StateError> {
let request = pb::ReleaseClaimRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
command: buffa::MessageField::some(release_to_wire(command)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.release_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
receipt(reply.receipt, "a successful release returns its receipt")
}
pub async fn status(
&self,
declared: &DeclaredCall,
request: &GetClaimStatus,
) -> Result<ClaimStatus, StateError> {
let request = pb::GetClaimStatusRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
scope: buffa::MessageField::some(scope_to_wire(request.scope())),
work: request.work().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.get_status_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
status_from_wire(required(
reply.status,
"a successful read returns claim status",
)?)
}
pub async fn history(
&self,
declared: &DeclaredCall,
request: &GetAttemptHistory,
) -> Result<AttemptHistory, StateError> {
let request = pb::GetClaimHistoryRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
scope: buffa::MessageField::some(scope_to_wire(request.scope())),
work: request.work().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.get_history_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
history_from_wire(required(
reply.history,
"a successful read returns attempt history",
)?)
}
pub async fn committed_receipt(
&self,
declared: &DeclaredCall,
scope: &CommandScope,
command_id: &CommandId,
) -> Result<Option<Receipt>, StateError> {
let request = pb::GetClaimReceiptRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
scope: buffa::MessageField::some(scope_to_wire(scope)),
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(attempted)))?
.into_owned();
reply
.receipt
.into_option()
.map(|receipt| Kernel::<Receipt>::try_from(receipt).map(Kernel::into_inner))
.transpose()
}
}
fn required<T: Default, P: buffa::ProtoBox<T>>(
field: buffa::MessageField<T, P>,
reason: &str,
) -> Result<T, StateError> {
field.into_option().ok_or_else(|| StateError::Malformed {
field: "response".to_owned(),
reason: reason.to_owned(),
})
}
fn receipt(field: impl Into<Option<pb::Receipt>>, reason: &str) -> Result<Receipt, StateError> {
let field = field.into().ok_or_else(|| StateError::Malformed {
field: "response".to_owned(),
reason: reason.to_owned(),
})?;
Kernel::<Receipt>::try_from(field).map(Kernel::into_inner)
}