Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 13 additions & 13 deletions src/in_memory_repo/projection_protocol/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,21 +8,21 @@ use std::future::Future;

use super::InMemoryRepository;
use crate::projection_protocol::{
ProjectionChange, ProjectionChangeCursor, ProjectionChangeKind, ProjectionChangeRead,
ProjectionChangeRetention, ProjectionCheckpoint, ProjectionCommitBatch,
ProjectionCommitOutcome, ProjectionCommitResult, ProjectionEpoch, ProjectionFailure,
ProjectionFailureBatch, ProjectionFailureLocation, ProjectionGeneration, ProjectionInputCursor,
ProjectionInputDisposition, ProjectionInputFingerprint, ProjectionLiveRecordBatch,
ProjectionLiveRecordBatchRequest, ProjectionMutationKind, ProjectionObligationEvidence,
ProjectionObligationEvidenceBatch, ProjectionObligationEvidenceBatchRequest,
ProjectionObservation, ProjectionObservationKind, ProjectionObservationTarget,
ProjectionPartition, ProjectionPartitionRuntimeState, ProjectionPendingRetry,
ProjectionProtocolError, ProjectionProtocolStore, ProjectionQuerySnapshot,
ProjectionQuerySnapshotBatch, ProjectionQuerySnapshotBatchRequest,
change_kind_for_mutation, checked_next, table_model_name, ProjectionChange,
ProjectionChangeCursor, ProjectionChangeKind, ProjectionChangeRead, ProjectionChangeRetention,
ProjectionCheckpoint, ProjectionCommitBatch, ProjectionCommitOutcome, ProjectionCommitResult,
ProjectionEpoch, ProjectionFailure, ProjectionFailureBatch, ProjectionFailureLocation,
ProjectionGeneration, ProjectionInputCursor, ProjectionInputDisposition,
ProjectionInputFingerprint, ProjectionLiveRecordBatch, ProjectionLiveRecordBatchRequest,
ProjectionMutationKind, ProjectionObligationEvidence, ProjectionObligationEvidenceBatch,
ProjectionObligationEvidenceBatchRequest, ProjectionObservation, ProjectionObservationKind,
ProjectionObservationTarget, ProjectionPartition, ProjectionPartitionRuntimeState,
ProjectionPendingRetry, ProjectionProtocolError, ProjectionProtocolStore,
ProjectionQuerySnapshot, ProjectionQuerySnapshotBatch, ProjectionQuerySnapshotBatchRequest,
ProjectionQuerySnapshotRequest, ProjectionRecordExpectation, ProjectionRecordMetadata,
ProjectionRecordScope, ProjectionSource, ProjectorTopologyId, RecordRevision,
RevisionComparison, SameTransactionProjectionBatch, SameTransactionProjectionEvidence,
TrustedProjectionInput, MAX_PROJECTION_POSITION,
TrustedProjectionInput,
};
use crate::read_model::in_memory::{
apply_read_model_write_plan, relational_storage_key, StoredRow,
Expand All @@ -45,7 +45,7 @@ use read_helpers::{
read_projection_query_snapshot_from_state,
};
use state::*;
use util::{checked_next, failure_matches_batch, storage_key_belongs_to_table, table_model_name};
use util::{failure_matches_batch, storage_key_belongs_to_table};

#[cfg(test)]
mod tests;
7 changes: 1 addition & 6 deletions src/in_memory_repo/projection_protocol/store_impl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -207,15 +207,10 @@ impl ProjectionProtocolStore for InMemoryRepository {
mutation.kind,
staged_rows.contains_key(&mutation.mutation.lock_key()),
)?;
let kind = match mutation.kind {
ProjectionMutationKind::Upsert => ProjectionChangeKind::RecordUpsert,
ProjectionMutationKind::Delete => ProjectionChangeKind::RecordDelete,
ProjectionMutationKind::Recreate => ProjectionChangeKind::RecordRecreate,
};
let change = staged_protocol.append_change(
&partition_key,
PendingChange {
kind,
kind: change_kind_for_mutation(mutation.kind),
causation_id: batch.input.causation_id.clone(),
observation_kind: None,
scope: Some(mutation.scope.clone()),
Expand Down
18 changes: 0 additions & 18 deletions src/in_memory_repo/projection_protocol/util.rs
Original file line number Diff line number Diff line change
@@ -1,13 +1,5 @@
use super::*;

pub(super) fn table_model_name(mutation: &TableMutation) -> &str {
match mutation {
TableMutation::UpsertRow(mutation) => &mutation.schema.model_name,
TableMutation::PatchRow(mutation) => &mutation.schema.model_name,
TableMutation::DeleteRow(mutation) => &mutation.schema.model_name,
}
}

pub(super) fn storage_key_belongs_to_table(storage_key: &str, table: &str) -> bool {
let Some(mut fingerprint) = storage_key
.strip_prefix(table)
Expand Down Expand Up @@ -39,16 +31,6 @@ pub(super) fn storage_key_belongs_to_table(storage_key: &str, table: &str) -> bo
true
}

pub(super) fn checked_next(
value: u64,
domain: &'static str,
) -> Result<u64, ProjectionProtocolError> {
if value >= MAX_PROJECTION_POSITION {
return Err(ProjectionProtocolError::PositionOverflow { domain });
}
Ok(value + 1)
}

pub(super) fn failure_matches_batch(
failure: &ProjectionFailure,
batch: &ProjectionFailureBatch,
Expand Down
29 changes: 29 additions & 0 deletions src/projection_protocol/store/backend_helpers.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
use super::{ProjectionChangeKind, ProjectionMutationKind, ProjectionProtocolError};
use crate::projection_protocol::MAX_PROJECTION_POSITION;
use crate::table::TableMutation;

pub(crate) fn checked_next(
value: u64,
domain: &'static str,
) -> Result<u64, ProjectionProtocolError> {
if value >= MAX_PROJECTION_POSITION {
return Err(ProjectionProtocolError::PositionOverflow { domain });
}
Ok(value + 1)
}

pub(crate) fn table_model_name(mutation: &TableMutation) -> &str {
match mutation {
TableMutation::UpsertRow(mutation) => &mutation.schema.model_name,
TableMutation::PatchRow(mutation) => &mutation.schema.model_name,
TableMutation::DeleteRow(mutation) => &mutation.schema.model_name,
}
}

pub(crate) fn change_kind_for_mutation(kind: ProjectionMutationKind) -> ProjectionChangeKind {
match kind {
ProjectionMutationKind::Upsert => ProjectionChangeKind::RecordUpsert,
ProjectionMutationKind::Delete => ProjectionChangeKind::RecordDelete,
ProjectionMutationKind::Recreate => ProjectionChangeKind::RecordRecreate,
}
}
2 changes: 2 additions & 0 deletions src/projection_protocol/store/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ use crate::table::{
RowKey, RowValues, TableMutation, TableSchema, TableStoreError, TableWritePlan,
};

mod backend_helpers;
mod commit;
mod error;
mod helpers;
Expand All @@ -43,6 +44,7 @@ use identity::{
MAX_FAILURE_DETAIL_BYTES, MAX_FAILURE_ID_BYTES, MAX_MESSAGE_ID_BYTES,
};

pub(crate) use backend_helpers::{change_kind_for_mutation, checked_next, table_model_name};
pub(crate) use commit::{
ProjectionCommitBatch, ProjectionFailureBatch, SameTransactionProjectionBatch,
};
Expand Down
26 changes: 0 additions & 26 deletions src/sqlx_repo/projection_protocol/helpers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -72,24 +72,6 @@ pub(super) fn verify_digest(
verify_bytes(actual, &expected, field)
}

pub(super) fn checked_next(
value: u64,
domain: &'static str,
) -> Result<u64, ProjectionProtocolError> {
if value >= MAX_PROJECTION_POSITION {
return Err(ProjectionProtocolError::PositionOverflow { domain });
}
Ok(value + 1)
}

pub(super) fn table_model_name(mutation: &TableMutation) -> &str {
match mutation {
TableMutation::UpsertRow(mutation) => &mutation.schema.model_name,
TableMutation::PatchRow(mutation) => &mutation.schema.model_name,
TableMutation::DeleteRow(mutation) => &mutation.schema.model_name,
}
}

pub(super) async fn physical_row_exists_in_tx<DB>(
tx: &mut Transaction<'_, DB>,
mutation: &TableMutation,
Expand All @@ -109,14 +91,6 @@ where
Ok(row_version_in_tx(tx, schema, key).await?.is_some())
}

pub(super) fn change_kind_for_mutation(kind: ProjectionMutationKind) -> ProjectionChangeKind {
match kind {
ProjectionMutationKind::Upsert => ProjectionChangeKind::RecordUpsert,
ProjectionMutationKind::Delete => ProjectionChangeKind::RecordDelete,
ProjectionMutationKind::Recreate => ProjectionChangeKind::RecordRecreate,
}
}

pub(super) fn decode_change_kind(
value: &str,
) -> Result<ProjectionChangeKind, ProjectionProtocolError> {
Expand Down
30 changes: 15 additions & 15 deletions src/sqlx_repo/projection_protocol/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,21 +18,21 @@ use std::pin::Pin;
use sqlx::{Encode, Executor, IntoArguments, Pool, QueryBuilder, Row, Transaction, Type};

use crate::projection_protocol::{
ProjectionChange, ProjectionChangeCursor, ProjectionChangeKind, ProjectionChangeRead,
ProjectionChangeRetention, ProjectionCheckpoint, ProjectionCommitBatch,
ProjectionCommitOutcome, ProjectionCommitResult, ProjectionEpoch, ProjectionFailure,
ProjectionFailureBatch, ProjectionFailureLocation, ProjectionGeneration, ProjectionInputCursor,
ProjectionInputDisposition, ProjectionInputFingerprint, ProjectionLiveRecordBatch,
ProjectionLiveRecordBatchRequest, ProjectionModelOwnership, ProjectionMutationKind,
ProjectionObligationEvidence, ProjectionObligationEvidenceBatch,
ProjectionObligationEvidenceBatchRequest, ProjectionObservation, ProjectionObservationKind,
ProjectionObservationTarget, ProjectionPartition, ProjectionPartitionRuntimeState,
ProjectionPartitionSnapshot, ProjectionPendingRetry, ProjectionProtocolError,
ProjectionProtocolStore, ProjectionQuerySnapshot, ProjectionQuerySnapshotBatch,
ProjectionQuerySnapshotBatchRequest, ProjectionQuerySnapshotRequest,
ProjectionRecordExpectation, ProjectionRecordMetadata, ProjectionRecordScope, ProjectionSource,
ProjectorTopologyId, RecordRevision, SameTransactionProjectionBatch,
SameTransactionProjectionEvidence, TrustedProjectionInput, MAX_PROJECTION_POSITION,
change_kind_for_mutation, checked_next, table_model_name, ProjectionChange,
ProjectionChangeCursor, ProjectionChangeKind, ProjectionChangeRead, ProjectionChangeRetention,
ProjectionCheckpoint, ProjectionCommitBatch, ProjectionCommitOutcome, ProjectionCommitResult,
ProjectionEpoch, ProjectionFailure, ProjectionFailureBatch, ProjectionFailureLocation,
ProjectionGeneration, ProjectionInputCursor, ProjectionInputDisposition,
ProjectionInputFingerprint, ProjectionLiveRecordBatch, ProjectionLiveRecordBatchRequest,
ProjectionModelOwnership, ProjectionMutationKind, ProjectionObligationEvidence,
ProjectionObligationEvidenceBatch, ProjectionObligationEvidenceBatchRequest,
ProjectionObservation, ProjectionObservationKind, ProjectionObservationTarget,
ProjectionPartition, ProjectionPartitionRuntimeState, ProjectionPartitionSnapshot,
ProjectionPendingRetry, ProjectionProtocolError, ProjectionProtocolStore,
ProjectionQuerySnapshot, ProjectionQuerySnapshotBatch, ProjectionQuerySnapshotBatchRequest,
ProjectionQuerySnapshotRequest, ProjectionRecordExpectation, ProjectionRecordMetadata,
ProjectionRecordScope, ProjectionSource, ProjectorTopologyId, RecordRevision,
SameTransactionProjectionBatch, SameTransactionProjectionEvidence, TrustedProjectionInput,
};
use crate::repository::RepositoryError;
use crate::sqlx_repo::read_model::{
Expand Down
Loading