use crate::bridge::envelope::{PhysicalPlan, Response};
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TenantId, TraceId, VShardId};
use super::submit_write::{
ChangeFeedOwner, SubmitWrite, WalDurability, WriteOrdering, submit_write,
};
use super::types::{AutocommitWrite, DataPlaneDispatch, WriteDispatch};
pub async fn dispatch_to_data_plane(
shared: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
vshard_id: VShardId,
plan: PhysicalPlan,
trace_id: TraceId,
) -> crate::Result<Response> {
dispatch_to_data_plane_with_source(
shared,
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
crate::event::EventSource::User,
)
.await
}
pub async fn dispatch_to_data_plane_with_source(
shared: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
vshard_id: VShardId,
plan: PhysicalPlan,
trace_id: TraceId,
event_source: crate::event::EventSource,
) -> crate::Result<Response> {
dispatch_to_data_plane_inner(
shared,
DataPlaneDispatch {
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
event_source,
txn_id: None,
durability: WalDurability::CallerSupplied {
wal_lsn: None,
resolved_now_ms: None,
},
},
)
.await
}
pub(crate) async fn dispatch_write_to_data_plane(
shared: &SharedState,
write: WriteDispatch,
) -> crate::Result<Response> {
let WriteDispatch {
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
event_source,
txn_id,
wal_lsn,
resolved_now_ms,
} = write;
dispatch_to_data_plane_inner(
shared,
DataPlaneDispatch {
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
event_source,
txn_id,
durability: WalDurability::CallerSupplied {
wal_lsn,
resolved_now_ms,
},
},
)
.await
}
pub(crate) async fn dispatch_autocommit_write(
shared: &SharedState,
write: AutocommitWrite,
) -> crate::Result<Response> {
let AutocommitWrite {
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
event_source,
txn_id,
} = write;
dispatch_to_data_plane_inner(
shared,
DataPlaneDispatch {
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
event_source,
txn_id,
durability: WalDurability::AppendHere { now_override: None },
},
)
.await
}
pub async fn dispatch_to_data_plane_with_txn(
shared: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
vshard_id: VShardId,
plan: PhysicalPlan,
trace_id: TraceId,
txn_id: Option<crate::types::TxnId>,
) -> crate::Result<Response> {
dispatch_to_data_plane_inner(
shared,
DataPlaneDispatch {
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
event_source: crate::event::EventSource::User,
txn_id,
durability: WalDurability::CallerSupplied {
wal_lsn: None,
resolved_now_ms: None,
},
},
)
.await
}
async fn dispatch_to_data_plane_inner(
shared: &SharedState,
params: DataPlaneDispatch,
) -> crate::Result<Response> {
let DataPlaneDispatch {
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
event_source,
txn_id,
durability,
} = params;
let plan = match crate::control::server::exchange::resolve_exchange_in_plan(
shared,
database_id,
tenant_id,
plan,
trace_id,
None,
)
.await?
{
crate::control::server::exchange::Resolved::Gathered(
resp,
_shard_watermarks,
_shuffle_reads,
) => {
return Ok(resp);
}
crate::control::server::exchange::Resolved::Plan(p) => p,
crate::control::server::exchange::Resolved::Stream(s) => {
return crate::control::server::exchange::gather::stream_to_response(s).await;
}
};
submit_write(
shared,
SubmitWrite {
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
event_source,
txn_id,
user_id: None,
durability,
ordering: WriteOrdering::Gate,
change_feed: ChangeFeedOwner::Funnel,
},
)
.await
.map(|outcome| outcome.response)
}