use crate::NullableKeyFilter;
use crate::changelog::{ChangeId, CommitId};
use crate::common::{LixTimestamp, SharedStr};
use crate::row_payload::TypedRow as WasmTypedRow;
use crate::row_pk::RowPk;
use bytes::Bytes;
use std::sync::Arc;
pub(crate) const TRACKED_STATE_HASH_BYTES: usize = 32;
pub(crate) const COMMIT_STATE_MAX_REPLAY_DEPTH: u16 = 32;
pub(crate) const COMMIT_STATE_MAX_REPLAY_BYTES: u64 = 256 * 1024 * 1024;
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, musli::Encode, musli::Decode)]
pub(crate) struct TrackedStateRootId(#[musli(bytes)] [u8; TRACKED_STATE_HASH_BYTES]);
impl TrackedStateRootId {
pub(crate) fn new(bytes: [u8; TRACKED_STATE_HASH_BYTES]) -> Self {
Self(bytes)
}
pub(crate) fn as_bytes(&self) -> &[u8; TRACKED_STATE_HASH_BYTES] {
&self.0
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub(crate) struct TrackedStateKey {
pub(crate) schema_key: String,
pub(crate) file_id: Option<String>,
pub(crate) row_pk: RowPk,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub(crate) struct TrackedStateKeyRef<'a> {
pub(crate) schema_key: &'a str,
pub(crate) file_id: Option<&'a str>,
pub(crate) row_pk: &'a RowPk,
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct TrackedStateDeltaRef<'a> {
pub(crate) schema_key: &'a str,
pub(crate) file_id: Option<&'a str>,
pub(crate) row_pk: &'a RowPk,
pub(crate) change_id: ChangeId,
pub(crate) commit_id: CommitId,
pub(crate) deleted: bool,
pub(crate) created_at: LixTimestamp,
pub(crate) updated_at: LixTimestamp,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct TrackedStateBaseCoordinate {
pub(crate) base_commit_id: CommitId,
pub(crate) group_index: u32,
pub(crate) row_index: u32,
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct TrackedStateCommitDeltaRef<'a> {
pub(crate) delta: TrackedStateDeltaRef<'a>,
pub(crate) metadata: Option<&'a lix_schema::Jsonb>,
pub(crate) snapshot: Option<&'a [u8]>,
pub(crate) origin_key: Option<&'a str>,
pub(crate) base_coordinate: Option<TrackedStateBaseCoordinate>,
pub(crate) authored: bool,
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct TrackedStateSingleStringReplacementRef<'a> {
pub(crate) schema_key: &'a str,
pub(crate) file_id: Option<&'a str>,
pub(crate) row_pk: &'a str,
pub(crate) commit_id: CommitId,
pub(crate) created_at: LixTimestamp,
pub(crate) updated_at: LixTimestamp,
pub(crate) metadata: Option<&'a lix_schema::Jsonb>,
pub(crate) snapshot: &'a [u8],
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct TrackedStateRootMutationRef<'a> {
pub(crate) delta: TrackedStateDeltaRef<'a>,
pub(crate) require_absence: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct TrackedStateIndexValue {
pub(crate) change_id: ChangeId,
pub(crate) commit_id: CommitId,
pub(crate) deleted: bool,
pub(crate) created_at: LixTimestamp,
pub(crate) updated_at: LixTimestamp,
}
impl TrackedStateIndexValue {
pub(crate) fn created_at(&self) -> LixTimestamp {
self.created_at
}
pub(crate) fn updated_at(&self) -> LixTimestamp {
self.updated_at
}
pub(crate) fn deleted(&self) -> bool {
self.deleted
}
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct TrackedStateIndexValueRef {
pub(crate) change_id: ChangeId,
pub(crate) commit_id: CommitId,
pub(crate) deleted: bool,
pub(crate) created_at: LixTimestamp,
pub(crate) updated_at: LixTimestamp,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct TrackedStateCommitRoot {
pub(crate) commit_id: CommitId,
pub(crate) root_id: TrackedStateRootId,
pub(crate) parent_roots: Vec<TrackedStateCommitRootParent>,
pub(crate) changed_key_count: u64,
pub(crate) row_count_estimate: u64,
pub(crate) tree_height: u32,
pub(crate) complete_state_fence: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct TrackedStateCommitRootParent {
pub(crate) commit_id: CommitId,
pub(crate) root_id: TrackedStateRootId,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct CommitStateReplayDebt {
pub(crate) depth: u16,
pub(crate) rows: u64,
pub(crate) bytes: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct CommitStateMutationPart {
#[musli(bytes)]
pub(crate) first_key: Vec<u8>,
#[musli(bytes)]
pub(crate) last_key: Vec<u8>,
pub(crate) content_digest: [u8; 32],
#[musli(with = crate::storage_codec::option)]
pub(crate) replacement_part: Option<StoredReplacementPart>,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct ColumnarMutationPartSet {
pub(crate) owner_commit_id: [u8; 16],
pub(crate) row_group_set_id: [u8; 16],
pub(crate) manifest_digest: [u8; 32],
pub(crate) schema_key: String,
pub(crate) row_count: u32,
pub(crate) group_row_counts: Vec<u32>,
#[musli(bytes)]
pub(crate) first_key: Vec<u8>,
#[musli(bytes)]
pub(crate) last_key: Vec<u8>,
pub(crate) page_first_keys: Vec<Vec<u8>>,
pub(crate) page_last_keys: Vec<Vec<u8>>,
pub(crate) uniform_created_at: LixTimestamp,
pub(crate) uniform_updated_at: LixTimestamp,
#[musli(with = crate::storage_codec::option)]
pub(crate) origin_key: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct CommitDeltaReplacementScope {
pub(crate) schema_key: String,
#[musli(with = crate::storage_codec::option)]
pub(crate) file_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct CommitDeltaLifecycleSummary {
pub(crate) scope: CommitDeltaReplacementScope,
pub(crate) ordered_identity_digest: [u8; 32],
pub(crate) uniform_created_at: LixTimestamp,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct StoredCommitDeltaReplacementGeneration {
pub(crate) owner_commit_id: [u8; 16],
pub(crate) scope: CommitDeltaReplacementScope,
#[musli(with = crate::storage_codec::option)]
pub(crate) fallback_commit_id: Option<[u8; 16]>,
pub(crate) integrity_digest: [u8; 32],
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct StoredReplacementPartsAuthority {
pub(crate) directory_digest: [u8; 32],
pub(crate) uniform_updated_at: LixTimestamp,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct StoredReplacementPart {
pub(crate) content_digest: [u8; 32],
pub(crate) owner_commit_id: [u8; 16],
pub(crate) first_address: u32,
pub(crate) uniform_created_at: LixTimestamp,
pub(crate) uniform_updated_at: LixTimestamp,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct CurrentStatePartDescriptor {
#[musli(bytes)]
pub(crate) first_key: Vec<u8>,
#[musli(bytes)]
pub(crate) last_key: Vec<u8>,
pub(crate) content_digest: [u8; 32],
pub(crate) source: CurrentStatePartSource,
pub(crate) source_row_offset: u16,
pub(crate) row_count: u16,
pub(crate) fragmented: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, musli::Encode, musli::Decode)]
pub(crate) enum CurrentStatePartSource {
Replacement(ReplacementPartSource),
NativeDataPart,
ColumnarPage(ColumnarPageSource),
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct ReplacementPartSource {
pub(crate) owner_commit_id: [u8; 16],
pub(crate) part_index: u32,
pub(crate) uniform_created_at: LixTimestamp,
pub(crate) uniform_updated_at: LixTimestamp,
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct ColumnarPageSource {
pub(crate) source_id: [u8; 16],
pub(crate) owner_commit_id: [u8; 16],
pub(crate) part_index: u32,
pub(crate) source_page_index: u16,
pub(crate) uniform_created_at: LixTimestamp,
pub(crate) uniform_updated_at: LixTimestamp,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct CurrentStateScopedRangeRoot {
pub(crate) tree: super::scoped_range::ScopedRangeRoot,
#[musli(with = crate::storage_codec::option)]
pub(crate) serving_base_commit_id: Option<CommitId>,
#[musli(with = crate::storage_codec::option)]
pub(crate) serving_base_root_id: Option<[u8; 32]>,
pub(crate) transition_digest: [u8; 32],
}
#[derive(Debug, Clone, Default, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct CommitStateTouchedScopeFilter {
pub(crate) complete: bool,
#[musli(bytes)]
pub(crate) bits: Vec<u8>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct CommitStateMutationInventory {
#[musli(with = crate::storage_codec::option)]
pub(crate) selected_source_commit_id: Option<[u8; 16]>,
pub(crate) member_count: u32,
pub(crate) selection_fingerprint: [u8; 32],
pub(crate) direct_part_row_counts: Vec<u16>,
pub(crate) direct_part_ownership: Vec<Vec<u8>>,
pub(crate) replacement_part_digests: Vec<[u8; 32]>,
#[musli(with = crate::storage_codec::option)]
pub(crate) single_partition: Option<CommitDeltaReplacementScope>,
#[musli(with = crate::storage_codec::option)]
pub(crate) lifecycle_summary: Option<CommitDeltaLifecycleSummary>,
#[musli(with = crate::storage_codec::option)]
pub(crate) replacement_generation: Option<StoredCommitDeltaReplacementGeneration>,
#[musli(with = crate::storage_codec::option)]
pub(crate) replacement_parts: Option<StoredReplacementPartsAuthority>,
#[musli(with = crate::storage_codec::option)]
pub(crate) columnar_parts: Option<ColumnarMutationPartSet>,
#[musli(bytes)]
pub(crate) inline_part: Vec<u8>,
pub(crate) parts: Vec<CommitStateMutationPart>,
}
impl CommitStateMutationInventory {
pub(crate) fn selected_source_commit_id(&self) -> Option<CommitId> {
self.selected_source_commit_id
.map(|bytes| CommitId::new(uuid::Uuid::from_bytes(bytes)))
}
pub(crate) fn part_count(&self) -> usize {
self.columnar_parts
.as_ref()
.map_or(0, |parts| parts.group_row_counts.len())
+ usize::from(!self.inline_part.is_empty())
+ if self.replacement_part_digests.is_empty() {
self.parts.len()
} else {
self.replacement_part_digests.len()
}
}
pub(crate) fn direct_coordinate_owned(
&self,
part_index: usize,
local_row: u16,
) -> Option<bool> {
if let Some(&row_count) = self.direct_part_row_counts.get(part_index) {
if local_row >= row_count {
return None;
}
}
let ownership = self.direct_part_ownership.get(part_index)?;
let ordinal = usize::from(local_row);
Some(*ownership.get(ordinal / 8)? & (1 << (ordinal % 8)) != 0)
}
pub(crate) fn direct_addresses_are_fully_owned(&self) -> bool {
!self.direct_part_row_counts.is_empty()
&& self
.direct_part_row_counts
.iter()
.enumerate()
.all(|(part_index, &row_count)| {
(0..row_count).all(|local_row| {
self.direct_coordinate_owned(part_index, local_row) == Some(true)
})
})
}
pub(crate) fn may_contain_finite_selected_members(&self) -> bool {
self.member_count != 0
&& self.selected_source_commit_id.is_none()
&& !self.direct_addresses_are_fully_owned()
&& self.columnar_parts.is_none()
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, musli::Encode, musli::Decode)]
pub(crate) enum CommitStateIncorporation {
#[default]
None,
Complete(CommitId),
LegacyUnknown,
}
#[derive(Debug, Clone, PartialEq, Eq, musli::Encode, musli::Decode)]
#[musli(packed)]
pub(crate) struct CommitStateManifest {
pub(crate) commit_id: CommitId,
pub(crate) incorporation: CommitStateIncorporation,
pub(crate) change_account_id: String,
pub(crate) replay_debt: CommitStateReplayDebt,
pub(crate) mutations: CommitStateMutationInventory,
pub(crate) touched_scope_filter: CommitStateTouchedScopeFilter,
pub(crate) global_scope: bool,
#[musli(with = crate::storage_codec::option)]
pub(crate) current_state_scoped_ranges: Option<Box<CurrentStateScopedRangeRoot>>,
#[musli(with = crate::storage_codec::option)]
pub(crate) row_pk_index_root_id: Option<TrackedStateRootId>,
#[musli(with = crate::storage_codec::option)]
pub(crate) snapshot_root: Option<Box<TrackedStateCommitRoot>>,
}
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
pub(crate) struct MaterializedTrackedStateRow {
pub(crate) row_pk: RowPk,
pub(crate) schema_key: String,
pub(crate) file_id: Option<String>,
pub(crate) snapshot_content: Option<SharedStr>,
#[serde(skip)]
pub(crate) decoded_snapshot: Option<Arc<WasmTypedRow>>,
pub(crate) metadata: Option<SharedStr>,
pub(crate) deleted: bool,
pub(crate) created_at: String,
pub(crate) updated_at: String,
pub(crate) change_id: ChangeId,
pub(crate) commit_id: CommitId,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
pub(crate) struct TrackedStateFilter {
#[serde(default)]
pub(crate) schema_keys: Vec<String>,
#[serde(default)]
pub(crate) row_pks: Vec<RowPk>,
#[serde(default)]
pub(crate) row_pk_lower: Option<RowPkRangeBound>,
#[serde(default)]
pub(crate) row_pk_upper: Option<RowPkRangeBound>,
#[serde(default)]
pub(crate) file_ids: Vec<NullableKeyFilter<String>>,
#[serde(default)]
pub(crate) include_tombstones: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub(crate) struct RowPkRangeBound {
pub(crate) row_pk: RowPk,
pub(crate) inclusive: bool,
}
impl TrackedStateFilter {
pub(crate) fn matches_row_pk(&self, row_pk: &RowPk) -> bool {
(self.row_pks.is_empty() || self.row_pks.contains(row_pk))
&& row_pk_satisfies_bounds(
row_pk,
self.row_pk_lower.as_ref(),
self.row_pk_upper.as_ref(),
)
}
}
pub(crate) fn row_pk_satisfies_bounds(
row_pk: &RowPk,
lower: Option<&RowPkRangeBound>,
upper: Option<&RowPkRangeBound>,
) -> bool {
lower.is_none_or(|bound| row_pk > &bound.row_pk || (bound.inclusive && row_pk == &bound.row_pk))
&& upper.is_none_or(|bound| {
row_pk < &bound.row_pk || (bound.inclusive && row_pk == &bound.row_pk)
})
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
pub(crate) struct TrackedStateReadColumns {
#[serde(default)]
pub(crate) columns: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
pub(crate) struct TrackedStateScanRequest {
#[serde(default)]
pub(crate) filter: TrackedStateFilter,
#[serde(default)]
pub(crate) read_columns: TrackedStateReadColumns,
#[serde(default)]
pub(crate) limit: Option<usize>,
}
#[derive(Debug, PartialEq, Eq)]
pub(crate) struct TrackedStateMutation {
pub(crate) encoded_key: Bytes,
pub(crate) encoded_value: Bytes,
}
impl TrackedStateMutation {
#[cfg(test)]
pub(crate) fn put_encoded(encoded_key: Vec<u8>, encoded_value: Vec<u8>) -> Self {
Self {
encoded_key: Bytes::from(encoded_key),
encoded_value: Bytes::from(encoded_value),
}
}
pub(crate) fn from_shared(encoded_key: Bytes, encoded_value: Bytes) -> Self {
Self {
encoded_key,
encoded_value,
}
}
}
#[derive(Debug, Default, PartialEq, Eq)]
pub(crate) struct TrackedStateMutationBatch {
mutations: Vec<TrackedStateMutation>,
}
impl TrackedStateMutationBatch {
pub(crate) fn from_shared(mutations: Vec<TrackedStateMutation>) -> Self {
Self { mutations }
}
pub(crate) fn len(&self) -> usize {
self.mutations.len()
}
pub(crate) fn first_encoded_key(&self) -> Option<&[u8]> {
self.mutations
.first()
.map(|mutation| mutation.encoded_key.as_ref())
}
#[cfg(test)]
pub(crate) fn as_slice(&self) -> &[TrackedStateMutation] {
&self.mutations
}
pub(crate) fn into_mutations(self) -> Vec<TrackedStateMutation> {
self.mutations
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct TrackedStateTreeScanRequest {
pub(crate) schema_keys: Vec<String>,
pub(crate) row_pks: Vec<RowPk>,
pub(crate) row_pk_lower: Option<RowPkRangeBound>,
pub(crate) row_pk_upper: Option<RowPkRangeBound>,
pub(crate) file_ids: Vec<NullableKeyFilter<String>>,
pub(crate) include_tombstones: bool,
pub(crate) limit: Option<usize>,
}
impl Default for TrackedStateTreeScanRequest {
fn default() -> Self {
Self {
schema_keys: Vec::new(),
row_pks: Vec::new(),
row_pk_lower: None,
row_pk_upper: None,
file_ids: Vec::new(),
include_tombstones: true,
limit: None,
}
}
}
impl TrackedStateTreeScanRequest {
pub(crate) fn matches(&self, key: &TrackedStateKey, value: &TrackedStateIndexValue) -> bool {
self.matches_ref(
TrackedStateKeyRef {
schema_key: &key.schema_key,
file_id: key.file_id.as_deref(),
row_pk: &key.row_pk,
},
value,
)
}
pub(crate) fn matches_ref(
&self,
key: TrackedStateKeyRef<'_>,
value: &TrackedStateIndexValue,
) -> bool {
if !self.include_tombstones && value.deleted {
return false;
}
self.matches_key_ref(key)
}
pub(crate) fn matches_key_ref(&self, key: TrackedStateKeyRef<'_>) -> bool {
if !self.schema_keys.is_empty()
&& !self
.schema_keys
.iter()
.any(|schema_key| schema_key == key.schema_key)
{
return false;
}
if !self.row_pks.is_empty() && !self.row_pks.contains(key.row_pk) {
return false;
}
if !row_pk_satisfies_bounds(
key.row_pk,
self.row_pk_lower.as_ref(),
self.row_pk_upper.as_ref(),
) {
return false;
}
if !self.file_ids.is_empty()
&& !self.file_ids.iter().any(|filter| match filter {
NullableKeyFilter::Any => true,
NullableKeyFilter::Null => key.file_id.is_none(),
NullableKeyFilter::Value(value) => key.file_id == Some(value.as_str()),
})
{
return false;
}
true
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct TrackedStateApplyResult {
pub(crate) root_id: TrackedStateRootId,
pub(crate) row_count: usize,
pub(crate) tree_height: usize,
pub(crate) chunk_count: usize,
pub(crate) chunk_bytes: usize,
}
#[cfg(test)]
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct TrackedStateTreeDiffEntry {
pub(crate) key: TrackedStateKey,
pub(crate) before: Option<TrackedStateIndexValue>,
pub(crate) after: Option<TrackedStateIndexValue>,
}