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 Query's audit capability.

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

use connectrpc::{ConnectError, RequestContext, Response, Router, ServiceRequest, ServiceResult};
use polyc_proto::proto::polychrome::state::v1::{
    BeginQueryAuditReply, BeginQueryAuditRequest, CompleteQueryAuditReply,
    CompleteQueryAuditRequest, GetQueryAuditReceiptReply, GetQueryAuditReceiptRequest,
    GetQueryAuditReply, GetQueryAuditRequest, ListUnmatchedQueryAuditsReply,
    ListUnmatchedQueryAuditsRequest, StateQueryAuditService, StateQueryAuditServiceExt,
};
use polyc_state::{
    command::CommandEnvelope,
    context::CallContext,
    error::StateError,
    id::{Audience, NamespaceId, Purpose},
    query_audit::{
        self, AuditPhase, BeginOutcome, BeginQueryAudit, ListUnmatchedIntents,
        QueryAuditCompletionFactory, QueryAuditHistory, QueryAuditRead, QueryAuditWrite, QueryId,
        ReadQueryAudit, RecordedCompletionInput, RequesterId,
    },
};

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 as state_to_connect_error,
    query_audit::{error::to_connect_error, wire::digest},
    trace::adopt_caller_trace,
    wire::{DeclaredCall, Kernel, declared_call},
};

/// The complete query-audit authority port.
///
/// The history read is part of the authority, not a second one: the ordinals
/// it serves are assigned by the same writes. The separation the design asks
/// for is on the client side, where the history is its own client type and its
/// own route, so a projector can fold the history without holding the client
/// that records phases.
pub trait QueryAuditAuthority:
    QueryAuditRead + QueryAuditWrite + QueryAuditCompletionFactory + QueryAuditHistory
{
}

impl<T> QueryAuditAuthority for T where
    T: QueryAuditRead + QueryAuditWrite + QueryAuditCompletionFactory + QueryAuditHistory
{
}

/// Query's sole durable-write surface.
pub struct QueryAuditSvc {
    audit: Arc<dyn QueryAuditAuthority>,
    draining: Arc<AtomicBool>,
    binding: Arc<AudienceBinding>,
}

impl QueryAuditSvc {
    /// Builds the service over one query-audit authority.
    #[must_use]
    pub const fn new(
        audit: Arc<dyn QueryAuditAuthority>,
        draining: Arc<AtomicBool>,
        binding: Arc<AudienceBinding>,
    ) -> Self {
        Self {
            audit,
            draining,
            binding,
        }
    }

    /// Registers the capability-specific service.
    #[must_use]
    pub fn register_on(self, router: Router) -> Router {
        Arc::new(self).register(router)
    }

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

    fn admit(
        &self,
        ctx: &RequestContext,
        context: impl Into<Option<polyc_proto::proto::polychrome::state::v1::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(context).map_err(|error| state_to_connect_error(&error))?;
        let family = query_audit::family();
        check_call_context_version(declared.version)
            .map_err(|error| state_to_connect_error(&error))?;
        check_audience_binding(
            &Self::peer(ctx),
            &declared.audience,
            &state_audience(),
            &self.binding,
            &family,
        )
        .map_err(|error| state_to_connect_error(&error))?;
        check_transport_deadline(ctx.time_remaining(), &family)
            .map_err(|error| state_to_connect_error(&error))?;
        Ok((
            declared.origin_relative_context(),
            adopt_caller_trace(ctx.headers(), method),
            slot_wait(ctx.time_remaining(), &declared),
        ))
    }
}

fn envelope(purpose: String, audience: String) -> CommandEnvelope {
    CommandEnvelope::new(
        Purpose::new(purpose),
        Audience::new(audience),
        query_audit::command_bounds(),
    )
}

fn required_field(field: &str, reason: &str) -> ConnectError {
    to_connect_error(
        &StateError::Malformed {
            field: field.to_owned(),
            reason: reason.to_owned(),
        }
        .into(),
    )
}

