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)))
}
}
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());
}
}