use std::sync::Arc;
use std::time::Duration;
use crate::bridge::envelope::{PhysicalPlan, Response, Status};
use crate::control::state::SharedState;
use crate::control::wal_replication::{AsyncRaftProposer, ReplicatedEntry, to_replicated_entry};
use crate::event::EventSource;
use crate::types::{DatabaseId, Lsn, TenantId, TraceId, VShardId};
pub async fn propose_sync_write(
state: &SharedState,
entry: ReplicatedEntry,
proposer: &Arc<AsyncRaftProposer>,
) -> crate::Result<Vec<u8>> {
let idempotency_key = entry.idempotency_key;
let data = entry.to_bytes();
let vshard_id = entry.vshard_id;
const BACKOFF_MS: [u64; 5] = [10, 25, 50, 100, 200];
let mut payload: Option<Vec<u8>> = None;
let mut last_err: Option<crate::Error> = None;
for (attempt, backoff_ms) in BACKOFF_MS.iter().enumerate() {
match proposer(vshard_id, idempotency_key, data.clone()).await {
Ok((p, _committed_version)) => {
payload = Some(p);
break;
}
Err(crate::Error::RetryableLeaderChange {
group_id,
log_index,
}) => {
state
.raft_propose_leader_change_retries
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::warn!(
attempt,
group_id,
log_index,
"raft entry overwritten by leader change — re-proposing"
);
last_err = Some(crate::Error::RetryableLeaderChange {
group_id,
log_index,
});
tokio::time::sleep(Duration::from_millis(*backoff_ms)).await;
continue;
}
Err(other) => {
return Err(crate::Error::Dispatch {
detail: format!("raft propose failed: {other}"),
});
}
}
}
payload.ok_or_else(|| {
last_err.unwrap_or_else(|| crate::Error::Dispatch {
detail: "raft propose retries exhausted".into(),
})
})
}
pub async fn dispatch_sync_response(
state: &SharedState,
tenant_id: TenantId,
vshard_id: VShardId,
plan: PhysicalPlan,
trace_id: TraceId,
event_source: EventSource,
) -> crate::Result<Response> {
if let Some(proposer) = state.async_raft_proposer.get()
&& let Some(entry) = to_replicated_entry(tenant_id, DatabaseId::DEFAULT, vshard_id, &plan)
{
let payload = propose_sync_write(state, entry, proposer).await?;
let request_id = state.next_request_id();
return Ok(Response {
request_id,
status: Status::Ok,
attempt: 1,
partial: false,
payload: payload.into(),
watermark_lsn: Lsn::new(0),
error_code: None,
read_set_valid: None,
read_version_lsn: crate::types::Lsn::ZERO,
write_set: Vec::new(),
});
}
crate::control::server::dispatch_utils::dispatch_to_data_plane_with_source(
state,
tenant_id,
DatabaseId::DEFAULT,
vshard_id,
plan,
trace_id,
event_source,
)
.await
}
pub async fn dispatch_sync_payload(
state: &SharedState,
tenant_id: TenantId,
vshard_id: VShardId,
plan: PhysicalPlan,
) -> crate::Result<Vec<u8>> {
let response = dispatch_sync_response(
state,
tenant_id,
vshard_id,
plan,
TraceId::ZERO,
EventSource::CrdtSync,
)
.await?;
Ok(response.payload.to_vec())
}
pub fn noop_dispatch_error(op: &str) -> crate::Error {
crate::Error::Internal {
detail: format!(
"{op} routed through path lacking SharedState; \
check listener wiring — {op} was ACKed but NOT applied"
),
}
}
pub async fn dispatch_sync_bytes(
state: &SharedState,
tenant_id: TenantId,
collection: &str,
plan: PhysicalPlan,
timeout: Duration,
event_source: EventSource,
) -> crate::Result<Vec<u8>> {
dispatch_write_replicated(
state,
tenant_id,
DatabaseId::DEFAULT,
collection,
plan,
timeout,
event_source,
)
.await
}
pub async fn dispatch_write_replicated(
state: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
collection: &str,
plan: PhysicalPlan,
timeout: Duration,
event_source: EventSource,
) -> crate::Result<Vec<u8>> {
let vshard_id = VShardId::from_collection_in_database(database_id, collection);
if let Some(proposer) = state.async_raft_proposer.get()
&& let Some(entry) = to_replicated_entry(tenant_id, database_id, vshard_id, &plan)
{
return propose_sync_write(state, entry, proposer).await;
}
let resp =
crate::control::server::shared::ddl::sync_dispatch::dispatch_async_response_with_source(
state,
tenant_id,
database_id,
collection,
plan,
timeout,
event_source,
)
.await?;
if resp.status != Status::Ok {
return Err(match resp.error_code {
Some(code) => crate::Error::DataPlane(*code),
None => crate::Error::Internal {
detail: String::from_utf8_lossy(&resp.payload).into_owned(),
},
});
}
state.advance_tenant_write_hlc(tenant_id.as_u64());
Ok(resp.payload.to_vec())
}