use connectrpc::client::{ClientConfig, ClientTransport};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
error::StateError,
persona_memory::journal::{
AppendMemoryRecords, DestroyMemoryRecords, MemoryAppendOutcome, MemoryJournalPartition,
MemoryMigrationOutcome, MemoryPartitionPage, MemoryReplayPage, MemoryReplayRange,
MemoryRewriteOutcome, MigrateMemoryRecords, RewriteMemoryRecords,
},
receipt::Receipt,
};
use crate::{
MAX_PERSONA_MEMORY_JOURNAL_WIRE_MESSAGE_BYTES,
error::{TransportFallback, from_connect_error},
trace::bounded_traced_options,
wire::{DeclaredCall, Kernel},
};
use super::wire::{
append_outcome_from_wire, append_to_wire, destroy_to_wire, migrate_to_wire,
migration_outcome_from_wire, partition_page_from_wire, replay_from_wire,
rewrite_outcome_from_wire, rewrite_to_wire,
};
pub struct PersonaMemoryJournalClient<T> {
inner: pb::StatePersonaMemoryJournalServiceClient<T>,
}
impl<T> PersonaMemoryJournalClient<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::StatePersonaMemoryJournalServiceClient::new(
transport,
config.with_default_max_message_size(MAX_PERSONA_MEMORY_JOURNAL_WIRE_MESSAGE_BYTES),
),
}
}
fn fallback(attempted: usize) -> TransportFallback {
TransportFallback::new(
polyc_state::persona_memory::journal::family(),
MAX_PERSONA_MEMORY_JOURNAL_WIRE_MESSAGE_BYTES as u64,
attempted as u64,
)
}
pub async fn append(
&self,
declared: &DeclaredCall,
command: &AppendMemoryRecords,
) -> Result<MemoryAppendOutcome, StateError> {
let request = pb::AppendPersonaMemoryRecordsRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
command: buffa::MessageField::some(append_to_wire(command)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.append_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
append_outcome_from_wire(
required(reply.receipt, "an append returns its receipt")?,
reply.positions,
)
}
pub async fn rewrite(
&self,
declared: &DeclaredCall,
command: &RewriteMemoryRecords,
) -> Result<MemoryRewriteOutcome, StateError> {
let request = pb::RewritePersonaMemoryRecordsRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
command: buffa::MessageField::some(rewrite_to_wire(command)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.rewrite_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
rewrite_outcome_from_wire(
required(reply.receipt, "a rewrite returns its receipt")?,
reply.dropped,
)
}
pub async fn migrate(
&self,
declared: &DeclaredCall,
command: &MigrateMemoryRecords,
) -> Result<MemoryMigrationOutcome, StateError> {
let request = pb::MigratePersonaMemoryRecordsRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
command: buffa::MessageField::some(migrate_to_wire(command)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.migrate_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
migration_outcome_from_wire(
required(reply.receipt, "a migration returns its receipt")?,
reply.copied,
)
}
pub async fn destroy(
&self,
declared: &DeclaredCall,
command: &DestroyMemoryRecords,
) -> Result<Receipt, StateError> {
let request = pb::DestroyPersonaMemoryRecordsRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
command: buffa::MessageField::some(destroy_to_wire(command)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.destroy_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
Kernel::<Receipt>::try_from(required(
reply.receipt,
"a destruction returns its receipt",
)?)
.map(Kernel::into_inner)
}
pub async fn replay(
&self,
declared: &DeclaredCall,
namespace: &str,
partition: &MemoryJournalPartition,
range: MemoryReplayRange,
) -> Result<MemoryReplayPage, StateError> {
let request = pb::ReplayPersonaMemoryRecordsRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
namespace: namespace.to_owned(),
partition: partition.as_str().to_owned(),
start: range.start(),
max_records: range.max_records(),
max_bytes: range.max_bytes(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.replay_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
let page = replay_from_wire(reply)?;
if page.signed_head().is_some_and(|signed| {
signed.namespace().as_str() != namespace || signed.partition() != partition
}) {
return Err(StateError::Malformed {
field: "signed_head".into(),
reason: "a signed replay head belongs to the requested namespace and partition"
.into(),
});
}
Ok(page)
}
pub async fn list(
&self,
declared: &DeclaredCall,
namespace: &str,
after: Option<&str>,
limit: u32,
) -> Result<MemoryPartitionPage, StateError> {
let request = pb::ListPersonaMemoryPartitionsRequest {
context: buffa::MessageField::some(Kernel(declared).into()),
namespace: namespace.to_owned(),
after: after.map(str::to_owned),
limit,
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let attempted = buffa::Message::encoded_len(&request) as usize;
let reply = self
.inner
.list_with_options(request, bounded_traced_options(declared))
.await
.map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
.into_owned();
partition_page_from_wire(reply)
}
}
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: "reply".into(),
reason: reason.into(),
})
}