polyc-state-connect 2026.9.0

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 the narrow projection catalog.

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

use connectrpc::{ConnectError, RequestContext, Response, Router, ServiceRequest, ServiceResult};
use polyc_proto::proto::polychrome::state::v1::{
    GetProjectionReceiptReply, GetProjectionReceiptRequest, PublishProjectionManifestReply,
    PublishProjectionManifestRequest, RecordProjectionManifestReply,
    RecordProjectionManifestRequest, RegisterProjectionPublisherReply,
    RegisterProjectionPublisherRequest, ResolveProjectionManifestReply,
    ResolveProjectionManifestRequest, StateProjectionCatalogService,
    StateProjectionCatalogServiceExt,
};
use polyc_state::{
    context::CallContext,
    error::StateError,
    id::{CommandId, OwnerId},
    projection::{
        self, ProjectionCatalog, ProjectionGeneration, ProjectionHead, ProjectionKey,
        ProjectionManifest, PublishManifest, PublisherFence, PublisherId, ReadProjectionReceipt,
        RecordManifest, RegisterPublisher, ResolveManifest,
    },
    revision::JournalSource,
};

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

/// State's narrow projection catalog service.
pub struct ProjectionCatalogSvc {
    catalog: Arc<dyn ProjectionCatalog>,
    draining: Arc<AtomicBool>,
    binding: Arc<AudienceBinding>,
}

impl ProjectionCatalogSvc {
    /// Builds the service.
    #[must_use]
    pub const fn new(
        catalog: Arc<dyn ProjectionCatalog>,
        draining: Arc<AtomicBool>,
        binding: Arc<AudienceBinding>,
    ) -> Self {
        Self {
            catalog,
            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), 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 = projection::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),
        ))
    }
}

fn required<T: Default, P: buffa::ProtoBox<T>>(
    field: &str,
    reason: &str,
    value: buffa::MessageField<T, P>,
) -> Result<T, ConnectError> {
    value.into_option().ok_or_else(|| {
        to_connect_error(
            &StateError::Malformed {
                field: field.to_owned(),
                reason: reason.to_owned(),
            }
            .into(),
        )
    })
}