#[allow(refining_impl_trait)]
impl StateQueryAuditService for QueryAuditSvc {
    fn begin(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, BeginQueryAuditRequest>,
    ) -> impl Future<Output = ServiceResult<BeginQueryAuditReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) = self.admit(&ctx, message.context, "BeginQueryAudit")?;
            let command = BeginQueryAudit::new(
                QueryId::new(message.query),
                NamespaceId::new(message.namespace),
                RequesterId::new(message.requester),
                digest("shape", &message.shape)
                    .map_err(|error| required_field("shape", &error.to_string()))?,
                Kernel::try_from(message.source.into_option().ok_or_else(|| {
                    required_field("source", "an intent carries its exact source vector")
                })?)
                .map_err(|error: StateError| to_connect_error(&error.into()))?
                .into_inner(),
                digest("digest", &message.digest)
                    .map_err(|error| required_field("digest", &error.to_string()))?,
                envelope(message.purpose, message.command_audience),
            );
            let audit = Arc::clone(&self.audit);
            let outcome = blocking::mutation(span, query_audit::family(), wait, move || {
                audit.begin(command, &context)
            })
            .await?;
            let receipt = match &outcome {
                BeginOutcome::Granted(permit) => permit.receipt(),
                BeginOutcome::AlreadyRecorded(receipt) => receipt,
            };
            Response::ok(BeginQueryAuditReply {
                receipt: buffa::MessageField::some(Kernel(receipt).into()),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn complete(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, CompleteQueryAuditRequest>,
    ) -> impl Future<Output = ServiceResult<CompleteQueryAuditReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) = self.admit(&ctx, message.context, "CompleteQueryAudit")?;
            let query = QueryId::new(message.query);
            let namespace = NamespaceId::new(message.namespace);
            let intent_source =
                Kernel::try_from(message.intent_source.into_option().ok_or_else(|| {
                    required_field(
                        "intent_source",
                        "a completion carries the source vector admitted by its intent",
                    )
                })?)
                .map_err(|error: StateError| to_connect_error(&error.into()))?
                .into_inner();
            let completion =
                Kernel::try_from(message.completion.into_option().ok_or_else(|| {
                    required_field("completion", "a completion call carries what the query did")
                })?)
                .map_err(|error: StateError| to_connect_error(&error.into()))?
                .into_inner();
            let input = RecordedCompletionInput::new(
                query,
                namespace,
                intent_source,
                completion,
                digest("digest", &message.digest)
                    .map_err(|error| required_field("digest", &error.to_string()))?,
                envelope(message.purpose, message.command_audience),
            );
            // The pair runs inside one closure because the whole authority
            // body belongs on one pool thread: one slot, one thread hop.
            // Each call still takes the authority's mutex separately, so a
            // second audit call can interleave between them exactly as it
            // could when this ran inline.
            let audit = Arc::clone(&self.audit);
            let receipt = blocking::mutation(span, query_audit::family(), wait, move || {
                let command = audit.prepare_recorded_completion(input, &context)?;
                audit.complete(command, &context)
            })
            .await?;
            Response::ok(CompleteQueryAuditReply {
                receipt: buffa::MessageField::some(Kernel(&receipt).into()),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn get_audit(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, GetQueryAuditRequest>,
    ) -> impl Future<Output = ServiceResult<GetQueryAuditReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) = self.admit(&ctx, message.context, "GetQueryAudit")?;
            let authority = Arc::clone(&self.audit);
            let audit = blocking::read(span, query_audit::family(), wait, move || {
                authority.audit(
                    ReadQueryAudit::new(
                        QueryId::new(message.query),
                        NamespaceId::new(message.namespace),
                    ),
                    &context,
                )
            })
            .await?;
            Response::ok(GetQueryAuditReply {
                audit: audit
                    .as_ref()
                    .map_or_else(buffa::MessageField::default, |audit| {
                        buffa::MessageField::some(Kernel(audit).into())
                    }),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn list_unmatched(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, ListUnmatchedQueryAuditsRequest>,
    ) -> impl Future<Output = ServiceResult<ListUnmatchedQueryAuditsReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) =
                self.admit(&ctx, message.context, "ListUnmatchedQueryAudits")?;
            let page = Kernel::try_from(message.page.into_option().ok_or_else(|| {
                required_field("page", "a bounded listing carries its page request")
            })?)
            .map_err(|error: StateError| to_connect_error(&error.into()))?
            .into_inner();
            let authority = Arc::clone(&self.audit);
            let result = blocking::read(span, query_audit::family(), wait, move || {
                authority.unmatched_intents(
                    ListUnmatchedIntents::new(NamespaceId::new(message.namespace), page),
                    &context,
                )
            })
            .await?;
            Response::ok(ListUnmatchedQueryAuditsReply {
                page: buffa::MessageField::some(Kernel(&result).into()),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn get_receipt(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, GetQueryAuditReceiptRequest>,
    ) -> impl Future<Output = ServiceResult<GetQueryAuditReceiptReply>> {
        let message = request.to_owned_message();
        async move {
            let (_context, span, wait) =
                self.admit(&ctx, message.context, "GetQueryAuditReceipt")?;
            let phase = if message.completion {
                AuditPhase::Completion
            } else {
                AuditPhase::Intent
            };
            let authority = Arc::clone(&self.audit);
            let receipt = blocking::read(span, query_audit::family(), wait, move || {
                authority.recorded_receipt(
                    &QueryId::new(message.query),
                    &NamespaceId::new(message.namespace),
                    phase,
                )
            })
            .await?;
            Response::ok(GetQueryAuditReceiptReply {
                receipt: receipt
                    .as_ref()
                    .map_or_else(buffa::MessageField::default, |receipt| {
                        buffa::MessageField::some(Kernel(receipt).into())
                    }),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }
}