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).
//! Explicit query-audit contract encoding.

use std::time::Duration;

use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
    digest::ContentDigest,
    error::StateError,
    immutable::Generation,
    page::Page,
    query_audit::{
        AuditIntent, ErrorClass, ProjectionFamily, ProjectionPin, QueryAudit, QueryCompletion,
        QueryOutcome, RequesterId, RowCount, SourceSnapshot, Truncation,
    },
    revision::{JournalPosition, Revision},
};

use crate::wire::{
    Kernel, completeness, fixed_bytes, known, malformed, nanos_from_duration, required,
};

impl From<Kernel<QueryOutcome>> for pb::QueryOutcome {
    fn from(value: Kernel<QueryOutcome>) -> Self {
        use pb::__buffa::oneof::query_outcome::Outcome;
        let outcome = match value.0 {
            QueryOutcome::Succeeded => Outcome::from(pb::QuerySucceeded {
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            }),
            QueryOutcome::Failed(class) => Outcome::from(pb::QueryFailed {
                error_class: pb::QueryErrorClass::from(Kernel(class)).into(),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            }),
        };
        Self {
            outcome: Some(outcome),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

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

    fn try_from(value: pb::QueryOutcome) -> Result<Self, Self::Error> {
        use pb::__buffa::oneof::query_outcome::Outcome;
        let outcome = match value.outcome {
            Some(Outcome::Succeeded(_)) => QueryOutcome::Succeeded,
            Some(Outcome::Failed(failed)) => {
                QueryOutcome::Failed(error_class("error_class", failed.error_class)?)
            }
            None => {
                return Err(malformed(
                    "outcome",
                    "a completion names how the query ended",
                ));
            }
        };
        Ok(Self(outcome))
    }
}

impl From<Kernel<ErrorClass>> for pb::QueryErrorClass {
    fn from(value: Kernel<ErrorClass>) -> Self {
        match value.0 {
            ErrorClass::Denied => Self::QUERY_ERROR_CLASS_DENIED,
            ErrorClass::Deadline => Self::QUERY_ERROR_CLASS_DEADLINE,
            ErrorClass::Cancelled => Self::QUERY_ERROR_CLASS_CANCELLED,
            ErrorClass::Bounds => Self::QUERY_ERROR_CLASS_BOUNDS,
            ErrorClass::Unavailable => Self::QUERY_ERROR_CLASS_UNAVAILABLE,
            ErrorClass::Malformed => Self::QUERY_ERROR_CLASS_MALFORMED,
            ErrorClass::Internal => Self::QUERY_ERROR_CLASS_INTERNAL,
        }
    }
}

fn error_class(
    field: &str,
    value: buffa::EnumValue<pb::QueryErrorClass>,
) -> Result<ErrorClass, StateError> {
    match known(field, value)? {
        pb::QueryErrorClass::QUERY_ERROR_CLASS_DENIED => Ok(ErrorClass::Denied),
        pb::QueryErrorClass::QUERY_ERROR_CLASS_DEADLINE => Ok(ErrorClass::Deadline),
        pb::QueryErrorClass::QUERY_ERROR_CLASS_CANCELLED => Ok(ErrorClass::Cancelled),
        pb::QueryErrorClass::QUERY_ERROR_CLASS_BOUNDS => Ok(ErrorClass::Bounds),
        pb::QueryErrorClass::QUERY_ERROR_CLASS_UNAVAILABLE => Ok(ErrorClass::Unavailable),
        pb::QueryErrorClass::QUERY_ERROR_CLASS_MALFORMED => Ok(ErrorClass::Malformed),
        pb::QueryErrorClass::QUERY_ERROR_CLASS_INTERNAL => Ok(ErrorClass::Internal),
        pb::QueryErrorClass::QUERY_ERROR_CLASS_UNSPECIFIED => Err(malformed(
            field,
            "a failed query names one known error class",
        )),
    }
}

impl From<Kernel<Truncation>> for pb::QueryTruncation {
    fn from(value: Kernel<Truncation>) -> Self {
        use pb::__buffa::oneof::query_truncation::Kind;
        let kind = match value.0 {
            Truncation::Complete => Kind::from(pb::QueryCompleteResult {
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            }),
            Truncation::TruncatedAt(limit) => Kind::from(pb::QueryTruncatedResult {
                limit,
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            }),
        };
        Self {
            kind: Some(kind),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

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

    fn try_from(value: pb::QueryTruncation) -> Result<Self, Self::Error> {
        use pb::__buffa::oneof::query_truncation::Kind;
        Ok(Self(match value.kind {
            Some(Kind::Complete(_)) => Truncation::Complete,
            Some(Kind::Truncated(value)) => Truncation::TruncatedAt(value.limit),
            None => {
                return Err(malformed(
                    "truncation",
                    "a completion says whether every matching row was returned",
                ));
            }
        }))
    }
}

impl From<Kernel<&ProjectionPin>> for pb::QueryProjectionPin {
    fn from(value: Kernel<&ProjectionPin>) -> Self {
        Self {
            family: value.0.family().as_str().to_owned(),
            generation: value.0.generation().get(),
            cursor: value.0.cursor().get(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

impl From<pb::QueryProjectionPin> for Kernel<ProjectionPin> {
    fn from(value: pb::QueryProjectionPin) -> Self {
        Self(ProjectionPin::new(
            ProjectionFamily::new(value.family),
            Generation::new(value.generation),
            JournalPosition::new(value.cursor),
        ))
    }
}

impl From<Kernel<&SourceSnapshot>> for pb::QuerySourceSnapshot {
    fn from(value: Kernel<&SourceSnapshot>) -> Self {
        use pb::__buffa::oneof::query_source_snapshot::Source;
        let source = match value.0 {
            SourceSnapshot::Authoritative(revision) => Source::from(pb::AuthoritativeQuerySource {
                revision: revision.get(),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            }),
            SourceSnapshot::Projected(pins) => Source::from(pb::ProjectedQuerySource {
                pins: pins
                    .iter()
                    .map(|pin| pb::QueryProjectionPin::from(Kernel(pin)))
                    .collect(),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            }),
        };
        Self {
            source: Some(source),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

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

    fn try_from(value: pb::QuerySourceSnapshot) -> Result<Self, Self::Error> {
        use pb::__buffa::oneof::query_source_snapshot::Source;
        Ok(Self(match value.source {
            Some(Source::Authoritative(value)) => {
                SourceSnapshot::Authoritative(Revision::new(value.revision))
            }
            Some(Source::Projected(value)) => SourceSnapshot::Projected(
                value
                    .pins
                    .into_iter()
                    .map(|pin| Kernel::<ProjectionPin>::from(pin).into_inner())
                    .collect(),
            ),
            None => return Err(malformed("source", "a completion names what it read")),
        }))
    }
}

impl From<Kernel<&QueryCompletion>> for pb::QueryCompletion {
    fn from(value: Kernel<&QueryCompletion>) -> Self {
        Self {
            outcome: buffa::MessageField::some(pb::QueryOutcome::from(Kernel(value.0.outcome()))),
            duration_nanos: nanos_from_duration(value.0.duration()),
            rows: value.0.rows().get(),
            truncation: buffa::MessageField::some(pb::QueryTruncation::from(Kernel(
                value.0.truncation(),
            ))),
            source: buffa::MessageField::some(pb::QuerySourceSnapshot::from(Kernel(
                value.0.source(),
            ))),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

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

    fn try_from(value: pb::QueryCompletion) -> Result<Self, Self::Error> {
        let outcome = Kernel::<QueryOutcome>::try_from(required(
            "outcome",
            "a completion names how the query ended",
            value.outcome,
        )?)?
        .into_inner();
        let truncation = Kernel::<Truncation>::try_from(required(
            "truncation",
            "a completion says whether its result was complete",
            value.truncation,
        )?)?
        .into_inner();
        let source = Kernel::<SourceSnapshot>::try_from(required(
            "source",
            "a completion names what it read",
            value.source,
        )?)?
        .into_inner();
        Ok(Self(QueryCompletion::new(
            outcome,
            Duration::from_nanos(value.duration_nanos),
            RowCount::new(value.rows),
            truncation,
            source,
        )))
    }
}

impl From<Kernel<&AuditIntent>> for pb::QueryAuditIntent {
    fn from(value: Kernel<&AuditIntent>) -> Self {
        Self {
            query: value.0.query().as_str().to_owned(),
            requester: value.0.requester().as_str().to_owned(),
            shape: value.0.shape().as_bytes().to_vec(),
            recorded_at_nanos: value.0.recorded_at().as_nanos(),
            position: value.0.position().get(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

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

    fn try_from(value: pb::QueryAuditIntent) -> Result<Self, Self::Error> {
        Ok(Self(AuditIntent::new(
            polyc_state::query_audit::QueryId::new(value.query),
            RequesterId::new(value.requester),
            ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
                "shape",
                &value.shape,
            )?),
            polyc_state::deadline::MonotonicInstant::from_nanos(value.recorded_at_nanos),
            JournalPosition::new(value.position),
        )))
    }
}

impl From<Kernel<&QueryAudit>> for pb::QueryAuditRecord {
    fn from(value: Kernel<&QueryAudit>) -> Self {
        Self {
            intent: buffa::MessageField::some(pb::QueryAuditIntent::from(Kernel(value.0.intent()))),
            completion: value.0.completion().map_or_else(
                buffa::MessageField::default,
                |completion| {
                    buffa::MessageField::some(pb::QueryCompletion::from(Kernel(completion)))
                },
            ),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

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

    fn try_from(value: pb::QueryAuditRecord) -> Result<Self, Self::Error> {
        let intent = Kernel::<AuditIntent>::try_from(required(
            "intent",
            "an audit record carries its intent",
            value.intent,
        )?)?
        .into_inner();
        let completion = value
            .completion
            .into_option()
            .map(|value| Kernel::<QueryCompletion>::try_from(value).map(Kernel::into_inner))
            .transpose()?;
        Ok(Self(QueryAudit::new(intent, completion)))
    }
}

impl From<Kernel<&Page<QueryAudit>>> for pb::QueryAuditPage {
    fn from(value: Kernel<&Page<QueryAudit>>) -> Self {
        let page = value.0;
        Self {
            records: page
                .records()
                .iter()
                .map(|record| pb::QueryAuditRecord::from(Kernel(record)))
                .collect(),
            next: page
                .next_cursor()
                .map_or_else(buffa::MessageField::default, |cursor| {
                    buffa::MessageField::some(pb::Cursor::from(Kernel(cursor)))
                }),
            completeness: pb::PageCompleteness::from(Kernel(page.completeness())).into(),
            consistency: pb::Consistency::from(Kernel(page.consistency())).into(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }
}

impl TryFrom<pb::QueryAuditPage> for Kernel<Page<QueryAudit>> {
    type Error = StateError;

    fn try_from(value: pb::QueryAuditPage) -> Result<Self, Self::Error> {
        let records = value
            .records
            .into_iter()
            .map(|record| Kernel::<QueryAudit>::try_from(record).map(Kernel::into_inner))
            .collect::<Result<Vec<_>, _>>()?;
        let next = value
            .next
            .into_option()
            .map(|cursor| Kernel::from(cursor).into_inner());
        let completeness = completeness("completeness", value.completeness)?;
        let consistency = crate::wire::consistency("consistency", value.consistency)?;
        Ok(Self(Page::new(records, next, completeness, consistency)))
    }
}

/// Reads one digest field.
pub(crate) fn digest(field: &str, bytes: &[u8]) -> Result<ContentDigest, StateError> {
    Ok(ContentDigest::from_bytes(fixed_bytes::<
        { ContentDigest::LEN },
    >(field, bytes)?))
}

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

    #[test]
    fn completions_and_audits_round_trip() {
        let completion = cases::succeeded(19);
        let wire = pb::QueryCompletion::from(Kernel(&completion));
        let back = Kernel::<QueryCompletion>::try_from(wire)
            .expect("decode")
            .into_inner();
        assert_eq!(back, completion);

        let audit = QueryAudit::new(
            AuditIntent::new(
                cases::query(1),
                cases::requester(),
                cases::digest_of(b"shape"),
                polyc_state::deadline::MonotonicInstant::from_nanos(7),
                JournalPosition::new(9),
            ),
            Some(cases::failed(ErrorClass::Denied)),
        );
        let wire = pb::QueryAuditRecord::from(Kernel(&audit));
        let back = Kernel::<QueryAudit>::try_from(wire)
            .expect("decode")
            .into_inner();
        assert_eq!(back, audit);
    }

    #[test]
    fn required_variants_and_digest_width_fail_closed() {
        let error = Kernel::<QueryCompletion>::try_from(pb::QueryCompletion::default())
            .expect_err("missing fields must fail");
        assert!(matches!(error, StateError::Malformed { .. }));
        assert!(digest("digest", &[1, 2]).is_err());
    }
}