use std::sync::Arc;
use std::time::{Duration, Instant};
use crate::bridge::envelope::{PhysicalPlan, Priority, Request, Response, Status};
use crate::control::server::wal_dispatch::{self, WalAppendRequest};
use crate::control::state::SharedState;
use crate::types::{DatabaseId, Lsn, ReadConsistency, TenantId, TraceId, TxnId, VShardId};
use super::change_events::{extract_write_change_set, publish_change_set};
use super::collect::{DispatchCollectError, collect_bounded_response};
pub(crate) enum WalDurability {
AppendHere { now_override: Option<u64> },
CallerSupplied {
wal_lsn: Option<Lsn>,
resolved_now_ms: Option<u64>,
},
}
pub(crate) enum WriteOrdering {
Gate,
AlreadyOrdered,
}
pub(crate) enum ChangeFeedOwner {
Funnel,
Unowned,
}
pub(crate) struct SubmitOutcome {
pub response: Response,
pub wal_lsn: Option<Lsn>,
}
pub(crate) struct SubmitWrite {
pub tenant_id: TenantId,
pub database_id: DatabaseId,
pub vshard_id: VShardId,
pub plan: PhysicalPlan,
pub trace_id: TraceId,
pub event_source: crate::event::EventSource,
pub txn_id: Option<TxnId>,
pub user_id: Option<Arc<str>>,
pub durability: WalDurability,
pub ordering: WriteOrdering,
pub change_feed: ChangeFeedOwner,
}
pub(crate) async fn submit_write(
shared: &SharedState,
params: SubmitWrite,
) -> crate::Result<SubmitOutcome> {
let SubmitWrite {
tenant_id,
database_id,
vshard_id,
mut plan,
trace_id,
event_source,
txn_id,
user_id,
durability,
ordering,
change_feed,
} = params;
let change_set = match change_feed {
ChangeFeedOwner::Funnel => Some(extract_write_change_set(&plan, tenant_id)),
ChangeFeedOwner::Unowned => None,
};
let post_apply = wal_dispatch::plan_post_apply_redo(&plan);
let appends_here = matches!(&durability, WalDurability::AppendHere { .. });
use crate::control::server::shared::write_admission::{
WriteAdmission, WriteTarget, admit, bare_ok_response, route_write_to_calvin,
};
let (admission, admission_guard, order_guard) = match ordering {
WriteOrdering::AlreadyOrdered => (
crate::bridge::envelope::Admission::Exempt(
crate::bridge::envelope::ExemptReason::AlreadyOrdered,
),
None,
None,
),
WriteOrdering::Gate => match admit(
shared,
&WriteTarget {
tenant_id,
database_id,
vshard_id,
plan: &plan,
},
) {
WriteAdmission::ExemptRead => (
crate::bridge::envelope::Admission::Exempt(
crate::bridge::envelope::ExemptReason::Read,
),
None,
None,
),
WriteAdmission::FastPath { guard } => {
(crate::bridge::envelope::Admission::Admitted, guard, None)
}
WriteAdmission::FastPathBlocking { key, keyed_lock } => {
let order_guard = keyed_lock.lock_owned(key).await;
(
crate::bridge::envelope::Admission::Admitted,
None,
Some(order_guard),
)
}
WriteAdmission::RouteToCalvin => {
let routed =
route_write_to_calvin(shared, tenant_id, database_id, vshard_id, plan).await?;
return Ok(SubmitOutcome {
response: routed
.unwrap_or_else(|| bare_ok_response(crate::types::RequestId::new(0))),
wal_lsn: None,
});
}
},
};
let (wal_lsn, resolved_now_ms) = match durability {
WalDurability::AppendHere { now_override } => {
let outcome = wal_dispatch::wal_append(WalAppendRequest {
wal: &shared.wal,
tenant_id,
vshard_id,
database_id,
plan: &plan,
credentials: None,
now_override,
})?;
(outcome.lsn, outcome.resolved_now_ms)
}
WalDurability::CallerSupplied {
wal_lsn,
resolved_now_ms,
} => (wal_lsn, resolved_now_ms),
};
if let Some(lsn) = wal_lsn {
wal_dispatch::stamp_minted_lsn(&mut plan, lsn);
}
let dispatch_started = Instant::now();
let vshard_u32 = vshard_id.as_u32();
let observe = |shared: &SharedState| {
let latency_us = dispatch_started.elapsed().as_micros().min(u64::MAX as u128) as u64;
shared.per_vshard_metrics.observe(vshard_u32, latency_us);
};
let request_id = shared.next_request_id();
let request = Request {
request_id,
tenant_id,
database_id,
vshard_id,
plan,
deadline: Instant::now() + Duration::from_secs(shared.tuning.network.default_deadline_secs),
priority: Priority::Normal,
trace_id,
consistency: ReadConsistency::Strong,
idempotency_key: None,
event_source,
user_roles: Vec::new(),
user_id,
statement_digest: None,
txn_id,
wal_lsn,
resolved_now_ms,
admission,
};
let mut rx = shared.tracker.register(request_id);
match shared.dispatcher.lock() {
Ok(mut d) => d.dispatch(request)?,
Err(poisoned) => poisoned.into_inner().dispatch(request)?,
};
let deferred_guards = if post_apply.is_some() {
Some((admission_guard, order_guard))
} else {
drop(admission_guard);
drop(order_guard);
None
};
let max_result_bytes = shared.tuning.network.max_query_result_bytes as usize;
let response = tokio::time::timeout(
Duration::from_secs(shared.tuning.network.default_deadline_secs),
collect_bounded_response(&mut rx, max_result_bytes),
)
.await
.map_err(|_| {
observe(shared);
crate::Error::DeadlineExceeded { request_id }
})?;
let response = match response {
Ok(r) => r,
Err(DispatchCollectError::OverBudget { bytes }) => {
shared.tracker.cancel(&request_id);
observe(shared);
return Err(crate::Error::ExecutionLimitExceeded {
detail: format!(
"query result exceeded max_query_result_bytes \
({bytes} > {max_result_bytes} bytes)"
),
});
}
Err(DispatchCollectError::ChannelClosed) => {
observe(shared);
return Err(crate::Error::Dispatch {
detail: "response channel closed".into(),
});
}
};
let post_apply_lsn = if let Some(collection) = &post_apply
&& appends_here
&& response.status == Status::Ok
{
wal_dispatch::append_write_set_redo(
&shared.wal,
tenant_id,
vshard_id,
database_id,
collection,
&response.write_set,
)?
} else {
None
};
drop(deferred_guards);
if response.status == Status::Ok {
let durable_target = match (wal_lsn, post_apply_lsn) {
(Some(a), Some(b)) => Some(a.max(b)),
(a, b) => a.or(b),
};
if let Some(lsn) = durable_target {
shared.wal.wait_durable(lsn).await?;
}
}
if response.status == Status::Ok {
if let Some(change_set) = change_set {
publish_change_set(shared, tenant_id, database_id, change_set, &response);
}
shared.advance_tenant_write_hlc(tenant_id.as_u64());
}
observe(shared);
Ok(SubmitOutcome { response, wal_lsn })
}