polyc-state-connect 2026.8.3

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 durable work-claims client.

use connectrpc::client::{ClientConfig, ClientTransport};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
    claims::{
        AcquireClaim, AttemptHistory, ClaimStatus, CompleteClaim, GetAttemptHistory,
        GetClaimStatus, Granted, ReleaseClaim, RenewClaim,
    },
    command::CommandScope,
    error::StateError,
    id::CommandId,
    receipt::Receipt,
};

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

use super::wire::{
    acquire_to_wire, complete_to_wire, granted_from_wire, history_from_wire, release_to_wire,
    renew_to_wire, scope_to_wire, status_from_wire,
};

/// Client bound to the durable work-claims capability.
pub struct ClaimsClient<T> {
    inner: pb::StateClaimsServiceClient<T>,
}

impl<T> ClaimsClient<T>
where
    T: ClientTransport,
    <T::ResponseBody as connectrpc::http_body::Body>::Error: std::fmt::Display,
{
    /// Builds a durable work-claims client.
    #[must_use]
    pub fn new(transport: T, config: ClientConfig) -> Self {
        Self {
            inner: pb::StateClaimsServiceClient::new(
                transport,
                config.with_default_max_message_size(MAX_CLAIMS_WIRE_MESSAGE_BYTES),
            ),
        }
    }

    fn fallback(attempted: usize) -> TransportFallback {
        TransportFallback::new(
            polyc_state::claims::family(),
            MAX_CLAIMS_WIRE_MESSAGE_BYTES as u64,
            attempted as u64,
        )
    }

    /// Acquires one work item.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn acquire(
        &self,
        declared: &DeclaredCall,
        command: &AcquireClaim,
    ) -> Result<Granted, StateError> {
        let request = pb::AcquireClaimRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            command: buffa::MessageField::some(acquire_to_wire(command)),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .acquire_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
            .into_owned();
        granted_from_wire(required(
            reply.granted,
            "a successful acquire returns its grant",
        )?)
    }

    /// Renews one held work item.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn renew(
        &self,
        declared: &DeclaredCall,
        command: &RenewClaim,
    ) -> Result<Granted, StateError> {
        let request = pb::RenewClaimRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            command: buffa::MessageField::some(renew_to_wire(command)),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .renew_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
            .into_owned();
        granted_from_wire(required(
            reply.granted,
            "a successful renewal returns its grant",
        )?)
    }

    /// Completes one held work item terminally.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn complete(
        &self,
        declared: &DeclaredCall,
        command: &CompleteClaim,
    ) -> Result<Receipt, StateError> {
        let request = pb::CompleteClaimRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            command: buffa::MessageField::some(complete_to_wire(command)),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .complete_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
            .into_owned();
        receipt(reply.receipt, "a successful completion returns its receipt")
    }

    /// Releases one held work item for another attempt.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn release(
        &self,
        declared: &DeclaredCall,
        command: &ReleaseClaim,
    ) -> Result<Receipt, StateError> {
        let request = pb::ReleaseClaimRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            command: buffa::MessageField::some(release_to_wire(command)),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .release_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
            .into_owned();
        receipt(reply.receipt, "a successful release returns its receipt")
    }

    /// Reads one work item's status.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn status(
        &self,
        declared: &DeclaredCall,
        request: &GetClaimStatus,
    ) -> Result<ClaimStatus, StateError> {
        let request = pb::GetClaimStatusRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            scope: buffa::MessageField::some(scope_to_wire(request.scope())),
            work: request.work().as_str().to_owned(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .get_status_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
            .into_owned();
        status_from_wire(required(
            reply.status,
            "a successful read returns claim status",
        )?)
    }

    /// Reads one work item's bounded attempt history.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn history(
        &self,
        declared: &DeclaredCall,
        request: &GetAttemptHistory,
    ) -> Result<AttemptHistory, StateError> {
        let request = pb::GetClaimHistoryRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            scope: buffa::MessageField::some(scope_to_wire(request.scope())),
            work: request.work().as_str().to_owned(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .get_history_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(attempted)))?
            .into_owned();
        history_from_wire(required(
            reply.history,
            "a successful read returns attempt history",
        )?)
    }

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

// 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: "response".to_owned(),
        reason: reason.to_owned(),
    })
}

fn receipt(field: impl Into<Option<pb::Receipt>>, reason: &str) -> Result<Receipt, StateError> {
    let field = field.into().ok_or_else(|| StateError::Malformed {
        field: "response".to_owned(),
        reason: reason.to_owned(),
    })?;
    Kernel::<Receipt>::try_from(field).map(Kernel::into_inner)
}