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 immutable-object metadata.

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

use connectrpc::{ConnectError, RequestContext, Response, Router, ServiceRequest, ServiceResult};
use polyc_proto::proto::polychrome::state::v1::{
    GetObjectGenerationReply, GetObjectGenerationRequest, GetObjectHeadReply, GetObjectHeadRequest,
    GetObjectReceiptReply, GetObjectReceiptRequest, ListObjectGenerationsReply,
    ListObjectGenerationsRequest, PruneObjectGenerationReply, PruneObjectGenerationRequest,
    RecordObjectGenerationReply, RecordObjectGenerationRequest, StateImmutableMetadataService,
    StateImmutableMetadataServiceExt, UpdateObjectHeadReply, UpdateObjectHeadRequest,
};
use polyc_state::{
    context::CallContext,
    error::StateError,
    id::{CommandId, OwnerId},
    immutable::{
        self, Generation, ImmutableObjectRead, ImmutableObjectWrite, ListObjectGenerations,
        ObjectId, PruneObjectGeneration, ReadObjectGeneration, ReadObjectHead,
        RecordObjectGeneration, UpdateObjectHead,
    },
};

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

/// The complete immutable-metadata authority port.
pub trait ImmutableMetadataAuthority: ImmutableObjectRead + ImmutableObjectWrite {}

impl<T> ImmutableMetadataAuthority for T where T: ImmutableObjectRead + ImmutableObjectWrite {}

/// State's metadata authority for immutable object generations.
pub struct ImmutableMetadataSvc {
    metadata: Arc<dyn ImmutableMetadataAuthority>,
    draining: Arc<AtomicBool>,
    binding: Arc<AudienceBinding>,
}

