use super::super::types::ReplicatedWrite;
use super::ctx::DecodeCtx;
use crate::bridge::envelope::PhysicalPlan;
use crate::engine::array::wal::ArrayPutCell;
use nodedb_array::types::ArrayId;
use nodedb_physical::physical_plan::ArrayOp;
use nodedb_types::sync::wire::SyncProvenance;
pub(super) fn decode_arm(ctx: &DecodeCtx, write: &ReplicatedWrite) -> crate::Result<PhysicalPlan> {
match write {
ReplicatedWrite::ArrayCellPut {
array,
cells_msgpack,
provenance,
} => cell_put(ctx, array, cells_msgpack, provenance),
ReplicatedWrite::ArrayCellDelete {
array,
coords_msgpack,
provenance,
} => cell_delete(ctx, array, coords_msgpack, provenance),
_ => Err(crate::Error::Internal {
detail: "entry_array::decode_arm called with a non-array-cell ReplicatedWrite \
variant (dispatch bug in decode/entry.rs's grouped array match arm)"
.into(),
}),
}
}
fn cell_put(
ctx: &DecodeCtx,
array: &str,
cells_msgpack: &[u8],
provenance: &Option<Vec<u8>>,
) -> crate::Result<PhysicalPlan> {
let cells: Vec<ArrayPutCell> =
zerompk::from_msgpack(cells_msgpack).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("array cell put decode: {e}"),
})?;
if let Some(assigner) = ctx.assigner {
for cell in &cells {
let pk_bytes =
zerompk::to_msgpack_vec(&cell.coord).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("array coord pk encode: {e}"),
})?;
assigner.bind(
ctx.database_id,
ctx.tenant_id,
array,
&pk_bytes,
cell.surrogate,
)?;
}
}
Ok(PhysicalPlan::Array(ArrayOp::Put {
array_id: ArrayId::new(ctx.tenant_id, array),
cells_msgpack: cells_msgpack.to_vec(),
wal_lsn: 0,
provenance: decode_provenance(provenance)?,
}))
}
fn cell_delete(
ctx: &DecodeCtx,
array: &str,
coords_msgpack: &[u8],
provenance: &Option<Vec<u8>>,
) -> crate::Result<PhysicalPlan> {
Ok(PhysicalPlan::Array(ArrayOp::Delete {
array_id: ArrayId::new(ctx.tenant_id, array),
coords_msgpack: coords_msgpack.to_vec(),
wal_lsn: 0,
provenance: decode_provenance(provenance)?,
}))
}
fn decode_provenance(bytes: &Option<Vec<u8>>) -> crate::Result<Option<SyncProvenance>> {
match bytes {
None => Ok(None),
Some(b) => zerompk::from_msgpack::<SyncProvenance>(b)
.map(Some)
.map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("array provenance decode: {e}"),
}),
}
}