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).
//! Server adapter for credential verifier metadata and lifecycle.

use std::future::Future;
use std::sync::{
    Arc,
    atomic::{AtomicBool, Ordering},
};

use connectrpc::{ConnectError, RequestContext, Response, Router, ServiceRequest, ServiceResult};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
    context::CallContext,
    credentials::{CredentialDirectory, CredentialLifecycleRead, CredentialRead, CredentialWrite},
    error::StateError,
    id::{CommandId, NamespaceId},
    page::PageCompleteness,
    revision::{JournalPosition, PartitionIncarnation},
    versioned::{VersionedRead, VersionedTransact},
};

use crate::{
    admission::{
        AudienceBinding, PeerIdentity, check_audience_binding, check_call_context_version,
        check_not_draining, check_transport_deadline, slot_wait, state_audience,
    },
    blocking,
    error::to_connect_error,
    trace::adopt_caller_trace,
    wire::{DeclaredCall, Kernel, declared_call, fixed_bytes},
};

use super::wire::{
    change_to_wire, command_from_wire, index_to_wire, record_to_wire, request_to_wire,
};

/// Complete credential authority port.
pub trait CredentialAuthority: CredentialRead + CredentialWrite + CredentialLifecycleRead {
    /// Namespace this authority owns.
    fn namespace(&self) -> &NamespaceId;
}

impl<V> CredentialAuthority for CredentialDirectory<V>
where
    V: VersionedRead + VersionedTransact,
{
    fn namespace(&self) -> &NamespaceId {
        self.namespace()
    }
}

/// State's capability-specific credential service.
pub struct CredentialSvc {
    authority: Arc<dyn CredentialAuthority>,
    draining: Arc<AtomicBool>,
    binding: Arc<AudienceBinding>,
}

impl CredentialSvc {
    /// Builds the service.
    #[must_use]
    pub const fn new(
        authority: Arc<dyn CredentialAuthority>,
        draining: Arc<AtomicBool>,
        binding: Arc<AudienceBinding>,
    ) -> Self {
        Self {
            authority,
            draining,
            binding,
        }
    }

    /// Registers the capability-specific route.
    #[must_use]
    pub fn register_on(self, router: Router) -> Router {
        use pb::StateCredentialServiceExt as _;
        Arc::new(self).register(router)
    }

    fn peer(context: &RequestContext) -> PeerIdentity {
        PeerIdentity::from_verified_leaf(
            context
                .peer_certs()
                .and_then(<[_]>::first)
                .map(|leaf| &**leaf),
        )
    }

    fn admit(
        &self,
        context: &RequestContext,
        wire: impl Into<Option<pb::CallContext>>,
        method: &'static str,
    ) -> Result<(CallContext, tracing::Span, std::time::Duration), ConnectError> {
        check_not_draining(self.draining.load(Ordering::Relaxed))?;
        let declared: DeclaredCall =
            declared_call(wire).map_err(|error| to_connect_error(&error))?;
        let family = polyc_state::versioned::family();
        check_call_context_version(declared.version).map_err(|error| to_connect_error(&error))?;
        check_audience_binding(
            &Self::peer(context),
            &declared.audience,
            &state_audience(),
            &self.binding,
            &family,
        )
        .map_err(|error| to_connect_error(&error))?;
        check_transport_deadline(context.time_remaining(), &family)
            .map_err(|error| to_connect_error(&error))?;
        Ok((
            declared.origin_relative_context(),
            adopt_caller_trace(context.headers(), method),
            slot_wait(context.time_remaining(), &declared),
        ))
    }

    fn namespace(&self, value: &str) -> Result<(), ConnectError> {
        if value == self.authority.namespace().as_str() {
            Ok(())
        } else {
            Err(to_connect_error(&StateError::Denied {
                family: polyc_state::versioned::family(),
            }))
        }
    }
}

