polyc-state-connect 2026.9.0

State plane transport adapter: capability-specific Connect clients and server-trait glue mapping the generated wire types onto the polyc-state kernel — typed outcomes, per-call admission, and the conformance surface the authenticated shell proves itself against (docs/proposals/separated-planes.md).
//! Capability-specific persona-memory journal client.

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,
};

/// Client bound to the authoritative persona-memory event journal.
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,
{
    /// Builds a bounded client.
    #[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,
        )
    }

    /// Appends one atomic batch.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    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,
        )
    }

    /// Rewrites one exact incarnation.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    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,
        )
    }

    /// Migrates two partitions atomically.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    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,
        )
    }

    /// Destroys one fenced partition.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    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)
    }

    /// Replays one bounded range.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    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)
    }

    /// Lists one bounded partition page.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    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)
    }
}

// Named over `P: buffa::ProtoBox<T>` rather than `impl Into<Option<T>>`: see
// `crate::wire::required`'s doc comment.
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(),
    })
}