use std::time::Duration;
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
digest::ContentDigest,
error::StateError,
page::Page,
query_audit::{
AuditIntent, ErrorClass, ProjectionPin, QueryAudit, QueryCompletion, QueryOutcome,
RecordedCompletionInput, RequesterId, RowCount, SourcePin, SourceSnapshot, Truncation,
},
revision::{JournalPosition, Revision},
};
use crate::wire::{
DeclaredCall, Kernel, completeness, fixed_bytes, known, malformed, nanos_from_duration,
required,
};
impl From<Kernel<(&DeclaredCall, &RecordedCompletionInput)>> for pb::CompleteQueryAuditRequest {
fn from(value: Kernel<(&DeclaredCall, &RecordedCompletionInput)>) -> Self {
let Kernel((declared, input)) = value;
let metadata = input.metadata();
Self {
context: buffa::MessageField::some(Kernel(declared).into()),
query: input.query().as_str().to_owned(),
namespace: input.namespace().as_str().to_owned(),
intent_source: buffa::MessageField::some(pb::QuerySourceSnapshot::from(Kernel(
input.intent_source(),
))),
completion: buffa::MessageField::some(pb::QueryCompletion::from(Kernel(
input.completion(),
))),
digest: metadata.digest().as_bytes().to_vec(),
purpose: metadata.envelope().purpose().as_str().to_owned(),
command_audience: metadata.envelope().audience().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
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 pb::QueryOutcome {
outcome,
__buffa_unknown_fields: _,
} = value;
let outcome = match outcome {
Some(Outcome::Succeeded(succeeded)) => {
let pb::QuerySucceeded {
__buffa_unknown_fields: _,
} = *succeeded;
QueryOutcome::Succeeded
}
Some(Outcome::Failed(failed)) => {
let pb::QueryFailed {
error_class: declared_error_class,
__buffa_unknown_fields: _,
} = *failed;
QueryOutcome::Failed(error_class("error_class", declared_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;
let pb::QueryTruncation {
kind,
__buffa_unknown_fields: _,
} = value;
Ok(Self(match kind {
Some(Kind::Complete(complete)) => {
let pb::QueryCompleteResult {
__buffa_unknown_fields: _,
} = *complete;
Truncation::Complete
}
Some(Kind::Truncated(truncated)) => {
let pb::QueryTruncatedResult {
limit,
__buffa_unknown_fields: _,
} = *truncated;
Truncation::TruncatedAt(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 {
manifest: buffa::MessageField::some(pb::ProjectionManifest::from(Kernel(
value.0.manifest(),
))),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::QueryProjectionPin> for Kernel<ProjectionPin> {
type Error = StateError;
fn try_from(value: pb::QueryProjectionPin) -> Result<Self, Self::Error> {
let pb::QueryProjectionPin {
manifest,
__buffa_unknown_fields: _,
} = value;
let manifest = Kernel::<polyc_state::projection::ProjectionManifest>::try_from(required(
"manifest",
"a projected source pin carries its complete descriptor",
manifest,
)?)?
.into_inner();
Ok(Self(ProjectionPin::new(manifest)))
}
}
impl From<Kernel<&SourcePin>> for pb::QuerySourcePin {
fn from(value: Kernel<&SourcePin>) -> Self {
use pb::__buffa::oneof::query_source_pin::Pin;
let pin = match value.0 {
SourcePin::Projected(projected) => {
Pin::from(pb::QueryProjectionPin::from(Kernel(projected)))
}
SourcePin::Journal(anchor) => Pin::from(pb::QueryJournalPin {
anchor: buffa::MessageField::some(pb::JournalAnchor::from(Kernel(anchor))),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
SourcePin::Authoritative(revision) => Pin::from(pb::QueryAuthoritativePin {
revision: revision.get(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}),
};
Self {
pin: Some(pin),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::QuerySourcePin> for Kernel<SourcePin> {
type Error = StateError;
fn try_from(value: pb::QuerySourcePin) -> Result<Self, Self::Error> {
use pb::__buffa::oneof::query_source_pin::Pin;
let pb::QuerySourcePin {
pin,
__buffa_unknown_fields: _,
} = value;
let pin = match pin {
Some(Pin::Projected(value)) => {
SourcePin::Projected(Kernel::<ProjectionPin>::try_from(*value)?.into_inner())
}
Some(Pin::Journal(value)) => {
let pb::QueryJournalPin {
anchor,
__buffa_unknown_fields: _,
} = *value;
SourcePin::Journal(
Kernel::<polyc_state::journal::JournalAnchor>::try_from(required(
"anchor",
"a journal source pin carries its exact anchor",
anchor,
)?)?
.into_inner(),
)
}
Some(Pin::Authoritative(value)) => {
let pb::QueryAuthoritativePin {
revision,
__buffa_unknown_fields: _,
} = *value;
SourcePin::Authoritative(Revision::new(revision))
}
None => return Err(malformed("source_pin", "a source pin names one known kind")),
};
Ok(Self(pin))
}
}
impl From<Kernel<&SourceSnapshot>> for pb::QuerySourceSnapshot {
fn from(value: Kernel<&SourceSnapshot>) -> Self {
Self {
pins: value
.0
.pins()
.iter()
.map(|pin| pb::QuerySourcePin::from(Kernel(pin)))
.collect(),
__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> {
let pb::QuerySourceSnapshot {
pins,
__buffa_unknown_fields: _,
} = value;
let pins = pins
.into_iter()
.map(|pin| Kernel::<SourcePin>::try_from(pin).map(Kernel::into_inner))
.collect::<Result<Vec<_>, _>>()?;
Ok(Self(SourceSnapshot::from_canonical(pins)?))
}
}
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 pb::QueryCompletion {
outcome,
duration_nanos,
rows,
truncation,
source,
__buffa_unknown_fields: _,
} = value;
let outcome = Kernel::<QueryOutcome>::try_from(required(
"outcome",
"a completion names how the query ended",
outcome,
)?)?
.into_inner();
let truncation = Kernel::<Truncation>::try_from(required(
"truncation",
"a completion says whether its result was complete",
truncation,
)?)?
.into_inner();
let source = Kernel::<SourceSnapshot>::try_from(required(
"source",
"a completion names what it read",
source,
)?)?
.into_inner();
Ok(Self(QueryCompletion::new(
outcome,
Duration::from_nanos(duration_nanos),
RowCount::new(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(),
namespace: value.0.namespace().as_str().to_owned(),
requester: value.0.requester().as_str().to_owned(),
shape: value.0.shape().as_bytes().to_vec(),
source: buffa::MessageField::some(pb::QuerySourceSnapshot::from(Kernel(
value.0.source(),
))),
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> {
let pb::QueryAuditIntent {
query,
namespace,
requester,
shape,
source,
recorded_at_nanos,
position,
__buffa_unknown_fields: _,
} = value;
Ok(Self(AuditIntent::new(
polyc_state::query_audit::QueryId::new(query),
polyc_state::id::NamespaceId::new(namespace),
RequesterId::new(requester),
ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>("shape", &shape)?),
Kernel::<SourceSnapshot>::try_from(required(
"source",
"an audit intent carries its exact source vector",
source,
)?)?
.into_inner(),
polyc_state::deadline::MonotonicInstant::from_nanos(recorded_at_nanos),
JournalPosition::new(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 pb::QueryAuditRecord {
intent,
completion,
__buffa_unknown_fields: _,
} = value;
let intent = Kernel::<AuditIntent>::try_from(required(
"intent",
"an audit record carries its intent",
intent,
)?)?
.into_inner();
let completion = completion
.into_option()
.map(|value| Kernel::<QueryCompletion>::try_from(value).map(Kernel::into_inner))
.transpose()?;
let audit = QueryAudit::new(intent, completion);
audit.validate_response_semantics()?;
Ok(Self(audit))
}
}
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 pb::QueryAuditPage {
records,
next,
completeness: declared_completeness,
consistency,
__buffa_unknown_fields: _,
} = value;
let records = records
.into_iter()
.map(|record| Kernel::<QueryAudit>::try_from(record).map(Kernel::into_inner))
.collect::<Result<Vec<_>, _>>()?;
let next = next
.into_option()
.map(|cursor| Kernel::from(cursor).into_inner());
let completeness = completeness("completeness", declared_completeness)?;
let consistency = crate::wire::consistency("consistency", 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::{
id::PartitionId,
immutable::{Generation, ObjectDescriptor},
journal::JournalAnchor,
projection::{
FamilyId, ProjectionGeneration, ProjectionKey, ProjectionManifest, PublisherFence,
},
query_audit::{MAX_SOURCE_PINS, cases},
revision::{JournalPosition, JournalSource, PartitionIncarnation},
};
fn mixed_source() -> SourceSnapshot {
let projected = cases::succeeded(0).source().pins()[0].clone();
SourceSnapshot::try_new(vec![
projected,
SourcePin::Journal(JournalAnchor::new(
JournalSource::new(
PartitionId::new("conv-journal"),
PartitionIncarnation::from_bytes([0xa5; PartitionIncarnation::LEN]),
),
JournalPosition::new(41),
)),
SourcePin::Authoritative(Revision::new(73)),
])
.expect("mixed source is canonical")
}
fn projected_manifest(ordinal: u32) -> ProjectionManifest {
let partition = PartitionId::new(format!("conv-wire-{ordinal}"));
let source = JournalSource::new(
partition.clone(),
PartitionIncarnation::from_bytes(
[u8::try_from(ordinal).expect("wire ordinal fits"); PartitionIncarnation::LEN],
),
);
let key = ProjectionKey::new(FamilyId::new(format!("wire.family-{ordinal}")), partition);
let base_source = polyc_state::projection::cases::source(1);
let base_fence = PublisherFence::new(
polyc_state::projection::cases::key(),
base_source.incarnation(),
1,
);
let base =
polyc_state::projection::cases::transport_record(&base_source, 1, 1, &base_fence);
let generation = ProjectionGeneration::new(u64::from(ordinal) + 1);
let descriptor = ObjectDescriptor::new(
key.object(),
Generation::new(generation.get()),
base.object_descriptor().digest(),
base.object_descriptor().owner().clone(),
base.object_descriptor().classification(),
base.object_descriptor().retention(),
base.object_descriptor().byte_len(),
base.object_descriptor().content_reference().clone(),
);
ProjectionManifest::new(
key.clone(),
generation,
polyc_state::projection::cases::checkpoint(&source, u64::from(ordinal) + 1),
1,
1,
descriptor,
base.artifact_object().clone(),
base.publisher().clone(),
PublisherFence::new(key, source.incarnation(), 1),
)
}
fn maximum_projected_source() -> SourceSnapshot {
SourceSnapshot::try_new(
(0..MAX_SOURCE_PINS)
.map(|ordinal| {
SourcePin::Projected(ProjectionPin::new(projected_manifest(ordinal)))
})
.collect(),
)
.expect("maximum projected source is canonical")
}
#[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::namespace(),
cases::requester(),
cases::digest_of(b"shape"),
cases::succeeded(0).source().clone(),
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());
}
#[test]
fn every_source_variant_round_trips_complete_exact_evidence() {
let source = mixed_source();
let wire = pb::QuerySourceSnapshot::from(Kernel(&source));
let back = Kernel::<SourceSnapshot>::try_from(wire)
.expect("decode complete source")
.into_inner();
assert_eq!(back, source);
assert_eq!(back.canonical_bytes(), source.canonical_bytes());
}
#[test]
fn unsorted_and_duplicate_wire_source_identities_fail_closed() {
let source = mixed_source();
let mut unsorted = pb::QuerySourceSnapshot::from(Kernel(&source));
unsorted.pins.reverse();
assert!(matches!(
Kernel::<SourceSnapshot>::try_from(unsorted),
Err(StateError::Malformed { .. })
));
let mut duplicate = pb::QuerySourceSnapshot::from(Kernel(&source));
duplicate.pins.push(duplicate.pins[0].clone());
assert!(matches!(
Kernel::<SourceSnapshot>::try_from(duplicate),
Err(StateError::Malformed { .. })
));
}
#[test]
fn maximum_complete_descriptors_fit_both_begin_and_complete_transport_requests() {
let source = maximum_projected_source();
let completion = QueryCompletion::new(
QueryOutcome::Succeeded,
Duration::from_secs(1),
RowCount::new(1),
Truncation::Complete,
source.clone(),
);
let begin = pb::BeginQueryAuditRequest {
context: buffa::MessageField::default(),
query: "q-max".to_owned(),
namespace: "tenant-max".to_owned(),
requester: "requester-max".to_owned(),
shape: vec![0; ContentDigest::LEN],
source: buffa::MessageField::some(pb::QuerySourceSnapshot::from(Kernel(&source))),
digest: vec![0; ContentDigest::LEN],
purpose: "query".to_owned(),
command_audience: "state".to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let complete = pb::CompleteQueryAuditRequest {
context: buffa::MessageField::default(),
query: "q-max".to_owned(),
namespace: "tenant-max".to_owned(),
completion: buffa::MessageField::some(pb::QueryCompletion::from(Kernel(&completion))),
intent_source: buffa::MessageField::some(pb::QuerySourceSnapshot::from(Kernel(
&source,
))),
digest: vec![0; ContentDigest::LEN],
purpose: "query".to_owned(),
command_audience: "state".to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
};
let begin_len = buffa::Message::encoded_len(&begin) as usize;
let complete_len = buffa::Message::encoded_len(&complete) as usize;
assert!(
begin_len <= crate::MAX_QUERY_AUDIT_WIRE_MESSAGE_BYTES,
"begin={begin_len}"
);
assert!(
complete_len <= crate::MAX_QUERY_AUDIT_WIRE_MESSAGE_BYTES,
"complete={complete_len}"
);
assert!(
complete_len > begin_len,
"completion carries both source vectors"
);
}
}