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).
//! Server adapter for durable sessions.

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

use buffa::EnumValue;
use connectrpc::{ConnectError, RequestContext, Response, Router, ServiceRequest, ServiceResult};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
    context::CallContext,
    error::StateError,
    id::{CommandId, NamespaceId},
    revision::Revision,
    sessions::{RefreshPreflight, SessionDirectory, SessionRead, SessionWrite},
    versioned::{VersionedRead, VersionedTransact, authority::SessionId},
};

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

use super::wire::{
    command_from_wire, expiration_to_wire, family_to_wire, principal_from_wire, record_to_wire,
};

/// Complete durable-session authority port.
pub trait SessionAuthority: SessionRead + SessionWrite {
    /// Namespace this authority owns.
    fn namespace(&self) -> &NamespaceId;
}

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

/// State's capability-specific durable-session service.
pub struct SessionSvc {
    authority: Arc<dyn SessionAuthority>,
    draining: Arc<AtomicBool>,
    binding: Arc<AudienceBinding>,
}

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

    /// Registers the capability-specific service.
    #[must_use]
    pub fn register_on(self, router: Router) -> Router {
        use pb::StateSessionServiceExt as _;
        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<pb::CallContext>>,
        method: &'static str,
    ) -> Result<(CallContext, tracing::Span), ConnectError> {
        check_not_draining(self.draining.load(Ordering::Relaxed))?;
        let declared: DeclaredCall =
            declared_call(context).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(ctx),
            &declared.audience,
            &state_audience(),
            &self.binding,
            &family,
        )
        .map_err(|error| to_connect_error(&error))?;
        check_transport_deadline(ctx.time_remaining(), &family)
            .map_err(|error| to_connect_error(&error))?;
        Ok((
            declared.origin_relative_context(),
            adopt_caller_trace(ctx.headers(), method),
        ))
    }

    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::StateSessionService for SessionSvc {
    async fn transact(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::TransactSessionRequest>,
    ) -> ServiceResult<pb::TransactSessionReply> {
        let message = request.to_owned_message();
        let (context, span) = self.admit(&ctx, message.context, "TransactSession")?;
        let _entered = span.enter();
        let metadata = message
            .metadata
            .into_option()
            .ok_or_else(|| malformed_connect("metadata", "session command carries metadata"))?;
        self.namespace(&metadata.namespace)?;
        let operation = message.operation.into_option().ok_or_else(|| {
            malformed_connect("operation", "session command carries one operation")
        })?;
        let receipt = self
            .authority
            .transact(
                command_from_wire(metadata, operation).map_err(|error| to_connect_error(&error))?,
                &context,
            )
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(pb::TransactSessionReply {
            receipt: buffa::MessageField::some(pb::Receipt::from(Kernel(&receipt))),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn get_family(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::GetSessionFamilyRequest>,
    ) -> ServiceResult<pb::GetSessionFamilyReply> {
        let message = request.to_owned_message();
        let (context, span) = self.admit(&ctx, message.context, "GetSessionFamily")?;
        let _entered = span.enter();
        self.namespace(&message.namespace)?;
        let fact = self
            .authority
            .family(&SessionId::new(message.family_id), &context)
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(pb::GetSessionFamilyReply {
            family: fact
                .value()
                .as_ref()
                .map_or_else(buffa::MessageField::none, |value| {
                    buffa::MessageField::some(family_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(),
        })
    }

    async fn inspect_refresh(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::InspectSessionRefreshRequest>,
    ) -> ServiceResult<pb::InspectSessionRefreshReply> {
        let message = request.to_owned_message();
        let (context, span) = self.admit(&ctx, message.context, "InspectSessionRefresh")?;
        let _entered = span.enter();
        self.namespace(&message.namespace)?;
        let result = self
            .authority
            .inspect_refresh(
                &SessionId::new(message.family_id),
                message.presented_generation,
                message.now_ms,
                &context,
            )
            .map_err(|error| to_connect_error(&error))?;
        let (classification, family, snapshot_revision, entry_revision) = match result {
            RefreshPreflight::Rotate(fact) => fact_parts(pb::SessionRefreshClass::Rotate, fact),
            RefreshPreflight::Replay(fact) => fact_parts(pb::SessionRefreshClass::Replay, fact),
            RefreshPreflight::ReuseDetected(fact) => {
                fact_parts(pb::SessionRefreshClass::ReuseDetected, fact)
            }
            RefreshPreflight::Expired(fact) => fact_parts(pb::SessionRefreshClass::Expired, fact),
            RefreshPreflight::UnknownFamily(fact) => (
                pb::SessionRefreshClass::UnknownFamily,
                None,
                fact.snapshot_revision().get(),
                fact.entry_revision().map(Revision::get),
            ),
        };
        Response::ok(pb::InspectSessionRefreshReply {
            classification: EnumValue::Known(classification),
            family: family.into(),
            snapshot_revision,
            entry_revision,
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn get_expiration_shard(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::GetSessionExpirationShardRequest>,
    ) -> ServiceResult<pb::GetSessionExpirationShardReply> {
        let message = request.to_owned_message();
        let (context, span) = self.admit(&ctx, message.context, "GetSessionExpirationShard")?;
        let _entered = span.enter();
        self.namespace(&message.namespace)?;
        let shard = u16::try_from(message.shard)
            .map_err(|_| malformed_connect("shard", "expiry shard fits u16"))?;
        let fact = self
            .authority
            .expiration_shard(shard, &context)
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(pb::GetSessionExpirationShardReply {
            expiration: buffa::MessageField::some(expiration_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(),
        })
    }

    async fn get_bearer_authorization(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::GetBearerAuthorizationRequest>,
    ) -> ServiceResult<pb::GetBearerAuthorizationReply> {
        let message = request.to_owned_message();
        let (context, span) = self.admit(&ctx, message.context, "GetBearerAuthorization")?;
        let _entered = span.enter();
        self.namespace(&message.namespace)?;
        let principal = principal_from_wire(message.principal.into_option().ok_or_else(|| {
            malformed_connect("principal", "authorization request carries a principal")
        })?)
        .map_err(|error| to_connect_error(&error))?;
        let fact = self
            .authority
            .bearer_authorization(&principal, &SessionId::new(message.session), &context)
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(pb::GetBearerAuthorizationReply {
            record: buffa::MessageField::some(record_to_wire(fact.value().record())),
            current_epoch: fact.value().current_epoch().get(),
            snapshot_revision: fact.snapshot_revision().get(),
            entry_revision: fact
                .entry_revision()
                .map(polyc_state::revision::Revision::get),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn get_bearer_record(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::GetBearerRecordRequest>,
    ) -> ServiceResult<pb::GetBearerRecordReply> {
        let message = request.to_owned_message();
        let (context, span) = self.admit(&ctx, message.context, "GetBearerRecord")?;
        let _entered = span.enter();
        self.namespace(&message.namespace)?;
        let fact = self
            .authority
            .bearer_record(&SessionId::new(message.session), &context)
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(pb::GetBearerRecordReply {
            record: buffa::MessageField::some(record_to_wire(*fact.value())),
            snapshot_revision: fact.snapshot_revision().get(),
            entry_revision: fact.entry_revision().map(Revision::get),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn get_authorization_epoch(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::GetSessionAuthorizationEpochRequest>,
    ) -> ServiceResult<pb::GetSessionAuthorizationEpochReply> {
        let message = request.to_owned_message();
        let (context, span) = self.admit(&ctx, message.context, "GetSessionAuthorizationEpoch")?;
        let _entered = span.enter();
        self.namespace(&message.namespace)?;
        let principal =
            principal_from_wire(message.principal.into_option().ok_or_else(|| {
                malformed_connect("principal", "epoch request carries a principal")
            })?)
            .map_err(|error| to_connect_error(&error))?;
        let fact = self
            .authority
            .authorization_epoch(&principal, &context)
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(pb::GetSessionAuthorizationEpochReply {
            epoch: fact.value().get(),
            snapshot_revision: fact.snapshot_revision().get(),
            entry_revision: fact
                .entry_revision()
                .map(polyc_state::revision::Revision::get),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn get_receipt(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::GetSessionReceiptRequest>,
    ) -> ServiceResult<pb::GetSessionReceiptReply> {
        let message = request.to_owned_message();
        let (_context, span) = self.admit(&ctx, message.context, "GetSessionReceipt")?;
        let _entered = span.enter();
        self.namespace(&message.namespace)?;
        let receipt = self
            .authority
            .committed_receipt(
                &NamespaceId::new(message.namespace),
                &CommandId::new(message.command_id),
            )
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(pb::GetSessionReceiptReply {
            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(),
        })
    }
}

fn fact_parts(
    classification: pb::SessionRefreshClass,
    fact: polyc_state::sessions::SessionFact<polyc_state::sessions::SessionFamily>,
) -> (
    pb::SessionRefreshClass,
    Option<pb::StateSessionFamily>,
    u64,
    Option<u64>,
) {
    let snapshot = fact.snapshot_revision().get();
    let entry = fact
        .entry_revision()
        .map(polyc_state::revision::Revision::get);
    let family = Some(family_to_wire(&fact.into_value()));
    (classification, family, snapshot, entry)
}

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