#[allow(refining_impl_trait)]
impl StateProjectionCatalogService for ProjectionCatalogSvc {
    async fn register_publisher(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, RegisterProjectionPublisherRequest>,
    ) -> ServiceResult<RegisterProjectionPublisherReply> {
        let RegisterProjectionPublisherRequest {
            context,
            metadata,
            key,
            source,
            publisher,
            caller,
            __buffa_unknown_fields: _,
        } = request.to_owned_message();
        let (context, span) = self.admit(&ctx, context, "RegisterProjectionPublisher")?;
        let _entered = span.enter();
        let (fence, receipt) = self
            .catalog
            .register(
                RegisterPublisher::new(
                    Kernel::try_from(required(
                        "metadata",
                        "a registration carries metadata",
                        metadata,
                    )?)
                    .map_err(|error: StateError| to_connect_error(&error.into()))?
                    .into_inner(),
                    Kernel::<ProjectionKey>::try_from(required(
                        "key",
                        "a registration names one key",
                        key,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    Kernel::<JournalSource>::try_from(required(
                        "source",
                        "a registration names one source",
                        source,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    PublisherId::new(publisher),
                    OwnerId::new(caller),
                ),
                &context,
            )
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(RegisterProjectionPublisherReply {
            fence: buffa::MessageField::some(Kernel(&fence).into()),
            receipt: buffa::MessageField::some(Kernel(&receipt).into()),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn record_manifest(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, RecordProjectionManifestRequest>,
    ) -> ServiceResult<RecordProjectionManifestReply> {
        let RecordProjectionManifestRequest {
            context,
            metadata,
            manifest,
            caller,
            signed_artifact,
            __buffa_unknown_fields: _,
        } = request.to_owned_message();
        let (context, span) = self.admit(&ctx, context, "RecordProjectionManifest")?;
        let _entered = span.enter();
        let receipt = self
            .catalog
            .record(
                RecordManifest::new(
                    Kernel::try_from(required("metadata", "a record carries metadata", metadata)?)
                        .map_err(|error: StateError| to_connect_error(&error.into()))?
                        .into_inner(),
                    Kernel::<ProjectionManifest>::try_from(required(
                        "manifest",
                        "a record carries one manifest",
                        manifest,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    signed_artifact,
                    OwnerId::new(caller),
                ),
                &context,
            )
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(RecordProjectionManifestReply {
            receipt: buffa::MessageField::some(Kernel(&receipt).into()),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn publish_manifest(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, PublishProjectionManifestRequest>,
    ) -> ServiceResult<PublishProjectionManifestReply> {
        let PublishProjectionManifestRequest {
            context,
            metadata,
            key,
            source,
            expected,
            publish,
            publisher,
            fence,
            caller,
            __buffa_unknown_fields: _,
        } = request.to_owned_message();
        let (context, span) = self.admit(&ctx, context, "PublishProjectionManifest")?;
        let _entered = span.enter();
        let receipt = self
            .catalog
            .publish(
                PublishManifest::new(
                    Kernel::try_from(required(
                        "metadata",
                        "a publication carries metadata",
                        metadata,
                    )?)
                    .map_err(|error: StateError| to_connect_error(&error.into()))?
                    .into_inner(),
                    Kernel::<ProjectionKey>::try_from(required(
                        "key",
                        "a publication names one key",
                        key,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    Kernel::<JournalSource>::try_from(required(
                        "source",
                        "a publication names one source",
                        source,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    ProjectionGeneration::new(expected),
                    ProjectionGeneration::new(publish),
                    PublisherId::new(publisher),
                    Kernel::<PublisherFence>::try_from(required(
                        "fence",
                        "a publication carries authority",
                        fence,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    OwnerId::new(caller),
                ),
                &context,
            )
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(PublishProjectionManifestReply {
            receipt: buffa::MessageField::some(Kernel(&receipt).into()),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn resolve_manifest(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, ResolveProjectionManifestRequest>,
    ) -> ServiceResult<ResolveProjectionManifestReply> {
        let ResolveProjectionManifestRequest {
            context,
            key,
            source,
            caller,
            __buffa_unknown_fields: _,
        } = request.to_owned_message();
        let (context, span) = self.admit(&ctx, context, "ResolveProjectionManifest")?;
        let _entered = span.enter();
        let resolution = self
            .catalog
            .resolve(
                ResolveManifest::new(
                    Kernel::<ProjectionKey>::try_from(required(
                        "key",
                        "a resolve names one key",
                        key,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    Kernel::<JournalSource>::try_from(required(
                        "source",
                        "a resolve names one source",
                        source,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    OwnerId::new(caller),
                ),
                &context,
            )
            .map_err(|error| to_connect_error(&error))?;
        // Every arm named. A superseded head sends its generation and its
        // lineage and withholds the manifest, so a reader on the other side
        // cannot receive rows for a journal that no longer exists.
        let (head, pending, highest_recorded) = resolution.into_parts();
        let (manifest, superseded_generation, superseded_source) = match &head {
            ProjectionHead::Absent => (
                buffa::MessageField::default(),
                None,
                buffa::MessageField::default(),
            ),
            ProjectionHead::Current(manifest) => (
                buffa::MessageField::some(Kernel(manifest.as_ref()).into()),
                None,
                buffa::MessageField::default(),
            ),
            ProjectionHead::Superseded { generation, source } => (
                buffa::MessageField::default(),
                Some(generation.get()),
                buffa::MessageField::some(Kernel(source.as_ref()).into()),
            ),
        };
        Response::ok(ResolveProjectionManifestReply {
            manifest,
            superseded_generation,
            superseded_source,
            pending_manifest: pending
                .as_deref()
                .map_or_else(buffa::MessageField::default, |manifest| {
                    buffa::MessageField::some(Kernel(manifest).into())
                }),
            highest_recorded_generation: highest_recorded.get(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn get_receipt(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, GetProjectionReceiptRequest>,
    ) -> ServiceResult<GetProjectionReceiptReply> {
        let GetProjectionReceiptRequest {
            context,
            key,
            source,
            command_id,
            caller,
            __buffa_unknown_fields: _,
        } = request.to_owned_message();
        let (context, span) = self.admit(&ctx, context, "GetProjectionReceipt")?;
        let _entered = span.enter();
        let receipt = self
            .catalog
            .receipt(
                ReadProjectionReceipt::new(
                    Kernel::<ProjectionKey>::try_from(required(
                        "key",
                        "a receipt read names one key",
                        key,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    Kernel::<JournalSource>::try_from(required(
                        "source",
                        "a receipt read names one source",
                        source,
                    )?)
                    .map_err(|error| to_connect_error(&error.into()))?
                    .into_inner(),
                    CommandId::new(command_id),
                    OwnerId::new(caller),
                ),
                &context,
            )
            .map_err(|error| to_connect_error(&error))?;
        Response::ok(GetProjectionReceiptReply {
            receipt: receipt
                .as_ref()
                .map_or_else(buffa::MessageField::default, |receipt| {
                    buffa::MessageField::some(Kernel(receipt).into())
                }),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }
}