use std::sync::Arc;
use tokio::sync::mpsc;
use tracing::debug;
use crate::bridge::envelope::{PhysicalPlan, Status};
use crate::control::array_sync::raft_apply::{
AppliedPosition, ArrayCellTarget, apply_array_cell_write, apply_array_op, apply_array_schema,
};
use crate::control::cluster::calvin::ReadResultEvent;
use crate::control::server::dispatch_utils::{
ChangeFeedOwner, SubmitWrite, WalDurability, WriteOrdering, submit_write,
};
use crate::control::state::SharedState;
use crate::control::wal_replication::{ReplicatedEntry, ReplicatedWrite, from_replicated_entry};
use crate::types::{DatabaseId, TenantId, TraceId};
use nodedb_physical::physical_plan::ArrayOp;
use super::applied_index::{AppliedPrefix, save_applied_index};
use super::applier::ApplyBatch;
use super::propose_tracker::{AppliedWrite, ProposeTracker};
pub async fn run_apply_loop(
mut apply_rx: mpsc::Receiver<ApplyBatch>,
state: Arc<SharedState>,
tracker: Arc<ProposeTracker>,
calvin_read_result_senders: Arc<
std::sync::Mutex<std::collections::BTreeMap<u32, mpsc::Sender<ReadResultEvent>>>,
>,
) {
while let Some(batch) = apply_rx.recv().await {
let mut prefix = AppliedPrefix::new();
for entry in &batch.entries {
let replicated_opt = ReplicatedEntry::from_bytes(&entry.data);
let applied_key = replicated_opt
.as_ref()
.map(|e| e.idempotency_key)
.unwrap_or(0);
let database_id = replicated_opt
.as_ref()
.map(|e| DatabaseId::new(e.database_id))
.unwrap_or(DatabaseId::DEFAULT);
if let Some(replicated) = replicated_opt {
let target_vshard = replicated.vshard_id;
match replicated.write {
ReplicatedWrite::ArrayOp {
ref array,
ref op_bytes,
ref provenance,
..
} => {
let applied_ok = apply_array_op(
&state,
&tracker,
AppliedPosition {
group_id: batch.group_id,
log_index: entry.index,
applied_key,
},
array,
op_bytes,
provenance.as_deref(),
)
.await;
prefix.record(entry.index, applied_ok);
continue;
}
ReplicatedWrite::ArraySchema {
ref array,
ref snapshot_payload,
schema_hlc_bytes,
} => {
let applied_ok = apply_array_schema(
&state,
&tracker,
AppliedPosition {
group_id: batch.group_id,
log_index: entry.index,
applied_key,
},
crate::control::array_sync::raft_apply::ArraySchemaPayload {
array,
snapshot_payload,
schema_hlc_bytes,
},
);
prefix.record(entry.index, applied_ok);
continue;
}
ReplicatedWrite::CalvinReadResult {
epoch,
position,
passive_vshard,
tenant_id,
ref values,
} => {
let decoded_values: Vec<(
nodedb_physical::physical_plan::meta::PassiveReadKeyId,
nodedb_types::Value,
)> = match zerompk::from_msgpack(values) {
Ok(decoded) => decoded,
Err(e) => {
tracing::warn!(
group_id = batch.group_id,
index = entry.index,
error = %e,
"failed to decode CalvinReadResult payload"
);
tracker.complete(
batch.group_id,
entry.index,
applied_key,
Err(crate::Error::Internal {
detail: format!("decode CalvinReadResult payload: {e}"),
}),
);
prefix.skip();
continue;
}
};
let event = ReadResultEvent {
epoch,
position,
passive_vshard,
tenant_id: TenantId::new(tenant_id),
values: decoded_values,
};
let send_result = calvin_read_result_senders
.lock()
.unwrap_or_else(|p| p.into_inner())
.get(&target_vshard)
.cloned()
.map(|sender| sender.try_send(event));
if let Some(Err(e)) = send_result {
tracing::warn!(
group_id = batch.group_id,
index = entry.index,
error = %e,
"failed to forward CalvinReadResult to scheduler"
);
}
tracker.complete(
batch.group_id,
entry.index,
applied_key,
Ok(AppliedWrite::unversioned(Vec::new())),
);
prefix.skip();
continue;
}
_ => {}
}
}
let decoded =
from_replicated_entry(&entry.data, Some(state.surrogate_assigner.as_ref()));
let (tenant_id, vshard_id, plan, resolved_now_ms) = match decoded {
Ok(Some(t)) => t,
Ok(None) => {
debug!(
group_id = batch.group_id,
index = entry.index,
"skipping non-ReplicatedEntry commit"
);
tracker.complete(
batch.group_id,
entry.index,
applied_key,
Ok(AppliedWrite::unversioned(Vec::new())),
);
prefix.skip();
continue;
}
Err(e) => {
tracing::warn!(
group_id = batch.group_id,
index = entry.index,
error = %e,
"failed to decode replicated entry (surrogate bind error)"
);
tracker.complete(
batch.group_id,
entry.index,
applied_key,
Err(crate::Error::Internal {
detail: format!("decode replicated entry: {e}"),
}),
);
prefix.record(entry.index, false);
continue;
}
};
if matches!(
plan,
PhysicalPlan::Array(ArrayOp::Put { .. } | ArrayOp::Delete { .. })
) {
let applied_ok = apply_array_cell_write(
&state,
&tracker,
AppliedPosition {
group_id: batch.group_id,
log_index: entry.index,
applied_key,
},
ArrayCellTarget {
tenant_id,
database_id,
vshard: vshard_id,
resolved_now_ms,
},
plan,
)
.await;
prefix.record(entry.index, applied_ok);
continue;
}
let submitted = submit_write(
&state,
SubmitWrite {
tenant_id,
database_id,
vshard_id,
plan,
trace_id: TraceId::generate(),
event_source: crate::event::EventSource::User,
txn_id: None,
user_id: None,
durability: WalDurability::AppendHere {
now_override: resolved_now_ms,
},
ordering: WriteOrdering::AlreadyOrdered,
change_feed: ChangeFeedOwner::Unowned,
},
)
.await
.map(|outcome| outcome.response);
let result = match submitted {
Ok(resp) if resp.status == Status::Ok => Ok(AppliedWrite::from_response(&resp)),
Ok(resp) => {
let reason = resp
.error_code
.as_ref()
.map(|c| format!("{c:?}"))
.unwrap_or_else(|| "execution error".into());
tracing::warn!(
group_id = batch.group_id,
index = entry.index,
reason = %reason,
"applying committed write failed"
);
Err(crate::Error::Internal { detail: reason })
}
Err(e) => {
tracing::warn!(
group_id = batch.group_id,
index = entry.index,
error = %e,
"applying committed write failed"
);
Err(crate::Error::Internal {
detail: e.to_string(),
})
}
};
let applied_ok = result.is_ok();
tracker.complete(batch.group_id, entry.index, applied_key, result);
prefix.record(entry.index, applied_ok);
}
if let Some(applied_index) = prefix.floor() {
record_durable_apply(&state, batch.group_id, applied_index);
}
}
}
fn record_durable_apply(state: &Arc<SharedState>, group_id: u64, applied_index: u64) {
save_applied_index(state, group_id, applied_index);
maybe_compact_log(state, group_id, applied_index);
}
fn maybe_compact_log(state: &Arc<SharedState>, group_id: u64, applied_index: u64) {
let Some(compactor) = state.raft_compactor.get() else {
return;
};
match compactor(group_id, applied_index) {
Ok(true) => {
debug!(
group_id,
applied_index, "raft log compacted past data-plane applied watermark"
);
}
Ok(false) => {}
Err(e) => {
tracing::warn!(
group_id,
applied_index,
error = %e,
"raft log compaction failed"
);
}
}
}