use std::sync::Arc;
use std::time::Duration;
use crate::bridge::envelope::{PhysicalPlan, Priority, Request, Response, Status};
use crate::control::distributed_applier::{AppliedWrite, ProposeResult};
use crate::control::server::dispatch_utils::{
ChangeFeedOwner, SubmitWrite, WalDurability, WriteOrdering, submit_write,
};
use crate::control::state::SharedState;
use crate::types::{DatabaseId, ReadConsistency, TenantId, TraceId, VShardId};
#[derive(Debug, Clone, Copy)]
pub(crate) struct AppliedPosition {
pub group_id: u64,
pub log_index: u64,
pub applied_key: u64,
}
pub(super) struct ArrayWriteSubmit {
pub tenant_id: TenantId,
pub database_id: DatabaseId,
pub vshard: VShardId,
pub plan: PhysicalPlan,
pub event_source: crate::event::EventSource,
pub resolved_now_ms: Option<u64>,
pub op_label: &'static str,
}
pub(super) async fn submit_array_write(
state: &Arc<SharedState>,
params: ArrayWriteSubmit,
) -> ProposeResult {
let ArrayWriteSubmit {
tenant_id,
database_id,
vshard,
plan,
event_source,
resolved_now_ms,
op_label,
} = params;
let outcome = submit_write(
state,
SubmitWrite {
tenant_id,
database_id,
vshard_id: vshard,
plan,
trace_id: TraceId::generate(),
event_source,
txn_id: None,
user_id: None,
durability: WalDurability::AppendHere {
now_override: resolved_now_ms,
},
ordering: WriteOrdering::AlreadyOrdered,
change_feed: ChangeFeedOwner::Unowned,
},
)
.await
.map_err(|e| crate::Error::Internal {
detail: format!("{op_label}: {e}"),
})?;
let response = outcome.response;
if response.status != Status::Ok {
let detail = response
.error_code
.as_ref()
.map(|c| format!("{op_label} error: {c:?}"))
.unwrap_or_else(|| format!("{op_label} returned error status"));
return Err(crate::Error::Internal { detail });
}
Ok(AppliedWrite::from_response(&response))
}
pub(super) fn vshard_for_array_op(
state: &Arc<SharedState>,
op: &nodedb_array::sync::op::ArrayOp,
) -> VShardId {
use nodedb_array::types::coord::value::CoordValue;
use nodedb_cluster::array_routing::{array_vshard_for_name, vshard_for_array_coord};
let tile_extents = state.array_sync_schemas.tile_extents(&op.header.array);
if let Some(extents) = tile_extents {
let coord_u64: Vec<u64> = op
.coord
.iter()
.map(|c| match c {
CoordValue::Int64(v) | CoordValue::TimestampMs(v) => *v as u64,
CoordValue::Float64(v) => v.to_bits(),
CoordValue::String(_) => 0,
})
.collect();
VShardId::new(vshard_for_array_coord(
&op.header.array,
&coord_u64,
&extents,
))
} else {
VShardId::new(array_vshard_for_name(&op.header.array))
}
}
pub(super) async fn ensure_array_open(
state: &Arc<SharedState>,
array_id: &nodedb_array::types::ArrayId,
vshard: crate::types::VShardId,
tenant_id: crate::types::TenantId,
) -> crate::Result<()> {
let (schema_msgpack, schema_hash, prefix_bits) = {
let cat = state
.array_catalog
.read()
.unwrap_or_else(|p| p.into_inner());
match cat.lookup_by_name(&array_id.name) {
Some(entry) => (
entry.schema_msgpack.clone(),
entry.schema_hash,
entry.prefix_bits,
),
None => {
return Err(crate::Error::Internal {
detail: format!(
"ensure_array_open: array '{}' not in catalog — register it before applying ops",
array_id.name
),
});
}
}
};
let open_plan = crate::bridge::envelope::PhysicalPlan::Array(
nodedb_physical::physical_plan::ArrayOp::OpenArray {
array_id: array_id.clone(),
schema_msgpack,
schema_hash,
prefix_bits,
},
);
let open_request = build_array_request(state, tenant_id, vshard, open_plan);
let open_request_id = open_request.request_id;
let mut open_rx = state.tracker.register(open_request_id);
let dispatch_result = match state.dispatcher.lock() {
Ok(mut d) => d.dispatch(open_request),
Err(poisoned) => poisoned.into_inner().dispatch(open_request),
};
if let Err(e) = dispatch_result {
return Err(crate::Error::Internal {
detail: format!("ensure_array_open: dispatch failed: {e}"),
});
}
await_data_plane(async move { open_rx.recv().await.ok_or(()) }, "OpenArray")
.await
.map(|_| ())
}
pub(super) fn build_array_request(
state: &Arc<SharedState>,
tenant_id: TenantId,
vshard_id: VShardId,
plan: crate::bridge::envelope::PhysicalPlan,
) -> Request {
Request {
request_id: state.next_request_id(),
tenant_id,
database_id: DatabaseId::DEFAULT,
vshard_id,
plan,
deadline: std::time::Instant::now() + Duration::from_secs(30),
priority: Priority::Normal,
trace_id: TraceId::generate(),
consistency: ReadConsistency::Strong,
idempotency_key: None,
event_source: crate::event::EventSource::CrdtSync,
user_roles: Vec::new(),
user_id: None,
statement_digest: None,
txn_id: None,
wal_lsn: None,
resolved_now_ms: None,
admission: crate::bridge::envelope::Admission::Exempt(
crate::bridge::envelope::ExemptReason::AlreadyOrdered,
),
}
}
pub(super) async fn await_data_plane(
rx: impl std::future::Future<Output = Result<Response, ()>>,
op_label: &str,
) -> ProposeResult {
match tokio::time::timeout(Duration::from_secs(30), rx).await {
Ok(Ok(resp)) if resp.status == Status::Ok => Ok(AppliedWrite::from_response(&resp)),
Ok(Ok(resp)) => {
let detail = resp
.error_code
.as_ref()
.map(|c| format!("{op_label} error: {c:?}"))
.unwrap_or_else(|| format!("{op_label} returned error status"));
Err(crate::Error::Internal { detail })
}
Ok(Err(_)) => Err(crate::Error::Internal {
detail: format!("{op_label}: response channel closed"),
}),
Err(_) => Err(crate::Error::Internal {
detail: format!("{op_label}: deadline exceeded"),
}),
}
}