use std::sync::Arc;
use tracing::warn;
use super::common::{AppliedPosition, ArrayWriteSubmit, ensure_array_open, submit_array_write};
use crate::bridge::envelope::PhysicalPlan;
use crate::control::distributed_applier::ProposeTracker;
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TenantId, VShardId};
use nodedb_physical::physical_plan::ArrayOp;
pub(crate) struct ArrayCellTarget {
pub tenant_id: TenantId,
pub database_id: DatabaseId,
pub vshard: VShardId,
pub resolved_now_ms: Option<u64>,
}
pub(crate) async fn apply_array_cell_write(
state: &Arc<SharedState>,
tracker: &Arc<ProposeTracker>,
pos: AppliedPosition,
target: ArrayCellTarget,
plan: PhysicalPlan,
) -> bool {
let AppliedPosition {
group_id,
log_index,
applied_key,
} = pos;
let ArrayCellTarget {
tenant_id,
database_id,
vshard,
resolved_now_ms,
} = target;
let array_id = match &plan {
PhysicalPlan::Array(ArrayOp::Put { array_id, .. })
| PhysicalPlan::Array(ArrayOp::Delete { array_id, .. }) => array_id.clone(),
other => {
let e = crate::Error::Internal {
detail: format!(
"apply_array_cell_write called with a non-array-cell plan: {other:?}"
),
};
tracker.complete(group_id, log_index, applied_key, Err(e));
return false;
}
};
if let Err(e) = ensure_array_open(state, &array_id, vshard, tenant_id).await {
warn!(
group_id, index = log_index, array = %array_id.name, error = %e,
"apply_array_cell_write: ensure_array_open failed"
);
tracker.complete(group_id, log_index, applied_key, Err(e));
return false;
}
let result = submit_array_write(
state,
ArrayWriteSubmit {
tenant_id,
database_id,
vshard,
plan,
event_source: crate::event::EventSource::User,
resolved_now_ms,
op_label: "array cell write",
},
)
.await;
if let Err(e) = &result {
warn!(
group_id, index = log_index, array = %array_id.name, error = %e,
"apply_array_cell_write: apply failed"
);
}
let applied_ok = result.is_ok();
tracker.complete(group_id, log_index, applied_key, result);
applied_ok
}