use super::super::types::{ReplicatedEntry, ReplicatedWrite};
use super::ctx::DecodeCtx;
use super::{
entry_array, entry_columnar_family, entry_crdt, entry_document, entry_graph, entry_kv, vector,
};
use crate::bridge::envelope::PhysicalPlan;
use crate::control::surrogate::SurrogateAssigner;
use crate::types::{DatabaseId, TenantId, VShardId};
pub type DecodedEntry = (TenantId, VShardId, PhysicalPlan, Option<u64>);
pub fn from_replicated_entry(
data: &[u8],
assigner: Option<&SurrogateAssigner>,
) -> crate::Result<Option<DecodedEntry>> {
let entry = match ReplicatedEntry::from_bytes(data) {
Some(e) => e,
None => return Ok(None),
};
match &entry.write {
ReplicatedWrite::ArrayOp { .. } | ReplicatedWrite::ArraySchema { .. } => {
return Ok(None);
}
_ => {}
}
let tenant_id = TenantId::new(entry.tenant_id);
let database_id = DatabaseId::new(entry.database_id);
let ctx = DecodeCtx {
assigner,
database_id,
tenant_id,
};
let (plan, resolved_now_ms) = to_physical_plan(&entry.write, &ctx)?;
Ok(Some((
tenant_id,
VShardId::new(entry.vshard_id),
plan,
resolved_now_ms,
)))
}
fn to_physical_plan(
write: &ReplicatedWrite,
ctx: &DecodeCtx,
) -> crate::Result<(PhysicalPlan, Option<u64>)> {
match write {
ReplicatedWrite::PointPut { .. }
| ReplicatedWrite::PointInsert { .. }
| ReplicatedWrite::PointDelete { .. }
| ReplicatedWrite::PointUpdate { .. }
| ReplicatedWrite::DocUpsert { .. }
| ReplicatedWrite::DocBatchInsert { .. }
| ReplicatedWrite::DocTruncate { .. }
| ReplicatedWrite::BulkDml { .. }
| ReplicatedWrite::InsertSelect { .. } => {
Ok((entry_document::decode_arm(ctx, write)?, None))
}
ReplicatedWrite::VectorInsert { .. }
| ReplicatedWrite::VectorBatchInsert { .. }
| ReplicatedWrite::VectorDelete { .. }
| ReplicatedWrite::SetVectorParams { .. }
| ReplicatedWrite::SparseInsert { .. }
| ReplicatedWrite::SparseDelete { .. }
| ReplicatedWrite::MultiVectorInsert { .. }
| ReplicatedWrite::MultiVectorDelete { .. }
| ReplicatedWrite::DeleteBySurrogate { .. }
| ReplicatedWrite::DirectUpsert { .. } => Ok((vector::decode_arm(ctx, write)?, None)),
ReplicatedWrite::CrdtApply { .. }
| ReplicatedWrite::CrdtImportCollection { .. }
| ReplicatedWrite::CrdtListInsert { .. }
| ReplicatedWrite::CrdtListDelete { .. }
| ReplicatedWrite::CrdtListMove { .. }
| ReplicatedWrite::CrdtDocUpsert { .. }
| ReplicatedWrite::CrdtDocDelete { .. }
| ReplicatedWrite::ConstraintChange { .. } => {
Ok((entry_crdt::decode_arm(ctx, write)?, None))
}
ReplicatedWrite::EdgePut { .. }
| ReplicatedWrite::EdgeDelete { .. }
| ReplicatedWrite::SetNodeLabels { .. }
| ReplicatedWrite::RemoveNodeLabels { .. }
| ReplicatedWrite::EdgePutBatch { .. }
| ReplicatedWrite::EdgeDeleteBatch { .. } => {
Ok((entry_graph::decode_arm(ctx, write)?, None))
}
ReplicatedWrite::KvTruncate { .. }
| ReplicatedWrite::KvPut { .. }
| ReplicatedWrite::KvDelete { .. }
| ReplicatedWrite::KvInsert { .. }
| ReplicatedWrite::KvInsertIfAbsent { .. }
| ReplicatedWrite::KvInsertOnConflictUpdate { .. }
| ReplicatedWrite::KvBatchPut { .. }
| ReplicatedWrite::KvExpire { .. }
| ReplicatedWrite::KvPersist { .. }
| ReplicatedWrite::KvIncr { .. }
| ReplicatedWrite::KvIncrFloat { .. }
| ReplicatedWrite::KvCas { .. }
| ReplicatedWrite::KvGetSet { .. }
| ReplicatedWrite::KvRegisterSortedIndex { .. }
| ReplicatedWrite::KvDropSortedIndex { .. }
| ReplicatedWrite::KvRegisterIndex { .. }
| ReplicatedWrite::KvDropIndex { .. }
| ReplicatedWrite::KvFieldSet { .. }
| ReplicatedWrite::KvTransfer { .. }
| ReplicatedWrite::KvTransferItem { .. } => entry_kv::decode_arm(ctx, write),
ReplicatedWrite::ColumnarIngest { .. }
| ReplicatedWrite::TimeseriesIngest { .. }
| ReplicatedWrite::FtsIndex { .. }
| ReplicatedWrite::FtsDelete { .. }
| ReplicatedWrite::SpatialInsert { .. }
| ReplicatedWrite::SpatialDelete { .. }
| ReplicatedWrite::ColumnarBulkDml { .. } => {
Ok((entry_columnar_family::decode_arm(write)?, None))
}
ReplicatedWrite::ArrayCellPut { .. } | ReplicatedWrite::ArrayCellDelete { .. } => {
Ok((entry_array::decode_arm(ctx, write)?, None))
}
ReplicatedWrite::ArrayOp { .. } => Err(crate::Error::Internal {
detail: "ArrayOp reached to_physical_plan (should have been intercepted)".into(),
}),
ReplicatedWrite::ArraySchema { .. } => Err(crate::Error::Internal {
detail: "ArraySchema reached to_physical_plan (should have been intercepted)".into(),
}),
ReplicatedWrite::CalvinReadResult { .. } => Err(crate::Error::Internal {
detail: "CalvinReadResult reached to_physical_plan (should have been intercepted)"
.into(),
}),
}
}