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).
//! Explicit projection catalog encoding.

use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
    command::{
        CommandEnvelope, CommandMetadata, CommandScope, FencingToken, Precondition, ResourceBounds,
    },
    digest::ContentDigest,
    error::StateError,
    id::{AggregateId, Audience, CommandId, NamespaceId, PartitionId, ProtocolVersion, Purpose},
    immutable::ContentReference,
    projection::{
        self, FamilyId, ProjectionGeneration, ProjectionKey, ProjectionManifest, PublisherFence,
        PublisherId,
        artifact::{ExactObjectRef, ObjectNamespace},
    },
    revision::PartitionIncarnation,
};

use crate::wire::{Kernel, fixed_bytes, required};

impl From<Kernel<&ProjectionKey>> for pb::ProjectionKey {
    fn from(value: Kernel<&ProjectionKey>) -> Self {
        Self {
            family: value.0.family().as_str().to_owned(),
            source_partition: value.0.source().as_str().to_owned(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

impl TryFrom<pb::ProjectionKey> for Kernel<ProjectionKey> {
    type Error = StateError;

    fn try_from(value: pb::ProjectionKey) -> Result<Self, Self::Error> {
        let pb::ProjectionKey {
            family,
            source_partition,
            __buffa_unknown_fields: _,
        } = value;
        if family.is_empty() || source_partition.is_empty() {
            return Err(StateError::Malformed {
                field: "projection_key".to_owned(),
                reason: "a projection key names a family and source partition".to_owned(),
            });
        }
        Ok(Self(ProjectionKey::new(
            FamilyId::new(family),
            PartitionId::new(source_partition),
        )))
    }
}

impl From<Kernel<&PublisherFence>> for pb::ProjectionPublisherFence {
    fn from(value: Kernel<&PublisherFence>) -> Self {
        Self {
            key: buffa::MessageField::some(pb::ProjectionKey::from(Kernel(value.0.key()))),
            source_incarnation: value.0.incarnation().as_bytes().to_vec(),
            term: value.0.term(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

impl TryFrom<pb::ProjectionPublisherFence> for Kernel<PublisherFence> {
    type Error = StateError;

    fn try_from(value: pb::ProjectionPublisherFence) -> Result<Self, Self::Error> {
        let pb::ProjectionPublisherFence {
            key,
            source_incarnation,
            term,
            __buffa_unknown_fields: _,
        } = value;
        if term == 0 {
            return Err(StateError::Malformed {
                field: "term".to_owned(),
                reason: "a publisher fence has a granted nonzero term".to_owned(),
            });
        }
        Ok(Self(PublisherFence::new(
            Kernel::<ProjectionKey>::try_from(required(
                "key",
                "a publisher fence names one projection key",
                key,
            )?)?
            .into_inner(),
            PartitionIncarnation::from_bytes(fixed_bytes::<{ PartitionIncarnation::LEN }>(
                "source_incarnation",
                &source_incarnation,
            )?),
            term,
        )))
    }
}

impl From<Kernel<&ProjectionManifest>> for pb::ProjectionManifest {
    fn from(value: Kernel<&ProjectionManifest>) -> Self {
        Self {
            key: buffa::MessageField::some(pb::ProjectionKey::from(Kernel(value.0.key()))),
            generation: value.0.generation().get(),
            checkpoint: buffa::MessageField::some(pb::SourceCheckpoint::from(Kernel(
                value.0.checkpoint(),
            ))),
            schema_version: value.0.schema_version(),
            fact_version: value.0.fact_version(),
            object: buffa::MessageField::some(pb::ObjectDescriptor::from(Kernel(
                value.0.object_descriptor(),
            ))),
            publisher: value.0.publisher().as_str().to_owned(),
            fence: buffa::MessageField::some(pb::ProjectionPublisherFence::from(Kernel(
                value.0.fence(),
            ))),
            artifact_object: buffa::MessageField::some(pb::ExactObjectRef::from(Kernel(
                value.0.artifact_object(),
            ))),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

impl From<Kernel<&ExactObjectRef>> for pb::ExactObjectRef {
    fn from(value: Kernel<&ExactObjectRef>) -> Self {
        Self {
            namespace: value.0.namespace().as_str().to_owned(),
            key: value.0.key().as_str().to_owned(),
            backend_generation: value.0.generation(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

impl TryFrom<pb::ExactObjectRef> for Kernel<ExactObjectRef> {
    type Error = StateError;

    fn try_from(value: pb::ExactObjectRef) -> Result<Self, Self::Error> {
        let pb::ExactObjectRef {
            namespace,
            key,
            backend_generation,
            __buffa_unknown_fields: _,
        } = value;
        Ok(Self(ExactObjectRef::try_new(
            ObjectNamespace::try_new(namespace)?,
            ContentReference::try_new(key)?,
            backend_generation,
        )?))
    }
}

impl TryFrom<pb::ProjectionManifest> for Kernel<ProjectionManifest> {
    type Error = StateError;

    fn try_from(value: pb::ProjectionManifest) -> Result<Self, Self::Error> {
        let pb::ProjectionManifest {
            key,
            generation,
            checkpoint,
            schema_version,
            fact_version,
            object,
            publisher,
            fence,
            artifact_object,
            __buffa_unknown_fields: _,
        } = value;
        if generation == 0 || schema_version == 0 || fact_version == 0 || publisher.is_empty() {
            return Err(StateError::Malformed {
                field: "manifest".to_owned(),
                reason: "a manifest names nonzero generation and versions and one publisher"
                    .to_owned(),
            });
        }
        let manifest = ProjectionManifest::new(
            Kernel::try_from(required("key", "a manifest names its projection key", key)?)?
                .into_inner(),
            ProjectionGeneration::new(generation),
            Kernel::try_from(required(
                "checkpoint",
                "a manifest carries its exact source checkpoint",
                checkpoint,
            )?)?
            .into_inner(),
            schema_version,
            fact_version,
            Kernel::try_from(required(
                "object",
                "a manifest carries its immutable descriptor",
                object,
            )?)?
            .into_inner(),
            Kernel::try_from(required(
                "artifact_object",
                "a manifest names where its signed artifact is stored",
                artifact_object,
            )?)?
            .into_inner(),
            PublisherId::new(publisher),
            Kernel::try_from(required(
                "fence",
                "a manifest carries its authority",
                fence,
            )?)?
            .into_inner(),
        );
        manifest
            .validate_structure()
            .map_err(|_| StateError::Malformed {
                field: "manifest".to_owned(),
                reason: "a manifest binds its key, source, generation, descriptor, and fence"
                    .to_owned(),
            })?;
        Ok(Self(manifest))
    }
}

impl From<Kernel<&CommandMetadata>> for pb::ProjectionCommandMetadata {
    fn from(value: Kernel<&CommandMetadata>) -> Self {
        let metadata = value.0;
        Self {
            command_id: metadata.command_id().as_str().to_owned(),
            aggregate: metadata.scope().aggregate().as_str().to_owned(),
            partition: metadata.scope().partition().as_str().to_owned(),
            namespace: metadata.scope().namespace().as_str().to_owned(),
            purpose: metadata.envelope().purpose().as_str().to_owned(),
            command_audience: metadata.envelope().audience().as_str().to_owned(),
            digest: metadata.digest().as_bytes().to_vec(),
            precondition: buffa::MessageField::some(pb::Precondition::from(Kernel(
                metadata.precondition(),
            ))),
            fence: metadata.fence().map(FencingToken::get),
            max_payload_bytes: metadata.envelope().bounds().max_payload_bytes(),
            max_records: metadata.envelope().bounds().max_records(),
            protocol_version: metadata.protocol_version().get(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

impl TryFrom<pb::ProjectionCommandMetadata> for Kernel<CommandMetadata> {
    type Error = StateError;

    fn try_from(value: pb::ProjectionCommandMetadata) -> Result<Self, Self::Error> {
        let pb::ProjectionCommandMetadata {
            command_id,
            aggregate,
            partition,
            namespace,
            purpose,
            command_audience,
            digest,
            precondition,
            fence,
            max_payload_bytes,
            max_records,
            protocol_version,
            __buffa_unknown_fields: _,
        } = value;
        let mut metadata = CommandMetadata::new(
            CommandId::new(command_id),
            projection::family(),
            ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>("digest", &digest)?),
            CommandScope::new(
                AggregateId::new(aggregate),
                PartitionId::new(partition),
                NamespaceId::new(namespace),
            ),
            CommandEnvelope::new(
                Purpose::new(purpose),
                Audience::new(command_audience),
                ResourceBounds::new(max_payload_bytes, max_records),
            ),
        )
        .with_precondition(
            Kernel::<Precondition>::try_from(required(
                "precondition",
                "a projection command names its precondition",
                precondition,
            )?)?
            .into_inner(),
        )
        .with_protocol_version(ProtocolVersion::new(protocol_version));
        if let Some(fence) = fence {
            metadata = metadata.with_fence(FencingToken::new(fence));
        }
        Ok(Self(metadata))
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use polyc_state::projection::cases;

    #[test]
    fn command_metadata_round_trips_every_field() {
        let command = cases::register("wire-register", &cases::source(1));
        let wire = pb::ProjectionCommandMetadata::from(Kernel(command.metadata()));
        let decoded = Kernel::<CommandMetadata>::try_from(wire)
            .expect("decode metadata")
            .into_inner();
        assert_eq!(decoded, *command.metadata());
    }

    #[test]
    fn malformed_incarnation_is_refused() {
        let fence = PublisherFence::new(
            cases::key(),
            PartitionIncarnation::from_bytes([1; PartitionIncarnation::LEN]),
            1,
        );
        let mut wire = pb::ProjectionPublisherFence::from(Kernel(&fence));
        wire.source_incarnation.pop();
        assert!(Kernel::<PublisherFence>::try_from(wire).is_err());
    }

    #[test]
    fn manifest_round_trips_every_field() {
        let source = cases::source(1);
        let fence = PublisherFence::new(cases::key(), source.incarnation(), 1);
        let manifest = cases::transport_record(&source, 1, 7, &fence);
        let wire = pb::ProjectionManifest::from(Kernel(&manifest));
        let decoded = Kernel::<ProjectionManifest>::try_from(wire)
            .expect("decode manifest")
            .into_inner();
        assert_eq!(decoded, manifest);
    }

    #[test]
    fn contradictory_manifest_bindings_are_refused() {
        let source = cases::source(1);
        let fence = PublisherFence::new(cases::key(), source.incarnation(), 1);
        let manifest = cases::transport_record(&source, 1, 7, &fence);

        let mut wrong_partition = pb::ProjectionManifest::from(Kernel(&manifest));
        wrong_partition
            .checkpoint
            .as_option_mut()
            .expect("checkpoint")
            .source
            .as_option_mut()
            .expect("source")
            .partition = "other".to_owned();
        assert!(Kernel::<ProjectionManifest>::try_from(wrong_partition).is_err());

        let mut wrong_generation = pb::ProjectionManifest::from(Kernel(&manifest));
        wrong_generation
            .object
            .as_option_mut()
            .expect("object")
            .generation = 2;
        assert!(Kernel::<ProjectionManifest>::try_from(wrong_generation).is_err());

        let mut wrong_fence = pb::ProjectionManifest::from(Kernel(&manifest));
        wrong_fence
            .fence
            .as_option_mut()
            .expect("fence")
            .key
            .as_option_mut()
            .expect("key")
            .source_partition = "other".to_owned();
        assert!(Kernel::<ProjectionManifest>::try_from(wrong_fence).is_err());
    }
}