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).
//! Capability-specific client for the observation authority (7-O).
//!
//! Two production callers: the control-plane observer (writer) and the
//! projector's observed-source port (reader). Neither needs the other's
//! half, so both share this one client type rather than each reimplementing
//! the wire framing.

use connectrpc::client::{ClientConfig, ClientTransport};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
    error::StateError,
    observation::{
        CollectionId, CollectionKind, ObservationHead, ObservationOrdinal, ObservationRecord,
        ObservationSource, ObservedCollectionListing, RecordObservation,
    },
    receipt::Receipt,
};

use crate::{
    MAX_OBSERVATION_WIRE_MESSAGE_BYTES,
    error::{TransportFallback, from_connect_error},
    trace::bounded_traced_options,
    wire::{DeclaredCall, Kernel},
};

const fn wire_kind(kind: CollectionKind) -> pb::ObservedCollectionKind {
    match kind {
        CollectionKind::Routines => pb::ObservedCollectionKind::OBSERVED_COLLECTION_KIND_ROUTINES,
    }
}

/// Client bound to the observation authority and no other State capability.
pub struct ObservationClient<T> {
    inner: pb::StateObservationServiceClient<T>,
}

impl<T> ObservationClient<T>
where
    T: ClientTransport,
    <T::ResponseBody as connectrpc::http_body::Body>::Error: std::fmt::Display,
{
    /// Builds a client over `transport`.
    #[must_use]
    pub fn new(transport: T, config: ClientConfig) -> Self {
        Self {
            inner: pb::StateObservationServiceClient::new(
                transport,
                config.with_default_max_message_size(MAX_OBSERVATION_WIRE_MESSAGE_BYTES),
            ),
        }
    }

    fn fallback(declared: &DeclaredCall, attempted: usize) -> TransportFallback {
        TransportFallback::new(
            polyc_state::observation::family(),
            MAX_OBSERVATION_WIRE_MESSAGE_BYTES as u64,
            attempted as u64,
            declared.budget,
        )
    }

    /// Records one observation. The observer's relist loop calls this after
    /// every complete relist, whether or not the collection changed.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure. A [`StateError::RevisionConflict`]
    /// means a stale expected ordinal — the caller relists and retries with
    /// the current head.
    pub async fn record(
        &self,
        declared: &DeclaredCall,
        command: &RecordObservation,
    ) -> Result<Receipt, StateError> {
        let draft = command.draft();
        let request = pb::RecordObservationRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            collection_kind: wire_kind(command.collection().kind()).into(),
            collection_namespace: command.collection().namespace().to_owned(),
            resource_version: draft.resource_version().as_str().to_owned(),
            observed_at_nanos: draft.observed_at().as_nanos(),
            observer: draft.observer().as_str().to_owned(),
            row_count: draft.row_count(),
            payload: draft.payload().as_bytes().to_vec(),
            expected_ordinal: match command.metadata().precondition() {
                polyc_state::command::Precondition::JournalHead(position) => position.get(),
                _ => 0,
            },
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .record_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
            .into_owned();
        Kernel::<Receipt>::try_from(reply.receipt.into_option().ok_or_else(|| {
            StateError::Malformed {
                field: "receipt".into(),
                reason: "a successful observation commit returns a receipt".into(),
            }
        })?)
        .map(Kernel::into_inner)
    }

    /// Reads the latest recorded observation for `collection`, or [`None`]
    /// when nothing has ever been recorded for it.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn latest(
        &self,
        declared: &DeclaredCall,
        collection: &CollectionId,
    ) -> Result<Option<ObservationHead>, StateError> {
        let request = pb::ReadLatestObservationRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            collection_kind: wire_kind(collection.kind()).into(),
            collection_namespace: collection.namespace().to_owned(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .read_latest_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
            .into_owned();
        reply.head.into_option().map(head_from_wire).transpose()
    }

    /// Reads the observation `source` recorded at `ordinal`.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure — including
    /// [`StateError::CompactedRange`] for a pruned ordinal and
    /// [`StateError::SourceIncarnationChanged`] for a wiped lineage.
    pub async fn read(
        &self,
        declared: &DeclaredCall,
        source: &ObservationSource,
        ordinal: ObservationOrdinal,
    ) -> Result<ObservationRecord, StateError> {
        let request = pb::ReadObservationRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            source: buffa::MessageField::some(pb::ObservationSource::from(Kernel(source))),
            ordinal: ordinal.as_journal_position().get(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .read_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
            .into_owned();
        record_from_wire(
            reply
                .record
                .into_option()
                .ok_or_else(|| StateError::Malformed {
                    field: "record".into(),
                    reason: "a successful read returns the observation".into(),
                })?,
        )
    }

    /// Lists the collections of `kind` this authority currently holds
    /// recorded state for.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport failure.
    pub async fn collections(
        &self,
        declared: &DeclaredCall,
        kind: CollectionKind,
    ) -> Result<ObservedCollectionListing, StateError> {
        let request = pb::ListObservedCollectionsRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            collection_kind: wire_kind(kind).into(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&request) as usize;
        let reply = self
            .inner
            .list_collections_with_options(request, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
            .into_owned();
        let collections = reply
            .collection_namespaces
            .into_iter()
            .map(|namespace| CollectionId::try_new(kind, namespace))
            .collect::<Result<Vec<_>, _>>()?;
        Ok(ObservedCollectionListing::new(collections, reply.complete))
    }
}

fn head_from_wire(wire: pb::ObservationHeadMessage) -> Result<ObservationHead, StateError> {
    let source =
        Kernel::<ObservationSource>::try_from(wire.source.into_option().ok_or_else(|| {
            StateError::Malformed {
                field: "source".into(),
                reason: "an observation head names its collection and lineage".into(),
            }
        })?)?
        .into_inner();
    Ok(ObservationHead::from_parts(
        source,
        ObservationOrdinal::new(polyc_state::revision::JournalPosition::new(wire.ordinal)),
        polyc_state::digest::ContentDigest::from_bytes(crate::wire::fixed_bytes::<
            { polyc_state::digest::ContentDigest::LEN },
        >(
            "payload_digest", &wire.payload_digest
        )?),
        polyc_state::observation::ResourceVersion::try_new(wire.resource_version)?,
        polyc_state::deadline::MonotonicInstant::from_nanos(wire.observed_at_nanos),
        polyc_state::deadline::MonotonicInstant::from_nanos(wire.recorded_at_nanos),
    ))
}

fn record_from_wire(wire: pb::ObservationRecordMessage) -> Result<ObservationRecord, StateError> {
    let source =
        Kernel::<ObservationSource>::try_from(wire.source.into_option().ok_or_else(|| {
            StateError::Malformed {
                field: "source".into(),
                reason: "an observation record names its collection and lineage".into(),
            }
        })?)?
        .into_inner();
    let payload = polyc_state::observation::ObservationPayload::try_new(wire.payload)?;
    Ok(ObservationRecord::from_parts(
        source,
        ObservationOrdinal::new(polyc_state::revision::JournalPosition::new(wire.ordinal)),
        polyc_state::observation::ResourceVersion::try_new(wire.resource_version)?,
        polyc_state::deadline::MonotonicInstant::from_nanos(wire.observed_at_nanos),
        polyc_state::deadline::MonotonicInstant::from_nanos(wire.recorded_at_nanos),
        polyc_state::observation::ObserverId::try_new(wire.observer)?,
        wire.row_count,
        polyc_state::digest::ContentDigest::from_bytes(crate::wire::fixed_bytes::<
            { polyc_state::digest::ContentDigest::LEN },
        >(
            "payload_digest", &wire.payload_digest
        )?),
        payload,
    ))
}