use std::sync::Arc;
use tracing::warn;
use super::common::{
AppliedPosition, ArrayWriteSubmit, ensure_array_open, submit_array_write, vshard_for_array_op,
};
use crate::control::array_sync::OriginApplyEngine;
use crate::control::distributed_applier::{AppliedWrite, ProposeTracker};
use crate::control::state::SharedState;
pub(crate) async fn apply_array_op(
state: &Arc<SharedState>,
tracker: &Arc<ProposeTracker>,
pos: AppliedPosition,
array: &str,
op_bytes: &[u8],
provenance_bytes: Option<&[u8]>,
) -> bool {
let AppliedPosition {
group_id,
log_index,
applied_key,
} = pos;
use crate::types::TenantId;
use nodedb_array::sync::op_codec;
let provenance: Option<nodedb_types::sync::wire::SyncProvenance> = match provenance_bytes {
None => None,
Some(b) => match zerompk::from_msgpack::<nodedb_types::sync::wire::SyncProvenance>(b) {
Ok(p) => Some(p),
Err(e) => {
warn!(
group_id, index = log_index, array = %array, error = %e,
"apply_array_op: provenance decode failed; applying without epoch fence (version skew or corruption)"
);
None
}
},
};
let op = match op_codec::decode_op(op_bytes) {
Ok(op) => op,
Err(e) => {
warn!(
group_id, index = log_index, array = %array, error = %e,
"apply_array_op: decode failed"
);
tracker.complete(
group_id,
log_index,
applied_key,
Err(crate::Error::Internal {
detail: format!("array op decode: {e}"),
}),
);
return false;
}
};
let engine = OriginApplyEngine::new(
Arc::clone(&state.array_sync_schemas),
Arc::clone(&state.array_sync_op_log),
);
if engine.already_seen(&op.header.array, op.header.hlc) {
tracker.complete(
group_id,
log_index,
applied_key,
Ok(AppliedWrite::unversioned(Vec::new())),
);
return true;
}
let vshard = vshard_for_array_op(state, &op);
use nodedb_array::sync::op::ArrayOpKind;
use nodedb_physical::physical_plan::ArrayOp as DataArrayOp;
let tenant_id = TenantId::new(0); let array_id = nodedb_array::types::ArrayId::new(tenant_id, &op.header.array);
if let Err(e) = ensure_array_open(state, &array_id, vshard, tenant_id).await {
warn!(
group_id, index = log_index, array = %array, error = %e,
"apply_array_op: ensure_array_open failed"
);
tracker.complete(group_id, log_index, applied_key, Err(e));
return false;
}
let data_op = match op.kind {
ArrayOpKind::Put => {
let cells = vec![crate::engine::array::wal::ArrayPutCell {
coord: op.coord.clone(),
attrs: op.attrs.clone().unwrap_or_default(),
surrogate: nodedb_types::Surrogate::ZERO,
system_from_ms: op.header.system_from_ms,
valid_from_ms: op.header.valid_from_ms,
valid_until_ms: op.header.valid_until_ms,
}];
let cells_msgpack = match zerompk::to_msgpack_vec(&cells) {
Ok(b) => b,
Err(e) => {
warn!(group_id, index = log_index, error = %e, "apply_array_op: cells encode failed");
tracker.complete(
group_id,
log_index,
applied_key,
Err(crate::Error::Internal {
detail: format!("cells encode: {e}"),
}),
);
return false;
}
};
DataArrayOp::Put {
array_id,
cells_msgpack,
wal_lsn: 0,
provenance: provenance.clone(),
}
}
ArrayOpKind::Delete | ArrayOpKind::Erase => {
let coords = vec![op.coord.clone()];
let coords_msgpack = match zerompk::to_msgpack_vec(&coords) {
Ok(b) => b,
Err(e) => {
warn!(group_id, index = log_index, error = %e, "apply_array_op: coords encode failed");
tracker.complete(
group_id,
log_index,
applied_key,
Err(crate::Error::Internal {
detail: format!("coords encode: {e}"),
}),
);
return false;
}
};
DataArrayOp::Delete {
array_id,
coords_msgpack,
wal_lsn: 0,
provenance,
}
}
};
let plan = crate::bridge::envelope::PhysicalPlan::Array(data_op);
let result = submit_array_write(
state,
ArrayWriteSubmit {
tenant_id,
database_id: crate::types::DatabaseId::DEFAULT,
vshard,
plan,
event_source: crate::event::EventSource::CrdtSync,
resolved_now_ms: None,
op_label: "array op",
},
)
.await;
match result {
Ok(applied) => {
if let Err(e) = engine.record_applied(&op) {
tracing::error!(
group_id, index = log_index, array = %op.header.array,
error = %e,
"apply_array_op: op applied but op-log append failed"
);
}
tracker.complete(group_id, log_index, applied_key, Ok(applied));
true
}
Err(e) => {
warn!(
group_id, index = log_index, array = %op.header.array, error = %e,
"apply_array_op: apply failed"
);
tracker.complete(group_id, log_index, applied_key, Err(e));
false
}
}
}