#[allow(refining_impl_trait)]
impl pb::StateCredentialService for CredentialSvc {
    fn transact(
        &self,
        context: RequestContext,
        request: ServiceRequest<'_, pb::TransactStateCredentialRequest>,
    ) -> impl Future<Output = ServiceResult<pb::TransactStateCredentialReply>> {
        let message = request.to_owned_message();
        async move {
            let (call, span, wait) = self.admit(&context, message.context, "TransactCredential")?;
            let metadata = message.metadata.into_option().ok_or_else(|| {
                to_connect_error(&StateError::Malformed {
                    field: "metadata".into(),
                    reason: "credential command carries metadata".into(),
                })
            })?;
            self.namespace(&metadata.namespace)?;
            let operation = message.operation.into_option().ok_or_else(|| {
                to_connect_error(&StateError::Malformed {
                    field: "operation".into(),
                    reason: "credential command carries one operation".into(),
                })
            })?;
            let command =
                command_from_wire(metadata, operation, message.request_record.into_option())
                    .map_err(|error| to_connect_error(&error))?;
            let authority = Arc::clone(&self.authority);
            let receipt =
                blocking::mutation(span, polyc_state::versioned::family(), wait, move || {
                    authority.transact(command, &call)
                })
                .await?;
            Response::ok(pb::TransactStateCredentialReply {
                receipt: buffa::MessageField::some(pb::Receipt::from(Kernel(&receipt))),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn get_credential(
        &self,
        context: RequestContext,
        request: ServiceRequest<'_, pb::GetStateCredentialRequest>,
    ) -> impl Future<Output = ServiceResult<pb::GetStateCredentialReply>> {
        let message = request.to_owned_message();
        async move {
            let (call, span, wait) = self.admit(&context, message.context, "GetCredential")?;
            self.namespace(&message.namespace)?;
            let credential_id = message.credential_id;
            let authority = Arc::clone(&self.authority);
            let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
                authority.credential(&credential_id, &call)
            })
            .await?;
            Response::ok(pb::GetStateCredentialReply {
                record: fact
                    .value()
                    .as_ref()
                    .map_or_else(buffa::MessageField::none, |record| {
                        buffa::MessageField::some(record_to_wire(record))
                    }),
                snapshot_revision: fact.snapshot_revision().get(),
                entry_revision: fact
                    .entry_revision()
                    .map(polyc_state::revision::Revision::get),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn get_request_record(
        &self,
        context: RequestContext,
        request: ServiceRequest<'_, pb::GetStateCredentialRequestRecordRequest>,
    ) -> impl Future<Output = ServiceResult<pb::GetStateCredentialRequestRecordReply>> {
        let message = request.to_owned_message();
        async move {
            let (call, span, wait) =
                self.admit(&context, message.context, "GetCredentialRequestRecord")?;
            self.namespace(&message.namespace)?;
            let operation_id = message.operation_id;
            let authority = Arc::clone(&self.authority);
            let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
                authority.request(&operation_id, &call)
            })
            .await?;
            Response::ok(pb::GetStateCredentialRequestRecordReply {
                request_record: fact
                    .value()
                    .as_ref()
                    .map_or_else(buffa::MessageField::none, |value| {
                        buffa::MessageField::some(request_to_wire(value))
                    }),
                snapshot_revision: fact.snapshot_revision().get(),
                entry_revision: fact
                    .entry_revision()
                    .map(polyc_state::revision::Revision::get),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn get_index(
        &self,
        context: RequestContext,
        request: ServiceRequest<'_, pb::GetStateCredentialIndexRequest>,
    ) -> impl Future<Output = ServiceResult<pb::GetStateCredentialIndexReply>> {
        let message = request.to_owned_message();
        async move {
            let (call, span, wait) = self.admit(&context, message.context, "GetCredentialIndex")?;
            self.namespace(&message.namespace)?;
            let authority = Arc::clone(&self.authority);
            let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
                authority.index(&call)
            })
            .await?;
            Response::ok(pb::GetStateCredentialIndexReply {
                index: buffa::MessageField::some(index_to_wire(fact.value())),
                snapshot_revision: fact.snapshot_revision().get(),
                entry_revision: fact
                    .entry_revision()
                    .map(polyc_state::revision::Revision::get),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn list_credentials(
        &self,
        context: RequestContext,
        request: ServiceRequest<'_, pb::ListStateCredentialsRequest>,
    ) -> impl Future<Output = ServiceResult<pb::ListStateCredentialsReply>> {
        let message = request.to_owned_message();
        async move {
            let (call, span, wait) = self.admit(&context, message.context, "ListCredentials")?;
            self.namespace(&message.namespace)?;
            let authority = Arc::clone(&self.authority);
            let page = blocking::read(span, polyc_state::versioned::family(), wait, move || {
                authority.list(message.after.as_deref(), message.limit, &call)
            })
            .await?;
            Response::ok(pb::ListStateCredentialsReply {
                records: page.records().iter().map(record_to_wire).collect(),
                next_after: page.next_after().map(str::to_owned),
                snapshot_revision: page.snapshot_revision().get(),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn list_changes(
        &self,
        context: RequestContext,
        request: ServiceRequest<'_, pb::ListStateCredentialChangesRequest>,
    ) -> impl Future<Output = ServiceResult<pb::ListStateCredentialChangesReply>> {
        let message = request.to_owned_message();
        async move {
            let (call, span, wait) =
                self.admit(&context, message.context, "ListCredentialChanges")?;
            self.namespace(&message.namespace)?;
            let after = message.after.map(JournalPosition::new);
            let expected_lineage = message
                .expected_lineage
                .as_deref()
                .map(|bytes| {
                    fixed_bytes::<32>("expected_lineage", bytes)
                        .map(PartitionIncarnation::from_bytes)
                })
                .transpose()
                .map_err(|error| to_connect_error(&error))?;
            let limit = message.limit;
            let authority = Arc::clone(&self.authority);
            // An absent cursor reads from the beginning; a present one resumes
            // strictly after it. The authority refuses a position it never issued,
            // so this maps the request and does not second-guess the answer.
            let page = blocking::read(span, polyc_state::versioned::family(), wait, move || {
                authority.changes(after, expected_lineage, limit, &call)
            })
            .await?;
            Response::ok(pb::ListStateCredentialChangesReply {
                changes: page.changes().iter().map(change_to_wire).collect(),
                next_after: page.next().map(|cursor| cursor.position().get()),
                completeness: buffa::EnumValue::Known(match page.completeness() {
                    PageCompleteness::Complete => pb::PageCompleteness::Complete,
                    PageCompleteness::Truncated => pb::PageCompleteness::Truncated,
                }),
                frontier: page.frontier().get(),
                lineage: page.lineage().as_bytes().to_vec(),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn get_receipt(
        &self,
        context: RequestContext,
        request: ServiceRequest<'_, pb::GetStateCredentialReceiptRequest>,
    ) -> impl Future<Output = ServiceResult<pb::GetStateCredentialReceiptReply>> {
        let message = request.to_owned_message();
        async move {
            let (_call, span, wait) =
                self.admit(&context, message.context, "GetCredentialReceipt")?;
            self.namespace(&message.namespace)?;
            let namespace = NamespaceId::new(message.namespace);
            let command_id = CommandId::new(message.command_id);
            let authority = Arc::clone(&self.authority);
            let receipt = blocking::read(span, polyc_state::versioned::family(), wait, move || {
                authority.committed_receipt(&namespace, &command_id)
            })
            .await?;
            Response::ok(pb::GetStateCredentialReceiptReply {
                receipt: receipt
                    .as_ref()
                    .map_or_else(buffa::MessageField::none, |value| {
                        buffa::MessageField::some(pb::Receipt::from(Kernel(value)))
                    }),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }
}