impl ImmutableMetadataSvc {
    /// Builds the service over one metadata authority.
    #[must_use]
    pub const fn new(
        metadata: Arc<dyn ImmutableMetadataAuthority>,
        draining: Arc<AtomicBool>,
        binding: Arc<AudienceBinding>,
    ) -> Self {
        Self {
            metadata,
            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 = immutable::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 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 StateImmutableMetadataService for ImmutableMetadataSvc {
    fn record(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, RecordObjectGenerationRequest>,
    ) -> impl Future<Output = ServiceResult<RecordObjectGenerationReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) =
                self.admit(&ctx, message.context, "RecordObjectGeneration")?;
            let metadata = Kernel::try_from(message.metadata.into_option().ok_or_else(|| {
                required_field("metadata", "an object command carries its metadata")
            })?)
            .map_err(|error: StateError| to_connect_error(&error.into()))?
            .into_inner();
            let descriptor =
                Kernel::try_from(message.descriptor.into_option().ok_or_else(|| {
                    required_field("descriptor", "a record carries its object descriptor")
                })?)
                .map_err(|error: StateError| to_connect_error(&error.into()))?
                .into_inner();
            let authority = Arc::clone(&self.metadata);
            let receipt = blocking::mutation(span, immutable::family(), wait, move || {
                authority.record(
                    RecordObjectGeneration::new(metadata, descriptor, OwnerId::new(message.caller)),
                    &context,
                )
            })
            .await?;
            Response::ok(RecordObjectGenerationReply {
                receipt: buffa::MessageField::some(Kernel(&receipt).into()),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn update_head(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, UpdateObjectHeadRequest>,
    ) -> impl Future<Output = ServiceResult<UpdateObjectHeadReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) = self.admit(&ctx, message.context, "UpdateObjectHead")?;
            let metadata = Kernel::try_from(message.metadata.into_option().ok_or_else(|| {
                required_field("metadata", "an object command carries its metadata")
            })?)
            .map_err(|error: StateError| to_connect_error(&error.into()))?
            .into_inner();
            let authority = Arc::clone(&self.metadata);
            let receipt = blocking::mutation(span, immutable::family(), wait, move || {
                authority.update_head(
                    UpdateObjectHead::new(
                        metadata,
                        ObjectId::new(message.object),
                        Generation::new(message.expected),
                        Generation::new(message.publish),
                        OwnerId::new(message.caller),
                    ),
                    &context,
                )
            })
            .await?;
            Response::ok(UpdateObjectHeadReply {
                receipt: buffa::MessageField::some(Kernel(&receipt).into()),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn prune(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, PruneObjectGenerationRequest>,
    ) -> impl Future<Output = ServiceResult<PruneObjectGenerationReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) =
                self.admit(&ctx, message.context, "PruneObjectGeneration")?;
            let metadata = Kernel::try_from(message.metadata.into_option().ok_or_else(|| {
                required_field("metadata", "an object command carries its metadata")
            })?)
            .map_err(|error: StateError| to_connect_error(&error.into()))?
            .into_inner();
            let authority = Arc::clone(&self.metadata);
            let receipt = blocking::mutation(span, immutable::family(), wait, move || {
                authority.prune(
                    PruneObjectGeneration::new(
                        metadata,
                        ObjectId::new(message.object),
                        Generation::new(message.generation),
                        OwnerId::new(message.caller),
                    ),
                    &context,
                )
            })
            .await?;
            Response::ok(PruneObjectGenerationReply {
                receipt: buffa::MessageField::some(Kernel(&receipt).into()),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn get_head(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, GetObjectHeadRequest>,
    ) -> impl Future<Output = ServiceResult<GetObjectHeadReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) = self.admit(&ctx, message.context, "GetObjectHead")?;
            let authority = Arc::clone(&self.metadata);
            let head = blocking::read(span, immutable::family(), wait, move || {
                authority.head(
                    ReadObjectHead::new(
                        ObjectId::new(message.object),
                        OwnerId::new(message.caller),
                    ),
                    &context,
                )
            })
            .await?;
            Response::ok(GetObjectHeadReply {
                head: head
                    .as_ref()
                    .map_or_else(buffa::MessageField::default, |head| {
                        buffa::MessageField::some(Kernel(head).into())
                    }),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn get_generation(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, GetObjectGenerationRequest>,
    ) -> impl Future<Output = ServiceResult<GetObjectGenerationReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) = self.admit(&ctx, message.context, "GetObjectGeneration")?;
            let authority = Arc::clone(&self.metadata);
            let metadata = blocking::read(span, immutable::family(), wait, move || {
                authority.generation(
                    ReadObjectGeneration::new(
                        ObjectId::new(message.object),
                        Generation::new(message.generation),
                        OwnerId::new(message.caller),
                    ),
                    &context,
                )
            })
            .await?;
            Response::ok(GetObjectGenerationReply {
                metadata: metadata
                    .as_ref()
                    .map_or_else(buffa::MessageField::default, |metadata| {
                        buffa::MessageField::some(Kernel(metadata).into())
                    }),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn list_generations(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, ListObjectGenerationsRequest>,
    ) -> impl Future<Output = ServiceResult<ListObjectGenerationsReply>> {
        let message = request.to_owned_message();
        async move {
            let (context, span, wait) =
                self.admit(&ctx, message.context, "ListObjectGenerations")?;
            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.metadata);
            let result = blocking::read(span, immutable::family(), wait, move || {
                authority.generations(
                    ListObjectGenerations::new(
                        ObjectId::new(message.object),
                        page,
                        OwnerId::new(message.caller),
                    ),
                    &context,
                )
            })
            .await?;
            Response::ok(ListObjectGenerationsReply {
                page: buffa::MessageField::some(Kernel(&result).into()),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }

    fn get_receipt(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, GetObjectReceiptRequest>,
    ) -> impl Future<Output = ServiceResult<GetObjectReceiptReply>> {
        let message = request.to_owned_message();
        async move {
            let (_context, span, wait) = self.admit(&ctx, message.context, "GetObjectReceipt")?;
            let authority = Arc::clone(&self.metadata);
            let receipt = blocking::read(span, immutable::family(), wait, move || {
                authority.recorded_receipt(
                    &ObjectId::new(message.object),
                    &CommandId::new(message.command_id),
                )
            })
            .await?;
            Response::ok(GetObjectReceiptReply {
                receipt: receipt
                    .as_ref()
                    .map_or_else(buffa::MessageField::default, |receipt| {
                        buffa::MessageField::some(Kernel(receipt).into())
                    }),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            })
        }
    }
}