use super::*;
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ProjectionRecordMetadata {
pub revision: RecordRevision,
pub tombstone: bool,
pub change: ProjectionChangeCursor,
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub(crate) struct ProjectionCheckpointProbe {
pub(crate) topology: ProjectorTopologyId,
pub(crate) partition: ProjectionPartition,
pub(crate) source: ProjectionSource,
pub(crate) epoch: ProjectionEpoch,
pub(crate) generation: ProjectionGeneration,
}
impl ProjectionCheckpointProbe {
pub(crate) fn new(
topology: ProjectorTopologyId,
partition: ProjectionPartition,
source: ProjectionSource,
epoch: ProjectionEpoch,
generation: ProjectionGeneration,
) -> Self {
Self {
topology,
partition,
source,
epoch,
generation,
}
}
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionQuerySnapshotRequest {
pub(crate) schema: Arc<TableSchema>,
pub(crate) key: RowKey,
pub(crate) scope: ProjectionRecordScope,
pub(crate) checkpoint_probes: Vec<ProjectionCheckpointProbe>,
}
impl ProjectionQuerySnapshotRequest {
pub(crate) fn new(
codec: &ProjectionScopeCodec,
partition: Option<&serde_json::Value>,
model: &str,
key: RowKey,
checkpoint_probes: Vec<ProjectionCheckpointProbe>,
) -> Result<Self, ProjectionProtocolError> {
let schema = codec.registered_schema_owned(model).map_err(|error| {
ProjectionProtocolError::InvalidBatch(format!(
"invalid projection query snapshot model: {error}"
))
})?;
let partition = codec.encode_partition(partition).map_err(|error| {
ProjectionProtocolError::InvalidBatch(format!(
"invalid projection query snapshot partition: {error}"
))
})?;
let scope = codec
.encode_row_scope_in_partition(model, partition, &key)
.map_err(|error| {
ProjectionProtocolError::InvalidBatch(format!(
"invalid projection query snapshot key: {error}"
))
})?;
let request = Self {
schema,
key,
scope,
checkpoint_probes,
};
request.validate()?;
Ok(request)
}
pub(crate) fn validate(&self) -> Result<(), ProjectionProtocolError> {
if self.checkpoint_probes.len() > MAX_PROJECTION_QUERY_CHECKPOINT_PROBES {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection query snapshot has {} checkpoint probes; maximum is {}",
self.checkpoint_probes.len(),
MAX_PROJECTION_QUERY_CHECKPOINT_PROBES
)));
}
if self.scope.model() != self.schema.model_name {
return Err(ProjectionProtocolError::ScopeMismatch {
field: "projection query model",
});
}
crate::table::validate_key(&self.schema, &self.key)?;
let mut probes = std::collections::HashSet::new();
for probe in &self.checkpoint_probes {
if &probe.topology != self.scope.topology() {
return Err(ProjectionProtocolError::ScopeMismatch {
field: "projection query checkpoint topology",
});
}
if &probe.partition != self.scope.projection_partition() {
return Err(ProjectionProtocolError::ScopeMismatch {
field: "projection query checkpoint partition",
});
}
if !probes.insert((probe.source.clone(), probe.generation)) {
return Err(ProjectionProtocolError::InvalidBatch(
"projection query snapshot repeats one source/generation checkpoint probe"
.into(),
));
}
}
Ok(())
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct ProjectionCheckpointSnapshot {
pub(crate) probe: ProjectionCheckpointProbe,
pub(crate) checkpoint: Option<ProjectionCheckpoint>,
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionQuerySnapshot {
pub(crate) row: Option<RowValues>,
pub(crate) record: Option<ProjectionRecordMetadata>,
pub(crate) checkpoints: Vec<ProjectionCheckpointSnapshot>,
pub(crate) change_head: Option<ProjectionChangeCursor>,
pub(crate) compacted_through: u64,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct ProjectionPartitionSnapshot {
pub(crate) head: Option<ProjectionChangeCursor>,
pub(crate) compacted_through: u64,
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionQuerySnapshotBatchRequest {
pub(crate) requests: Vec<ProjectionQuerySnapshotRequest>,
}
impl ProjectionQuerySnapshotBatchRequest {
pub(crate) fn new(
requests: Vec<ProjectionQuerySnapshotRequest>,
) -> Result<Self, ProjectionProtocolError> {
let request = Self { requests };
request.validate()?;
Ok(request)
}
pub(crate) fn validate(&self) -> Result<(), ProjectionProtocolError> {
if self.requests.len() > MAX_PROJECTION_QUERY_BATCH_ROWS {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection query snapshot batch has {} rows; maximum is {}",
self.requests.len(),
MAX_PROJECTION_QUERY_BATCH_ROWS
)));
}
let checkpoint_probes = self.requests.iter().try_fold(0usize, |total, request| {
total
.checked_add(request.checkpoint_probes.len())
.ok_or_else(|| {
ProjectionProtocolError::InvalidBatch(
"projection query snapshot batch checkpoint-probe count overflowed".into(),
)
})
})?;
if checkpoint_probes > MAX_PROJECTION_QUERY_BATCH_CHECKPOINT_PROBES {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection query snapshot batch has {checkpoint_probes} aggregate checkpoint probes; maximum is {}",
MAX_PROJECTION_QUERY_BATCH_CHECKPOINT_PROBES
)));
}
for request in &self.requests {
request.validate()?;
}
Ok(())
}
}
#[derive(Clone, Debug, Default, PartialEq)]
pub(crate) struct ProjectionQuerySnapshotBatch {
pub(crate) snapshots: Vec<ProjectionQuerySnapshot>,
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionScopedRowSnapshot {
pub(crate) scope: ProjectionRecordScope,
pub(crate) row: Option<RowValues>,
pub(crate) record: Option<ProjectionRecordMetadata>,
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionExecutionSnapshotBatchRequest {
pub(crate) requests: Vec<ProjectionQuerySnapshotRequest>,
}
impl ProjectionExecutionSnapshotBatchRequest {
pub(crate) fn new(
requests: Vec<ProjectionQuerySnapshotRequest>,
) -> Result<Self, ProjectionProtocolError> {
let request = Self { requests };
request.validate()?;
Ok(request)
}
pub(crate) fn validate(&self) -> Result<(), ProjectionProtocolError> {
if self.requests.len() > MAX_PROJECTION_QUERY_BATCH_ROWS {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection execution snapshot batch has {} scopes; maximum is {}",
self.requests.len(),
MAX_PROJECTION_QUERY_BATCH_ROWS
)));
}
let mut scopes = std::collections::HashSet::new();
for request in &self.requests {
request.validate()?;
if !scopes.insert(request.scope.clone()) {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection execution snapshot batch repeats model `{}` record scope",
request.scope.model()
)));
}
}
Ok(())
}
}
#[derive(Clone, Debug, Default, PartialEq)]
pub(crate) struct ProjectionExecutionSnapshotBatch {
pub(crate) snapshots: Vec<ProjectionScopedRowSnapshot>,
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionGraphIncludeRequest {
pub(crate) relationship: crate::table::RelationshipDef,
pub(crate) target_schema: std::sync::Arc<TableSchema>,
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionGraphSnapshotRequest {
pub(crate) root: ProjectionQuerySnapshotRequest,
pub(crate) includes: std::collections::BTreeMap<String, ProjectionGraphIncludeRequest>,
pub(crate) max_unique_record_scopes: usize,
}
impl ProjectionGraphSnapshotRequest {
pub(crate) fn new(
root: ProjectionQuerySnapshotRequest,
includes: impl IntoIterator<Item = (String, std::sync::Arc<TableSchema>)>,
max_unique_record_scopes: usize,
) -> Result<Self, ProjectionProtocolError> {
root.validate()?;
if max_unique_record_scopes == 0
|| max_unique_record_scopes > MAX_PROJECTION_QUERY_BATCH_ROWS
{
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection graph snapshot record-scope budget is {max_unique_record_scopes}; expected 1..={MAX_PROJECTION_QUERY_BATCH_ROWS}"
)));
}
let includes = includes.into_iter().collect::<Vec<_>>();
let query_scopes = includes.len().checked_add(1).ok_or_else(|| {
ProjectionProtocolError::InvalidBatch(
"projection graph snapshot query-scope count overflowed".into(),
)
})?;
if query_scopes > MAX_PROJECTION_QUERY_BATCH_ROWS {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection graph snapshot has {query_scopes} query scopes; maximum is {}",
MAX_PROJECTION_QUERY_BATCH_ROWS
)));
}
let mut validated = std::collections::BTreeMap::new();
for (include, target_schema) in includes {
let relationship = root
.schema
.relationships
.iter()
.find(|relationship| relationship.field_name == include)
.cloned()
.ok_or_else(|| {
ProjectionProtocolError::InvalidBatch(format!(
"projection graph snapshot model `{}` has no relationship `{include}`",
root.schema.model_name
))
})?;
if matches!(
relationship.kind,
crate::table::RelationshipKind::ManyToMany
) {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection graph relationship `{include}` is many-to-many; project an explicit join read model instead"
)));
}
target_schema.validate()?;
if relationship.target_model != target_schema.model_name {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection graph relationship `{include}` targets `{}` but registered schema is `{}`",
relationship.target_model, target_schema.model_name
)));
}
if validated
.insert(
include,
ProjectionGraphIncludeRequest {
relationship,
target_schema,
},
)
.is_some()
{
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection graph snapshot repeats an include for model `{}`",
root.schema.model_name
)));
}
}
Ok(Self {
root,
includes: validated,
max_unique_record_scopes,
})
}
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionGraphIncludeSnapshot {
pub(crate) relationship: crate::table::RelationshipDef,
pub(crate) target_schema: TableSchema,
pub(crate) rows: Vec<ProjectionScopedRowSnapshot>,
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionGraphSnapshot {
pub(crate) root: ProjectionScopedRowSnapshot,
pub(crate) includes: std::collections::BTreeMap<String, ProjectionGraphIncludeSnapshot>,
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub(crate) struct ProjectionObligationEvidenceRequest {
pub(crate) causation_id: String,
pub(crate) scope: ProjectionRecordScope,
pub(crate) kind: ProjectionObservationKind,
}
impl ProjectionObligationEvidenceRequest {
pub(crate) fn new(
causation_id: impl Into<String>,
scope: ProjectionRecordScope,
kind: ProjectionObservationKind,
) -> Result<Self, ProjectionProtocolError> {
let request = Self {
causation_id: bounded_opaque(
"projection causation ID",
causation_id,
MAX_CAUSATION_ID_BYTES,
)?,
scope,
kind,
};
request.validate()?;
Ok(request)
}
pub(crate) fn validate(&self) -> Result<(), ProjectionProtocolError> {
bounded_opaque(
"projection causation ID",
self.causation_id.clone(),
MAX_CAUSATION_ID_BYTES,
)?;
Ok(())
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) enum ProjectionObligationEvidence {
Pending,
Observed(ProjectionObservation),
TerminalFailure(ProjectionFailure),
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct ProjectionObligationEvidenceBatchRequest {
pub(crate) requests: Vec<ProjectionObligationEvidenceRequest>,
}
impl ProjectionObligationEvidenceBatchRequest {
pub(crate) fn new(
requests: Vec<ProjectionObligationEvidenceRequest>,
) -> Result<Self, ProjectionProtocolError> {
let request = Self { requests };
request.validate()?;
Ok(request)
}
pub(crate) fn validate(&self) -> Result<(), ProjectionProtocolError> {
if self.requests.len() > MAX_PROJECTION_EVIDENCE_BATCH_ITEMS {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection obligation evidence batch has {} probes; maximum is {}",
self.requests.len(),
MAX_PROJECTION_EVIDENCE_BATCH_ITEMS
)));
}
let mut exact = std::collections::HashSet::new();
for request in &self.requests {
request.validate()?;
if !exact.insert((request.causation_id.as_str(), &request.scope, request.kind)) {
return Err(ProjectionProtocolError::InvalidBatch(
"projection obligation evidence batch repeats an exact probe".into(),
));
}
}
Ok(())
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct ProjectionObligationEvidenceBatch {
pub(crate) evidence: Vec<ProjectionObligationEvidence>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct ProjectionCausationEvidenceRequest {
pub(crate) causation_id: String,
pub(crate) topologies: Vec<ProjectorTopologyId>,
}
impl ProjectionCausationEvidenceRequest {
pub(crate) fn new(
causation_id: impl Into<String>,
topologies: Vec<ProjectorTopologyId>,
) -> Result<Self, ProjectionProtocolError> {
let request = Self {
causation_id: bounded_opaque(
"projection causation ID",
causation_id,
MAX_CAUSATION_ID_BYTES,
)?,
topologies,
};
request.validate()?;
Ok(request)
}
pub(crate) fn validate(&self) -> Result<(), ProjectionProtocolError> {
bounded_opaque(
"projection causation ID",
self.causation_id.clone(),
MAX_CAUSATION_ID_BYTES,
)?;
if self.topologies.is_empty() || self.topologies.len() > MAX_PROJECTION_EVIDENCE_BATCH_ITEMS
{
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection causation evidence has {} topology filters; expected 1..={}",
self.topologies.len(),
MAX_PROJECTION_EVIDENCE_BATCH_ITEMS
)));
}
let mut exact = std::collections::HashSet::new();
if self
.topologies
.iter()
.any(|topology| !exact.insert(topology))
{
return Err(ProjectionProtocolError::InvalidBatch(
"projection causation evidence repeats an exact topology filter".into(),
));
}
Ok(())
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct ProjectionCausationEvidenceBatch {
pub(crate) observations: Vec<ProjectionObservation>,
pub(crate) terminal_failure_topologies: Vec<ProjectorTopologyId>,
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionLiveRecordRequest {
pub(crate) schema: Arc<TableSchema>,
pub(crate) topology: ProjectorTopologyId,
pub(crate) key: RowKey,
pub(crate) canonical_key_bytes: Vec<u8>,
pub(crate) canonical_key_hash: [u8; 32],
}
impl ProjectionLiveRecordRequest {
pub(crate) fn new(
codec: &ProjectionScopeCodec,
model: &str,
key: RowKey,
) -> Result<Self, ProjectionProtocolError> {
let schema = codec.registered_schema_owned(model).map_err(|error| {
ProjectionProtocolError::InvalidBatch(format!(
"invalid projection live-record model: {error}"
))
})?;
crate::table::validate_key(&schema, &key)?;
let canonical_key_bytes =
codec
.encode_unpartitioned_row_key(model, &key)
.map_err(|error| {
ProjectionProtocolError::InvalidBatch(format!(
"invalid projection live-record key: {error}"
))
})?;
let canonical_key_hash = ProjectionRecordScope::key_digest_for(&canonical_key_bytes);
Ok(Self {
schema,
topology: codec.topology().clone(),
key,
canonical_key_bytes,
canonical_key_hash,
})
}
pub(crate) fn model(&self) -> &str {
&self.schema.model_name
}
pub(crate) fn validate(&self) -> Result<(), ProjectionProtocolError> {
crate::table::validate_key(&self.schema, &self.key)?;
if self.canonical_key_bytes.is_empty()
|| self.canonical_key_bytes.len() > super::MAX_PROJECTION_RECORD_KEY_BYTES
{
return Err(ProjectionProtocolError::InvalidBatch(
"projection live-record canonical key is empty or oversized".into(),
));
}
if ProjectionRecordScope::key_digest_for(&self.canonical_key_bytes)
!= self.canonical_key_hash
{
return Err(ProjectionProtocolError::InvalidBatch(
"projection live-record canonical key bytes and digest disagree".into(),
));
}
Ok(())
}
}
#[derive(Clone, Debug, PartialEq)]
pub(crate) struct ProjectionLiveRecordBatchRequest {
pub(crate) requests: Vec<ProjectionLiveRecordRequest>,
}
impl ProjectionLiveRecordBatchRequest {
pub(crate) fn new(
requests: Vec<ProjectionLiveRecordRequest>,
) -> Result<Self, ProjectionProtocolError> {
let request = Self { requests };
request.validate()?;
Ok(request)
}
pub(crate) fn validate(&self) -> Result<(), ProjectionProtocolError> {
if self.requests.len() > MAX_PROJECTION_EVIDENCE_BATCH_ITEMS {
return Err(ProjectionProtocolError::InvalidBatch(format!(
"projection live-record batch has {} rows; maximum is {}",
self.requests.len(),
MAX_PROJECTION_EVIDENCE_BATCH_ITEMS
)));
}
let mut identities = std::collections::HashSet::new();
for request in &self.requests {
request.validate()?;
if !identities.insert((
&request.topology,
request.model(),
request.canonical_key_hash,
)) {
return Err(ProjectionProtocolError::InvalidBatch(
"projection live-record batch repeats a topology/model/key identity".into(),
));
}
}
Ok(())
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct ProjectionLiveRecordBatch {
pub(crate) records: Vec<Option<ProjectionRecordMetadata>>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct ProjectionPendingRetry {
pub(crate) failure_id: String,
pub(crate) input: ProjectionInputCursor,
pub(crate) input_fingerprint: ProjectionInputFingerprint,
pub(crate) message_id: String,
pub(crate) causation_id: String,
pub(crate) failed_generation: ProjectionGeneration,
pub(crate) gap_free: bool,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct ProjectionPartitionRuntimeState {
pub(crate) active_generation: ProjectionGeneration,
pub(crate) stopped_failure_id: Option<String>,
pub(crate) pending_retry: Option<ProjectionPendingRetry>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ProjectionObservation {
pub causation_id: String,
pub kind: ProjectionObservationKind,
pub revision: Option<RecordRevision>,
pub scope: ProjectionRecordScope,
pub change: ProjectionChangeCursor,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ProjectionFailure {
pub failure_id: String,
pub input: ProjectionInputCursor,
pub input_fingerprint: ProjectionInputFingerprint,
pub message_id: String,
pub causation_id: String,
pub generation: ProjectionGeneration,
pub gap_free: bool,
pub failure_code: String,
pub failure_bytes: Vec<u8>,
pub failure_digest: [u8; 32],
pub change: ProjectionChangeCursor,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct ProjectionFailureLocation {
pub(crate) topology: ProjectorTopologyId,
pub(crate) partition: ProjectionPartition,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub enum ProjectionChangeKind {
Checkpoint,
RecordUpsert,
RecordDelete,
RecordRecreate,
Observation,
Failure,
}
impl ProjectionChangeKind {
pub(crate) fn as_storage_str(self) -> &'static str {
match self {
Self::Checkpoint => "checkpoint",
Self::RecordUpsert => "record_upsert",
Self::RecordDelete => "record_delete",
Self::RecordRecreate => "record_recreate",
Self::Observation => "observation",
Self::Failure => "failure",
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ProjectionChange {
pub cursor: ProjectionChangeCursor,
pub kind: ProjectionChangeKind,
pub causation_id: String,
pub observation_kind: Option<ProjectionObservationKind>,
pub scope: Option<ProjectionRecordScope>,
pub revision: Option<RecordRevision>,
pub failure_id: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ProjectionCommitResult {
pub outcome: ProjectionCommitOutcome,
pub checkpoint: Option<ProjectionCheckpoint>,
pub records: Vec<ProjectionRecordMetadata>,
pub changes: Vec<ProjectionChange>,
}
impl ProjectionCommitResult {
pub(crate) fn not_applied(
outcome: ProjectionCommitOutcome,
checkpoint: Option<ProjectionCheckpoint>,
) -> Self {
debug_assert!(outcome != ProjectionCommitOutcome::Applied);
Self {
outcome,
checkpoint,
records: Vec::new(),
changes: Vec::new(),
}
}
}