use polyc_projection::family::{QUERY_AUDIT_MAX_ID_BYTES, QUERY_AUDIT_MAX_SOURCE_PINS};
use polyc_state::{
feed::SourceEvidence,
query_audit::{ErrorClass, QueryAuditHistoryEntry, QueryOutcome, SourcePin, Truncation},
revision::JournalPosition,
};
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct QueryAuditFacts {
pub intents: Vec<QueryAuditIntentFact>,
pub completions: Vec<QueryAuditCompletionFact>,
pub pins: Vec<QueryAuditSourcePinFact>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QueryAuditIntentFact {
pub position: u64,
pub namespace: String,
pub query_id: String,
pub requester: String,
pub shape_digest: [u8; 32],
pub recorded_at_nanos: u64,
pub source_pin_count: u64,
pub command_id: String,
pub command_digest: [u8; 32],
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QueryAuditCompletionFact {
pub position: u64,
pub namespace: String,
pub query_id: String,
pub intent_position: u64,
pub outcome: &'static str,
pub error_class: Option<&'static str>,
pub duration_nanos: u64,
pub rows: u64,
pub truncation: &'static str,
pub truncated_at: Option<u64>,
pub command_id: String,
pub command_digest: [u8; 32],
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QueryAuditSourcePinFact {
pub position: u64,
pub namespace: String,
pub query_id: String,
pub pin_index: u64,
pub pin_kind: &'static str,
pub family: Option<String>,
pub source_partition: Option<String>,
pub pin_source_incarnation: Option<[u8; 32]>,
pub projection_generation: Option<u64>,
pub evidence_kind: Option<&'static str>,
pub evidence_position: Option<u64>,
pub evidence_journal_position: Option<u64>,
pub schema_version: Option<u64>,
pub fact_version: Option<u64>,
pub artifact_digest: Option<[u8; 32]>,
pub anchor_head: Option<u64>,
pub revision: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum QueryAuditFoldError {
#[error("ordinal {position} does not follow {previous} in the history's own order")]
NonAscendingOrdinal {
position: u64,
previous: u64,
},
#[error("the completion at {position} names an intent at {intent_position}")]
CompletionBeforeItsIntent {
position: u64,
intent_position: u64,
},
#[error("the intent at {position} carries {pins} pins, past the bound of {bound}")]
TooManyPins {
position: u64,
pins: usize,
bound: usize,
},
#[error("{field} at ordinal {position} is {len} bytes, past the bound of {bound}")]
IdentifierTooLong {
position: u64,
field: &'static str,
len: usize,
bound: usize,
},
}
const fn outcome_label(outcome: QueryOutcome) -> &'static str {
match outcome {
QueryOutcome::Succeeded => "succeeded",
QueryOutcome::Failed(_) => "failed",
}
}
const fn error_class_label(class: ErrorClass) -> &'static str {
match class {
ErrorClass::Denied => "denied",
ErrorClass::Deadline => "deadline",
ErrorClass::Cancelled => "cancelled",
ErrorClass::Bounds => "bounds",
ErrorClass::Unavailable => "unavailable",
ErrorClass::Malformed => "malformed",
ErrorClass::Internal => "internal",
}
}
const fn truncation_label(truncation: Truncation) -> (&'static str, Option<u64>) {
match truncation {
Truncation::Complete => ("complete", None),
Truncation::TruncatedAt(limit) => ("truncated", Some(limit)),
}
}
const fn evidence_kind_label(evidence: &SourceEvidence) -> &'static str {
match evidence {
SourceEvidence::Journal(_) => "journal",
SourceEvidence::Versioned(_) => "versioned",
SourceEvidence::QueryAudit(_) => "query_audit",
SourceEvidence::PersonaMemory(_) => "persona_memory",
SourceEvidence::Observed(_) => "observed",
}
}
fn bounded(position: u64, field: &'static str, value: &str) -> Result<String, QueryAuditFoldError> {
if value.len() > QUERY_AUDIT_MAX_ID_BYTES {
return Err(QueryAuditFoldError::IdentifierTooLong {
position,
field,
len: value.len(),
bound: QUERY_AUDIT_MAX_ID_BYTES,
});
}
Ok(value.to_owned())
}
pub fn fold_query_audit_history(
entries: &[QueryAuditHistoryEntry],
) -> Result<QueryAuditFacts, QueryAuditFoldError> {
let mut facts = QueryAuditFacts::default();
let mut previous: Option<u64> = None;
for entry in entries {
let position = entry_ordinal(entry).get();
if let Some(previous) = previous
&& position <= previous
{
return Err(QueryAuditFoldError::NonAscendingOrdinal { position, previous });
}
previous = Some(position);
match entry {
QueryAuditHistoryEntry::Intent {
intent, receipt, ..
} => {
let pins = intent.source().pins();
if pins.len() > QUERY_AUDIT_MAX_SOURCE_PINS {
return Err(QueryAuditFoldError::TooManyPins {
position,
pins: pins.len(),
bound: QUERY_AUDIT_MAX_SOURCE_PINS,
});
}
let namespace = bounded(position, "namespace", intent.namespace().as_str())?;
let query_id = bounded(position, "query_id", intent.query().as_str())?;
let requester = bounded(position, "requester", intent.requester().as_str())?;
let command_id = bounded(position, "command_id", receipt.command_id().as_str())?;
for (index, pin) in pins.iter().enumerate() {
facts.pins.push(pin_fact(
position,
&namespace,
&query_id,
u64::try_from(index).unwrap_or(u64::MAX),
pin,
)?);
}
facts.intents.push(QueryAuditIntentFact {
position,
namespace,
query_id,
requester,
shape_digest: *intent.shape().as_bytes(),
recorded_at_nanos: intent.recorded_at().as_nanos(),
source_pin_count: intent.source().pin_count(),
command_id,
command_digest: *receipt.digest().as_bytes(),
});
}
QueryAuditHistoryEntry::Completion {
query,
namespace,
intent_ordinal,
completion,
receipt,
..
} => {
let intent_position = intent_ordinal.get();
if intent_position >= position {
return Err(QueryAuditFoldError::CompletionBeforeItsIntent {
position,
intent_position,
});
}
let (truncation, truncated_at) = truncation_label(completion.truncation());
facts.completions.push(QueryAuditCompletionFact {
position,
namespace: bounded(position, "namespace", namespace.as_str())?,
query_id: bounded(position, "query_id", query.as_str())?,
intent_position,
outcome: outcome_label(completion.outcome()),
error_class: completion.error_class().map(error_class_label),
duration_nanos: u64::try_from(completion.duration().as_nanos())
.unwrap_or(u64::MAX),
rows: completion.rows().get(),
truncation,
truncated_at,
command_id: bounded(position, "command_id", receipt.command_id().as_str())?,
command_digest: *receipt.digest().as_bytes(),
});
}
}
}
Ok(facts)
}
const fn entry_ordinal(entry: &QueryAuditHistoryEntry) -> JournalPosition {
match entry {
QueryAuditHistoryEntry::Intent { ordinal, .. }
| QueryAuditHistoryEntry::Completion { ordinal, .. } => *ordinal,
}
}
fn pin_fact(
position: u64,
namespace: &str,
query_id: &str,
pin_index: u64,
pin: &SourcePin,
) -> Result<QueryAuditSourcePinFact, QueryAuditFoldError> {
match pin {
SourcePin::Projected(projected) => {
let manifest = projected.manifest();
let evidence = manifest.evidence();
let evidence_journal_position = match evidence {
SourceEvidence::Journal(checkpoint) => Some(checkpoint.journal_position().get()),
SourceEvidence::Versioned(_)
| SourceEvidence::QueryAudit(_)
| SourceEvidence::PersonaMemory(_)
| SourceEvidence::Observed(_) => None,
};
Ok(QueryAuditSourcePinFact {
position,
namespace: namespace.to_owned(),
query_id: query_id.to_owned(),
pin_index,
pin_kind: "projected",
family: Some(bounded(
position,
"family",
manifest.key().family().as_str(),
)?),
source_partition: Some(bounded(
position,
"source_partition",
evidence.source().partition().as_str(),
)?),
pin_source_incarnation: Some(*evidence.source().incarnation().as_bytes()),
projection_generation: Some(manifest.generation().get()),
evidence_kind: Some(evidence_kind_label(evidence)),
evidence_position: Some(evidence.position().get()),
evidence_journal_position,
schema_version: Some(u64::from(manifest.schema_version())),
fact_version: Some(u64::from(manifest.fact_version())),
artifact_digest: Some(*manifest.object_descriptor().digest().as_bytes()),
anchor_head: None,
revision: None,
})
}
SourcePin::Journal(anchor) => Ok(QueryAuditSourcePinFact {
position,
namespace: namespace.to_owned(),
query_id: query_id.to_owned(),
pin_index,
pin_kind: "journal",
family: None,
source_partition: Some(bounded(
position,
"source_partition",
anchor.source().partition().as_str(),
)?),
pin_source_incarnation: Some(*anchor.source().incarnation().as_bytes()),
projection_generation: None,
evidence_kind: None,
evidence_position: None,
evidence_journal_position: None,
schema_version: None,
fact_version: None,
artifact_digest: None,
anchor_head: Some(anchor.head().get()),
revision: None,
}),
SourcePin::Authoritative(revision) => Ok(QueryAuditSourcePinFact {
position,
namespace: namespace.to_owned(),
query_id: query_id.to_owned(),
pin_index,
pin_kind: "authoritative",
family: None,
source_partition: None,
pin_source_incarnation: None,
projection_generation: None,
evidence_kind: None,
evidence_position: None,
evidence_journal_position: None,
schema_version: None,
fact_version: None,
artifact_digest: None,
anchor_head: None,
revision: Some(revision.get()),
}),
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::unwrap_used, clippy::expect_used)]
use super::*;
use polyc_state::{
deadline::MonotonicInstant,
digest::ContentDigest,
feed::SourceCheckpoint,
id::{CommandId, NamespaceId, OwnerId, PartitionId},
immutable::{
Classification, ContentReference, Generation as ObjectGeneration, ObjectDescriptor,
},
journal::{JournalAnchor, JournalAttestation},
projection::{
FamilyId, ProjectionGeneration, ProjectionKey, ProjectionManifest, PublisherFence,
PublisherId,
artifact::{ExactObjectRef, ObjectNamespace},
},
query_audit::ProjectionPin,
query_audit::{
AuditIntent, HistoryReceipt, QueryCompletion, QueryId, RequesterId, RowCount,
SourceSnapshot,
},
revision::{CommitRoot, JournalSource, PartitionIncarnation, Revision},
};
use std::time::Duration;
fn digest(byte: u8) -> ContentDigest {
ContentDigest::from_bytes([byte; 32])
}
fn receipt(name: &str, byte: u8) -> HistoryReceipt {
HistoryReceipt::new(CommandId::new(name), digest(byte))
}
fn snapshot(pins: Vec<SourcePin>) -> SourceSnapshot {
SourceSnapshot::try_new(pins).expect("the pins are within the authority's bound")
}
fn journal_pin(partition: &str, head: u64) -> SourcePin {
SourcePin::Journal(JournalAnchor::new(
JournalSource::new(
PartitionId::new(partition),
PartitionIncarnation::from_bytes([7; 32]),
),
JournalPosition::new(head),
))
}
fn intent_entry(ordinal: u64, query: &str, pins: Vec<SourcePin>) -> QueryAuditHistoryEntry {
QueryAuditHistoryEntry::Intent {
ordinal: JournalPosition::new(ordinal),
intent: AuditIntent::new(
QueryId::new(query),
NamespaceId::new("tenant-a"),
RequesterId::new("requester-1"),
digest(1),
snapshot(pins),
MonotonicInstant::from_nanos(1_000),
JournalPosition::new(ordinal),
),
receipt: receipt("begin-1", 2),
}
}
fn completion_entry(
ordinal: u64,
intent_ordinal: u64,
query: &str,
outcome: QueryOutcome,
truncation: Truncation,
) -> QueryAuditHistoryEntry {
QueryAuditHistoryEntry::Completion {
ordinal: JournalPosition::new(ordinal),
query: QueryId::new(query),
namespace: NamespaceId::new("tenant-a"),
intent_ordinal: JournalPosition::new(intent_ordinal),
completion: QueryCompletion::new(
outcome,
Duration::from_micros(5),
RowCount::new(12),
truncation,
snapshot(vec![]),
),
receipt: receipt("complete-1", 3),
}
}
#[test]
fn an_intent_and_its_completion_are_two_rows_that_join() {
let facts = fold_query_audit_history(&[
intent_entry(1, "q-1", vec![journal_pin("conv-a", 9)]),
completion_entry(2, 1, "q-1", QueryOutcome::Succeeded, Truncation::Complete),
])
.expect("the slice folds");
assert_eq!(facts.intents.len(), 1);
assert_eq!(facts.completions.len(), 1);
assert_eq!(facts.pins.len(), 1);
assert_eq!(facts.intents[0].position, 1);
assert_eq!(facts.intents[0].source_pin_count, 1);
assert_eq!(
facts.completions[0].intent_position,
facts.intents[0].position
);
assert_eq!(facts.pins[0].position, facts.intents[0].position);
assert_eq!(facts.pins[0].pin_kind, "journal");
assert_eq!(facts.pins[0].anchor_head, Some(9));
assert_eq!(facts.pins[0].revision, None);
}
#[test]
fn a_completion_without_its_intent_is_still_a_row() {
let facts = fold_query_audit_history(&[completion_entry(
9,
4,
"q-earlier",
QueryOutcome::Failed(ErrorClass::Deadline),
Truncation::TruncatedAt(500),
)])
.expect("a tail that splits a pair still folds");
assert!(facts.intents.is_empty());
assert_eq!(facts.completions.len(), 1);
assert_eq!(facts.completions[0].intent_position, 4);
assert_eq!(facts.completions[0].outcome, "failed");
assert_eq!(facts.completions[0].error_class, Some("deadline"));
assert_eq!(facts.completions[0].truncation, "truncated");
assert_eq!(facts.completions[0].truncated_at, Some(500));
assert_eq!(facts.completions[0].duration_nanos, 5_000);
}
#[test]
fn an_unmatched_intent_is_a_row_and_nothing_else() {
let facts = fold_query_audit_history(&[intent_entry(3, "q-open", vec![])])
.expect("the slice folds");
assert_eq!(facts.intents.len(), 1);
assert!(
facts.completions.is_empty(),
"an unmatched intent must not be normalised into a completion"
);
}
#[test]
fn a_slice_whose_ordinals_do_not_ascend_refuses() {
let repeated = fold_query_audit_history(&[
intent_entry(4, "q-1", vec![]),
intent_entry(4, "q-2", vec![]),
]);
assert!(matches!(
repeated,
Err(QueryAuditFoldError::NonAscendingOrdinal {
position: 4,
previous: 4
})
));
let descending = fold_query_audit_history(&[
intent_entry(5, "q-1", vec![]),
intent_entry(2, "q-2", vec![]),
]);
assert!(matches!(
descending,
Err(QueryAuditFoldError::NonAscendingOrdinal {
position: 2,
previous: 5
})
));
}
#[test]
fn a_completion_that_points_forward_refuses() {
for intent_ordinal in [6, 7] {
let refused = fold_query_audit_history(&[completion_entry(
6,
intent_ordinal,
"q-1",
QueryOutcome::Succeeded,
Truncation::Complete,
)]);
assert!(
matches!(
refused,
Err(QueryAuditFoldError::CompletionBeforeItsIntent { position: 6, .. })
),
"a completion naming an intent at {intent_ordinal} must refuse"
);
}
}
#[test]
fn an_identifier_past_the_bound_refuses() {
let long = "q".repeat(QUERY_AUDIT_MAX_ID_BYTES + 1);
let refused = fold_query_audit_history(&[intent_entry(1, &long, vec![])]);
assert!(matches!(
refused,
Err(QueryAuditFoldError::IdentifierTooLong {
position: 1,
field: "query_id",
..
})
));
}
#[test]
fn each_pin_kind_fills_only_its_own_columns() {
let facts = fold_query_audit_history(&[intent_entry(
1,
"q-1",
vec![SourcePin::Authoritative(Revision::new(41))],
)])
.expect("the slice folds");
let pin = &facts.pins[0];
assert_eq!(pin.pin_kind, "authoritative");
assert_eq!(pin.revision, Some(41));
assert_eq!(pin.family, None);
assert_eq!(pin.source_partition, None);
assert_eq!(pin.pin_source_incarnation, None);
assert_eq!(pin.projection_generation, None);
assert_eq!(pin.evidence_kind, None);
assert_eq!(pin.evidence_position, None);
assert_eq!(pin.evidence_journal_position, None);
assert_eq!(pin.schema_version, None);
assert_eq!(pin.fact_version, None);
assert_eq!(pin.artifact_digest, None);
assert_eq!(pin.anchor_head, None);
}
#[test]
fn the_pins_of_one_intent_are_indexed_in_the_authoritys_order() {
let pins = vec![journal_pin("conv-a", 1), journal_pin("conv-b", 2)];
let recorded = snapshot(pins.clone());
let facts =
fold_query_audit_history(&[intent_entry(1, "q-1", pins)]).expect("the slice folds");
assert_eq!(facts.pins.len(), 2);
for (index, pin) in recorded.pins().iter().enumerate() {
let SourcePin::Journal(anchor) = pin else {
unreachable!("both pins are journal anchors")
};
assert_eq!(facts.pins[index].pin_index, index as u64);
assert_eq!(facts.pins[index].anchor_head, Some(anchor.head().get()));
}
}
#[test]
fn a_projected_pin_names_its_generation_and_never_its_address() {
let source = JournalSource::new(
PartitionId::new("conv-a"),
PartitionIncarnation::from_bytes([3; PartitionIncarnation::LEN]),
);
let key = ProjectionKey::new(
FamilyId::new(polyc_projection::family::CONVERSATION_CORE),
PartitionId::new("conv-a"),
);
let reference = ExactObjectRef::try_new(
ObjectNamespace::try_new("fleet-artifacts").unwrap(),
ContentReference::try_new("manifests/secret-address.manifest").unwrap(),
11,
)
.unwrap();
let checkpoint = SourceCheckpoint::try_new(
source.clone(),
JournalPosition::new(14),
JournalPosition::new(20),
20,
JournalAttestation::new(
CommitRoot::from_bytes([7; CommitRoot::LEN]),
21,
vec![8; 64],
vec![9; 32],
),
)
.unwrap();
let manifest = ProjectionManifest::new(
key.clone(),
ProjectionGeneration::new(6),
SourceEvidence::Journal(checkpoint),
1,
1,
ObjectDescriptor::new(
key.object(),
ObjectGeneration::new(1),
digest(4),
OwnerId::new("projector"),
Classification::Confidential,
polyc_state::immutable::Retention::UntilReleased,
1,
reference.key().clone(),
),
reference,
PublisherId::new("projector/secret-publisher"),
PublisherFence::new(key, source.incarnation(), 1),
);
let facts = fold_query_audit_history(&[intent_entry(
1,
"q-1",
vec![SourcePin::Projected(ProjectionPin::new(manifest))],
)])
.expect("the slice folds");
let pin = &facts.pins[0];
assert_eq!(pin.pin_kind, "projected");
assert_eq!(
pin.family.as_deref(),
Some(polyc_projection::family::CONVERSATION_CORE)
);
assert_eq!(pin.source_partition.as_deref(), Some("conv-a"));
assert_eq!(pin.pin_source_incarnation, Some([3; 32]));
assert_eq!(pin.projection_generation, Some(6));
assert_eq!(pin.evidence_kind, Some("journal"));
assert_eq!(pin.evidence_position, Some(14));
assert_eq!(pin.evidence_journal_position, Some(20));
assert_eq!(pin.schema_version, Some(1));
assert_eq!(pin.fact_version, Some(1));
assert_eq!(pin.artifact_digest, Some([4; 32]));
assert_eq!(pin.anchor_head, None);
assert_eq!(pin.revision, None);
let rendered = format!("{pin:?}");
for prohibited in [
"fleet-artifacts",
"manifests/secret-address.manifest",
"secret-publisher",
] {
assert!(
!rendered.contains(prohibited),
"{prohibited} reached a published pin row: {rendered}"
);
}
}
#[test]
fn the_family_bounds_pins_exactly_as_the_authority_does() {
assert_eq!(
QUERY_AUDIT_MAX_SOURCE_PINS,
polyc_state::query_audit::MAX_SOURCE_PINS as usize
);
}
#[test]
fn the_vocabulary_this_fold_spells_is_the_authoritys_own() {
for class in [
ErrorClass::Denied,
ErrorClass::Deadline,
ErrorClass::Cancelled,
ErrorClass::Bounds,
ErrorClass::Unavailable,
ErrorClass::Malformed,
ErrorClass::Internal,
] {
assert_eq!(error_class_label(class), class.to_string());
}
assert_eq!(outcome_label(QueryOutcome::Succeeded), "succeeded");
assert_eq!(
outcome_label(QueryOutcome::Failed(ErrorClass::Internal)),
"failed"
);
assert_eq!(
truncation_label(Truncation::Complete),
("complete", None),
"a complete result names no limit"
);
assert_eq!(
truncation_label(Truncation::TruncatedAt(7)),
("truncated", Some(7))
);
}
}