use connectrpc::client::{ClientConfig, ClientTransport};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
error::StateError,
observation::{
CollectionId, CollectionKind, ObservationHead, ObservationOrdinal, ObservationRecord,
ObservationSource, ObservedCollectionListing, RecordObservation,
},
receipt::Receipt,
};
use crate::{
MAX_OBSERVATION_WIRE_MESSAGE_BYTES,
error::{TransportFallback, from_connect_error},
trace::bounded_traced_options,
wire::{DeclaredCall, Kernel},
};
const fn wire_kind(kind: CollectionKind) -> pb::ObservedCollectionKind {
match kind {
CollectionKind::Routines => pb::ObservedCollectionKind::OBSERVED_COLLECTION_KIND_ROUTINES,
}
}
pub struct ObservationClient<T> {
inner: pb::StateObservationServiceClient<T>,
}
impl<T> ObservationClient<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::StateObservationServiceClient::new(
transport,
config.with_default_max_message_size(MAX_OBSERVATION_WIRE_MESSAGE_BYTES),
),
}
}
fn fallback(declared: &DeclaredCall, attempted: usize) -> TransportFallback {
TransportFallback::new(
polyc_state::observation::family(),
MAX_OBSERVATION_WIRE_MESSAGE_BYTES as u64,
attempted as u64,
declared.budget,
)
}
pub async fn record(
&self,
declared: &DeclaredCall,
command: &RecordObservation,
) -> Result<Receipt, StateError> {
let draft = command.draft();
let request = pb::RecordObservationRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
collection_kind: wire_kind(command.collection().kind()).into(),
collection_namespace: command.collection().namespace().to_owned(),
resource_version: draft.resource_version().as_str().to_owned(),
observed_at_nanos: draft.observed_at().as_nanos(),
observer: draft.observer().as_str().to_owned(),
row_count: draft.row_count(),
payload: draft.payload().as_bytes().to_vec(),
expected_ordinal: match command.metadata().precondition() {
polyc_state::command::Precondition::JournalHead(position) => position.get(),
_ => 0,
},
__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(declared, attempted)))?
.into_owned();
Kernel::<Receipt>::try_from(reply.receipt.into_option().ok_or_else(|| {
StateError::Malformed {
field: "receipt".into(),
reason: "a successful observation commit returns a receipt".into(),
}
})?)
.map(Kernel::into_inner)
}
pub async fn latest(
&self,
declared: &DeclaredCall,
collection: &CollectionId,
) -> Result<Option<ObservationHead>, StateError> {
let request = pb::ReadLatestObservationRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
collection_kind: wire_kind(collection.kind()).into(),
collection_namespace: collection.namespace().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.read_latest_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
.into_owned();
reply.head.into_option().map(head_from_wire).transpose()
}
pub async fn read(
&self,
declared: &DeclaredCall,
source: &ObservationSource,
ordinal: ObservationOrdinal,
) -> Result<ObservationRecord, StateError> {
let request = pb::ReadObservationRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
source: buffa::MessageField::some(pb::ObservationSource::from(Kernel(source))),
ordinal: ordinal.as_journal_position().get(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.read_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
.into_owned();
record_from_wire(
reply
.record
.into_option()
.ok_or_else(|| StateError::Malformed {
field: "record".into(),
reason: "a successful read returns the observation".into(),
})?,
)
}
pub async fn collections(
&self,
declared: &DeclaredCall,
kind: CollectionKind,
) -> Result<ObservedCollectionListing, StateError> {
let request = pb::ListObservedCollectionsRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
collection_kind: wire_kind(kind).into(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.list_collections_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
.into_owned();
let collections = reply
.collection_namespaces
.into_iter()
.map(|namespace| CollectionId::try_new(kind, namespace))
.collect::<Result<Vec<_>, _>>()?;
Ok(ObservedCollectionListing::new(collections, reply.complete))
}
}
fn head_from_wire(wire: pb::ObservationHeadMessage) -> Result<ObservationHead, StateError> {
let source =
Kernel::<ObservationSource>::try_from(wire.source.into_option().ok_or_else(|| {
StateError::Malformed {
field: "source".into(),
reason: "an observation head names its collection and lineage".into(),
}
})?)?
.into_inner();
Ok(ObservationHead::from_parts(
source,
ObservationOrdinal::new(polyc_state::revision::JournalPosition::new(wire.ordinal)),
polyc_state::digest::ContentDigest::from_bytes(crate::wire::fixed_bytes::<
{ polyc_state::digest::ContentDigest::LEN },
>(
"payload_digest", &wire.payload_digest
)?),
polyc_state::observation::ResourceVersion::try_new(wire.resource_version)?,
polyc_state::deadline::MonotonicInstant::from_nanos(wire.observed_at_nanos),
polyc_state::deadline::MonotonicInstant::from_nanos(wire.recorded_at_nanos),
))
}
fn record_from_wire(wire: pb::ObservationRecordMessage) -> Result<ObservationRecord, StateError> {
let source =
Kernel::<ObservationSource>::try_from(wire.source.into_option().ok_or_else(|| {
StateError::Malformed {
field: "source".into(),
reason: "an observation record names its collection and lineage".into(),
}
})?)?
.into_inner();
let payload = polyc_state::observation::ObservationPayload::try_new(wire.payload)?;
Ok(ObservationRecord::from_parts(
source,
ObservationOrdinal::new(polyc_state::revision::JournalPosition::new(wire.ordinal)),
polyc_state::observation::ResourceVersion::try_new(wire.resource_version)?,
polyc_state::deadline::MonotonicInstant::from_nanos(wire.observed_at_nanos),
polyc_state::deadline::MonotonicInstant::from_nanos(wire.recorded_at_nanos),
polyc_state::observation::ObserverId::try_new(wire.observer)?,
wire.row_count,
polyc_state::digest::ContentDigest::from_bytes(crate::wire::fixed_bytes::<
{ polyc_state::digest::ContentDigest::LEN },
>(
"payload_digest", &wire.payload_digest
)?),
payload,
))
}