use connectrpc::client::{ClientConfig, ClientTransport};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
error::StateError,
id::CommandId,
immutable::{
ListObjectGenerations, ObjectError, ObjectHead, ObjectId, ObjectMetadata,
PruneObjectGeneration, ReadObjectGeneration, ReadObjectHead, RecordObjectGeneration,
UpdateObjectHead,
},
page::Page,
receipt::Receipt,
};
use crate::{
MAX_WIRE_MESSAGE_BYTES,
error::TransportFallback,
immutable::error::from_connect_error,
trace::bounded_traced_options,
wire::{DeclaredCall, Kernel},
};
pub struct ImmutableMetadataClient<T> {
inner: pb::StateImmutableMetadataServiceClient<T>,
}
impl<T> ImmutableMetadataClient<T>
where
T: ClientTransport,
<T::ResponseBody as connectrpc::http_body::Body>::Error: std::fmt::Display,
{
pub fn new(transport: T, config: ClientConfig) -> Self {
Self {
inner: pb::StateImmutableMetadataServiceClient::new(
transport,
config.with_default_max_message_size(MAX_WIRE_MESSAGE_BYTES),
),
}
}
fn fallback(attempted: usize) -> TransportFallback {
TransportFallback::new(
polyc_state::immutable::family(),
MAX_WIRE_MESSAGE_BYTES as u64,
attempted as u64,
)
}
pub async fn record(
&self,
declared: &DeclaredCall,
command: &RecordObjectGeneration,
) -> Result<Receipt, ObjectError> {
let request = pb::RecordObjectGenerationRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
metadata: buffa::MessageField::some(Kernel(command.metadata()).into()),
descriptor: buffa::MessageField::some(Kernel(command.descriptor()).into()),
caller: command.caller().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.record_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
receipt(reply.receipt)
}
pub async fn update_head(
&self,
declared: &DeclaredCall,
command: &UpdateObjectHead,
) -> Result<Receipt, ObjectError> {
let request = pb::UpdateObjectHeadRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
metadata: buffa::MessageField::some(Kernel(command.metadata()).into()),
object: command.object().as_str().to_owned(),
expected: command.expected().get(),
publish: command.publish().get(),
caller: command.caller().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.update_head_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
receipt(reply.receipt)
}
pub async fn prune(
&self,
declared: &DeclaredCall,
command: &PruneObjectGeneration,
) -> Result<Receipt, ObjectError> {
let request = pb::PruneObjectGenerationRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
metadata: buffa::MessageField::some(Kernel(command.metadata()).into()),
object: command.object().as_str().to_owned(),
generation: command.generation().get(),
caller: command.caller().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.prune_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
receipt(reply.receipt)
}
pub async fn head(
&self,
declared: &DeclaredCall,
request: &ReadObjectHead,
) -> Result<Option<ObjectHead>, ObjectError> {
let wire = pb::GetObjectHeadRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
object: request.object().as_str().to_owned(),
caller: request.caller().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&wire) as usize;
let reply = self
.inner
.get_head_with_options(wire, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
optional(reply.head)
}
pub async fn generation(
&self,
declared: &DeclaredCall,
request: &ReadObjectGeneration,
) -> Result<Option<ObjectMetadata>, ObjectError> {
let wire = pb::GetObjectGenerationRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
object: request.object().as_str().to_owned(),
generation: request.generation().get(),
caller: request.caller().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&wire) as usize;
let reply = self
.inner
.get_generation_with_options(wire, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
optional(reply.metadata)
}
pub async fn generations(
&self,
declared: &DeclaredCall,
request: &ListObjectGenerations,
) -> Result<Page<ObjectMetadata>, ObjectError> {
let wire = pb::ListObjectGenerationsRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
object: request.object().as_str().to_owned(),
caller: request.caller().as_str().to_owned(),
page: buffa::MessageField::some(Kernel(request.page()).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&wire) as usize;
let reply = self
.inner
.list_generations_with_options(wire, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
let page = reply
.page
.into_option()
.ok_or_else(|| malformed("page", "a successful generation listing carries its page"))?;
Kernel::<Page<ObjectMetadata>>::try_from(page)
.map(Kernel::into_inner)
.map_err(ObjectError::from)
}
pub async fn receipt(
&self,
declared: &DeclaredCall,
object: &ObjectId,
command_id: &CommandId,
) -> Result<Option<Receipt>, ObjectError> {
let wire = pb::GetObjectReceiptRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
object: object.as_str().to_owned(),
command_id: command_id.as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&wire) as usize;
let reply = self
.inner
.get_receipt_with_options(wire, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
optional(reply.receipt)
}
}
fn malformed(field: &str, reason: &str) -> ObjectError {
ObjectError::from(StateError::Malformed {
field: field.to_owned(),
reason: reason.to_owned(),
})
}
fn optional<T, W, P>(field: buffa::MessageField<W, P>) -> Result<Option<T>, ObjectError>
where
W: Default,
P: buffa::ProtoBox<W>,
Kernel<T>: TryFrom<W, Error = StateError>,
{
field
.into_option()
.map(|value| {
Kernel::<T>::try_from(value)
.map(Kernel::into_inner)
.map_err(ObjectError::from)
})
.transpose()
}
fn receipt(field: impl Into<Option<pb::Receipt>>) -> Result<Receipt, ObjectError> {
let value = field.into().ok_or_else(|| {
malformed(
"receipt",
"a successful object mutation carries its receipt",
)
})?;
Kernel::<Receipt>::try_from(value)
.map(Kernel::into_inner)
.map_err(ObjectError::from)
}