use crate::observability::{HandlerOutcome, StreamingOp};
use serde::{Deserialize, Serialize};
use std::sync::Arc;
use std::time::{Duration, Instant};
use crate::streaming::anchor::AnchorManager;
use crate::streaming::handle::StreamAnchorHandle;
pub const DETECTION_MULTIPLIER: u8 = 3;
fn default_heartbeat_interval_ms() -> u64 {
5_000
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct StreamCancelHandle(u128);
#[derive(Serialize, Deserialize)]
struct StreamCancelHandleWire {
hi: u64,
lo: u64,
}
impl StreamCancelHandle {
pub fn pack(worker_id: velo_ext::WorkerId, stream_id: u64) -> Self {
Self(((worker_id.as_u64() as u128) << 64) | (stream_id as u128))
}
pub fn unpack(self) -> (velo_ext::WorkerId, u64) {
let hi = (self.0 >> 64) as u64;
let lo = self.0 as u64;
(velo_ext::WorkerId::from_u64(hi), lo)
}
}
impl Serialize for StreamCancelHandle {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
StreamCancelHandleWire {
hi: (self.0 >> 64) as u64,
lo: self.0 as u64,
}
.serialize(serializer)
}
}
impl<'de> Deserialize<'de> for StreamCancelHandle {
fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
let wire = StreamCancelHandleWire::deserialize(deserializer)?;
Ok(Self(((wire.hi as u128) << 64) | (wire.lo as u128)))
}
}
#[derive(Debug, Serialize, Deserialize)]
pub struct StreamCancelRequest {
pub sender_stream_id: u64,
}
pub struct SenderEntry {
pub cancel_token: tokio_util::sync::CancellationToken,
pub rx_closer: std::sync::Mutex<Option<flume::Receiver<()>>>,
}
#[derive(Default)]
pub struct SenderRegistry {
pub senders: dashmap::DashMap<u64, SenderEntry>,
}
pub fn create_stream_cancel_handler(
sender_registry: Arc<SenderRegistry>,
) -> crate::messenger::Handler {
crate::messenger::Handler::am_handler(
"_stream_cancel",
move |ctx: crate::messenger::Context| {
let req = serde_json::from_slice::<StreamCancelRequest>(&ctx.payload)?;
if let Some((_, entry)) = sender_registry.senders.remove(&req.sender_stream_id) {
drop(entry.rx_closer.lock().unwrap().take());
entry.cancel_token.cancel();
}
Ok(())
},
)
.build()
}
#[derive(Debug, Serialize, Deserialize)]
pub struct AnchorAttachRequest {
pub handle: StreamAnchorHandle,
pub session_id: u64,
pub stream_cancel_handle: StreamCancelHandle,
#[serde(default)]
pub supported_transport_keys: Vec<velo_ext::TransportKey>,
}
#[derive(Debug, Serialize, Deserialize)]
pub enum AnchorAttachResponse {
Ok {
streaming_transport_key: velo_ext::TransportKey,
#[serde(default = "default_heartbeat_interval_ms")]
heartbeat_interval_ms: u64,
#[serde(default)]
routing_session_id: u64,
#[serde(default)]
initial_credit: u32,
#[serde(default)]
slot_byte_budget: u32,
},
Err { reason: String },
}
#[derive(Debug, Serialize, Deserialize)]
pub struct AnchorDetachRequest {
pub handle: StreamAnchorHandle,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct AnchorFinalizeRequest {
pub handle: StreamAnchorHandle,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct AnchorCancelRequest {
pub handle: StreamAnchorHandle,
}
pub(crate) async fn reader_pump(
transport_rx: flume::Receiver<Vec<u8>>,
frame_tx: flume::Sender<Vec<u8>>,
cancel_token: tokio_util::sync::CancellationToken,
ctx: crate::streaming::anchor::AnchorContext,
local_id: u64,
heartbeat_deadline: Duration,
) {
let crate::streaming::anchor::AnchorContext {
registry,
mpsc_registry,
metrics,
} = ctx;
let mut missed_heartbeats: u8 = 0;
loop {
tokio::select! {
_ = cancel_token.cancelled() => break,
result = tokio::time::timeout(heartbeat_deadline, transport_rx.recv_async()) => {
match result {
Ok(Ok(bytes)) => {
missed_heartbeats = 0;
match frame_tx.try_send(bytes) {
Ok(()) => {}
Err(flume::TrySendError::Full(b)) => {
if let Some(m) = metrics.as_ref() {
m.record_reader_pump_backpressure();
}
if frame_tx.send_async(b).await.is_err() {
break; }
}
Err(flume::TrySendError::Disconnected(_)) => break,
}
}
Ok(Err(_)) => break, Err(_timeout) => {
missed_heartbeats += 1;
if missed_heartbeats >= DETECTION_MULTIPLIER {
if let Some(m) = metrics.as_ref() {
m.record_heartbeat_watchdog_firing();
}
tracing::warn!(
local_id,
anchor_frame_tx_len = frame_tx.len(),
anchor_frame_tx_cap = frame_tx.capacity().unwrap_or_default(),
transport_rx_len = transport_rx.len(),
transport_rx_cap = transport_rx.capacity().unwrap_or_default(),
heartbeat_deadline_ms = heartbeat_deadline.as_millis() as u64,
detection_multiplier = DETECTION_MULTIPLIER,
"reader_pump: heartbeat watchdog fired, injecting Dropped \
(saturation indicator: see velo_streaming_*_backpressure_total)"
);
let dropped_bytes = crate::streaming::sender::cached_dropped().clone();
if frame_tx.try_send(dropped_bytes).is_err() {
tracing::warn!(
local_id,
"reader_pump: anchor channel saturated at watchdog-fire; \
Dropped sentinel could not be injected, consumer will see \
channel close (EOF) -- watchdog firing counter is the \
authoritative signal here"
);
}
if let Some((_, entry)) = registry.remove(&local_id) {
entry.cancel_token.cancel();
crate::streaming::anchor::set_active_anchor_gauge(
metrics.as_ref(),
®istry,
&mpsc_registry,
);
}
break;
}
}
}
}
}
}
cancel_token.cancel();
}
pub fn create_anchor_attach_handler(manager: Arc<AnchorManager>) -> crate::messenger::Handler {
crate::messenger::Handler::typed_unary_async(
"_anchor_attach",
move |ctx: crate::messenger::TypedContext<AnchorAttachRequest>| {
let manager = manager.clone();
async move {
let started = Instant::now();
let req = ctx.input;
if req.handle.is_mpsc_stream() {
manager.record_streaming_operation(
StreamingOp::Attach,
HandlerOutcome::Error,
"unknown",
started,
);
return Ok(AnchorAttachResponse::Err {
reason: format!("anchor {} is mpsc; use _mpsc_anchor_attach", req.handle),
});
}
let (_, local_id) = req.handle.unpack();
{
let entry = manager.registry.get(&local_id);
match entry {
None => {
manager.record_streaming_operation(
StreamingOp::Attach,
HandlerOutcome::Error,
"unknown",
started,
);
return Ok(AnchorAttachResponse::Err {
reason: format!("anchor {} not found", req.handle),
});
}
Some(e) if e.attachment => {
manager.record_streaming_operation(
StreamingOp::Attach,
HandlerOutcome::Error,
"unknown",
started,
);
return Ok(AnchorAttachResponse::Err {
reason: format!("anchor {} already attached", req.handle),
});
}
_ => {} }
}
let routing_session_id = manager
.next_routing_session_id
.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
+ 1;
let selection = manager.select_streaming_transport(&req.supported_transport_keys);
let receiver = match selection.transport.bind(local_id, routing_session_id).await {
Ok(rx) => rx,
Err(e) => {
manager.record_streaming_operation(
StreamingOp::Attach,
HandlerOutcome::Error,
"unknown",
started,
);
return Ok(AnchorAttachResponse::Err {
reason: format!("transport error: {}", e),
});
}
};
let streaming_transport_key = selection.key;
use dashmap::mapref::entry::Entry;
match manager.registry.entry(local_id) {
Entry::Vacant(_) => {
manager.record_streaming_operation(
StreamingOp::Attach,
HandlerOutcome::Error,
"unknown",
started,
);
Ok(AnchorAttachResponse::Err {
reason: format!("anchor {} removed during bind", req.handle),
})
}
Entry::Occupied(mut occ) => {
let entry = occ.get_mut();
if entry.attachment {
manager.record_streaming_operation(
StreamingOp::Attach,
HandlerOutcome::Error,
"unknown",
started,
);
Ok(AnchorAttachResponse::Err {
reason: format!("anchor {} already attached", req.handle),
})
} else {
let pump_cancel = entry.cancel_token.child_token();
entry.active_pump_token = Some(pump_cancel.clone());
let pump_frame_tx = entry.frame_tx.clone();
let heartbeat_interval = entry.heartbeat_interval;
entry.attachment = true;
entry.stream_cancel_handle = Some(req.stream_cancel_handle);
drop(occ);
let (_, local_id) = req.handle.unpack();
tokio::spawn(reader_pump(
receiver, pump_frame_tx, pump_cancel, manager.anchor_context(),
local_id, heartbeat_interval,
));
manager.record_streaming_operation(
StreamingOp::Attach,
HandlerOutcome::Success,
streaming_transport_key.as_str(),
started,
);
Ok(AnchorAttachResponse::Ok {
streaming_transport_key,
heartbeat_interval_ms: heartbeat_interval.as_millis() as u64,
routing_session_id,
initial_credit: selection.initial_credit,
slot_byte_budget: selection.slot_byte_budget,
})
}
}
}
}
},
)
.spawn()
.build()
}
pub fn create_anchor_detach_handler(manager: Arc<AnchorManager>) -> crate::messenger::Handler {
crate::messenger::Handler::typed_unary_async(
"_anchor_detach",
move |ctx: crate::messenger::TypedContext<AnchorDetachRequest>| {
let manager = manager.clone();
async move {
let started = Instant::now();
let req = ctx.input;
let (_, local_id) = req.handle.unpack();
use dashmap::mapref::entry::Entry;
let maybe_entry_info = match manager.registry.entry(local_id) {
Entry::Vacant(_) => None,
Entry::Occupied(mut occ) => {
let entry = occ.get_mut();
entry.attachment = false;
Some((entry.active_pump_token.take(), entry.frame_tx.clone()))
}
};
if let Some((maybe_pump_token, frame_tx)) = maybe_entry_info {
if let Some(pump_token) = maybe_pump_token {
pump_token.cancel();
}
let sentinel_bytes = crate::streaming::sender::cached_detached().clone();
let _ = frame_tx.try_send(sentinel_bytes);
manager.record_streaming_operation(
StreamingOp::Detach,
HandlerOutcome::Success,
"velo",
started,
);
} else {
manager.record_streaming_operation(
StreamingOp::Detach,
HandlerOutcome::Error,
"velo",
started,
);
}
Ok(())
}
},
)
.spawn()
.build()
}
pub fn create_anchor_finalize_handler(manager: Arc<AnchorManager>) -> crate::messenger::Handler {
crate::messenger::Handler::typed_unary_async(
"_anchor_finalize",
move |ctx: crate::messenger::TypedContext<AnchorFinalizeRequest>| {
let manager = manager.clone();
async move {
let started = Instant::now();
let req = ctx.input;
let (_, local_id) = req.handle.unpack();
if let Some(entry) = manager.remove_anchor(local_id) {
let sentinel_bytes = crate::streaming::sender::cached_finalized().clone();
let _ = entry.frame_tx.try_send(sentinel_bytes);
manager.record_streaming_operation(
StreamingOp::Finalize,
HandlerOutcome::Success,
"velo",
started,
);
} else {
manager.record_streaming_operation(
StreamingOp::Finalize,
HandlerOutcome::Error,
"velo",
started,
);
}
Ok(())
}
},
)
.spawn()
.build()
}
pub fn create_anchor_cancel_handler(manager: Arc<AnchorManager>) -> crate::messenger::Handler {
crate::messenger::Handler::typed_unary_async(
"_anchor_cancel",
move |ctx: crate::messenger::TypedContext<AnchorCancelRequest>| {
let manager = manager.clone();
async move {
let started = Instant::now();
let req = ctx.input;
let (_, local_id) = req.handle.unpack();
if let Some(entry) = manager.remove_anchor(local_id) {
entry.cancel_token.cancel();
manager.record_streaming_operation(
StreamingOp::Cancel,
HandlerOutcome::Success,
"velo",
started,
);
} else {
manager.record_streaming_operation(
StreamingOp::Cancel,
HandlerOutcome::Error,
"velo",
started,
);
}
Ok(())
}
},
)
.spawn()
.build()
}
#[cfg(test)]
mod tests;