polyc-state-connect 2026.10.2

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 per-persona usage client.
//!
//! Async end to end. The generated client is async, every method here awaits
//! it, and nothing in this family reaches for a blocking bridge or a blocking
//! pool.
//!
//! Every call carries the caller's remaining budget in its
//! [`DeclaredCall`](crate::wire::DeclaredCall), and the caller is expected to
//! supply a real one. This family gates a quota: a read that never returns is
//! a turn that hangs, which is the outcome failing closed exists to avoid.

use connectrpc::client::{ClientConfig, ClientTransport};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
    error::StateError,
    id::{CommandId, NamespaceId},
    receipt::Receipt,
    revision::Revision,
    usage::{
        AccrualView, ConversationId, ConversationUsageRecord, PersonaId, UsageCommand, UsageFact,
        UsageRollupRecord,
    },
};

use crate::{
    MAX_WIRE_MESSAGE_BYTES,
    error::{TransportFallback, from_connect_error},
    trace::bounded_traced_options,
    wire::{DeclaredCall, Kernel},
};

use super::wire::{conversation_from_wire, metadata_to_wire, operation_to_wire, rollup_from_wire};

/// Client bound to one usage namespace and no other State capability.
pub struct UsageClient<T> {
    inner: pb::StateUsageServiceClient<T>,
    namespace: NamespaceId,
}

impl<T> UsageClient<T>
where
    T: ClientTransport,
    <T::ResponseBody as connectrpc::http_body::Body>::Error: std::fmt::Display,
{
    /// Builds a namespace-bound client.
    #[must_use]
    pub fn new(transport: T, config: ClientConfig, namespace: NamespaceId) -> Self {
        Self {
            inner: pb::StateUsageServiceClient::new(
                transport,
                config.with_default_max_message_size(MAX_WIRE_MESSAGE_BYTES),
            ),
            namespace,
        }
    }

    fn fallback(declared: &DeclaredCall, attempted: usize) -> TransportFallback {
        TransportFallback::new(
            polyc_state::usage::family(),
            MAX_WIRE_MESSAGE_BYTES as u64,
            attempted as u64,
            declared.budget,
        )
    }

    /// Commits one typed usage mutation.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure. A turn already folded into
    /// its conversation's row comes back as a refusal, never as a success that
    /// moved a total.
    pub async fn transact(
        &self,
        declared: &DeclaredCall,
        command: &UsageCommand,
    ) -> Result<Receipt, StateError> {
        if command.metadata().scope().namespace() != &self.namespace {
            return Err(StateError::Denied {
                family: polyc_state::usage::family(),
            });
        }
        let request = pb::TransactStateUsageRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            metadata: buffa::MessageField::some(metadata_to_wire(command)),
            persona_id: command.persona().as_str().to_owned(),
            operation: buffa::MessageField::some(operation_to_wire(command.operation())),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .transact_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 usage mutation returns a receipt".into(),
            }
        })?)
        .map(Kernel::into_inner)
    }

    /// Reads one persona's maintained totals.
    ///
    /// The value is [`None`] for a persona holding no rollup. A quota gate
    /// takes the zero with `unwrap_or_default`; a merge needs the distinction.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn rollup(
        &self,
        declared: &DeclaredCall,
        persona: &PersonaId,
    ) -> Result<UsageFact<Option<UsageRollupRecord>>, StateError> {
        let request = pb::GetStateUsageRollupRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            namespace: self.namespace.as_str().to_owned(),
            persona_id: persona.as_str().to_owned(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .get_rollup_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
            .into_owned();
        Ok(UsageFact::observed(
            reply.rollup.into_option().as_ref().map(rollup_from_wire),
            Revision::new(reply.snapshot_revision),
            reply.entry_revision.map(Revision::new),
        ))
    }

    /// Reads one conversation's accrual row for one persona.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn conversation(
        &self,
        declared: &DeclaredCall,
        persona: &PersonaId,
        conversation: &ConversationId,
    ) -> Result<UsageFact<ConversationUsageRecord>, StateError> {
        let request = pb::GetStateUsageConversationRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            namespace: self.namespace.as_str().to_owned(),
            persona_id: persona.as_str().to_owned(),
            conversation_id: conversation.as_str().to_owned(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .get_conversation_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
            .into_owned();
        Ok(UsageFact::observed(
            conversation_from_wire(reply.record.into_option().ok_or_else(|| {
                StateError::Malformed {
                    field: "record".into(),
                    reason: "a usage reply carries its accrual row".into(),
                }
            })?),
            Revision::new(reply.snapshot_revision),
            reply.entry_revision.map(Revision::new),
        ))
    }

    /// Reads the two rows one accrual moves, at one revision.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn accrual_view(
        &self,
        declared: &DeclaredCall,
        persona: &PersonaId,
        conversation: &ConversationId,
    ) -> Result<AccrualView, StateError> {
        let request = pb::GetStateUsageAccrualViewRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            namespace: self.namespace.as_str().to_owned(),
            persona_id: persona.as_str().to_owned(),
            conversation_id: conversation.as_str().to_owned(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .get_accrual_view_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
            .into_owned();
        // An absent rollup is a real answer here, not a malformed reply: it is
        // how a persona that holds none, or whose totals a fold absorbed,
        // crosses the wire.
        let rollup = reply.rollup.into_option();
        Ok(AccrualView::observed(
            rollup.as_ref().map(rollup_from_wire),
            conversation_from_wire(reply.record.into_option().ok_or_else(|| {
                StateError::Malformed {
                    field: "record".into(),
                    reason: "an accrual view carries its accrual row".into(),
                }
            })?),
            Revision::new(reply.snapshot_revision),
            reply.rollup_revision.map(Revision::new),
            reply.record_revision.map(Revision::new),
        ))
    }

    /// Settles an ambiguous command from its durable receipt.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn committed_receipt(
        &self,
        declared: &DeclaredCall,
        command_id: &CommandId,
    ) -> Result<Option<Receipt>, StateError> {
        let request = pb::GetStateUsageReceiptRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            namespace: self.namespace.as_str().to_owned(),
            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(declared, attempted)))?
            .into_owned();
        reply
            .receipt
            .into_option()
            .map(|value| Kernel::<Receipt>::try_from(value).map(Kernel::into_inner))
            .transpose()
    }
}