use std::pin::Pin;
use tokio::sync::{broadcast, mpsc};
use tokio_stream::Stream;
use tonic::Status;
use crate::ir::LogicalFilter;
use crate::proto::udb::core::livequery::services::v1 as lq_pb;
use super::budget::StreamSlot;
use super::errors::livequery_backpressure_status;
use super::predicate::{
change_frame, change_row, event_matches_tenant_scope, filter_matches_row, topic_matches_source,
};
pub(crate) type LiveQueryStream =
Pin<Box<dyn Stream<Item = Result<lq_pb::SubscribeResponse, Status>> + Send + 'static>>;
pub(crate) async fn run_delta_forward(
mut rx: broadcast::Receiver<crate::cdc::CdcEnvelope>,
tx: mpsc::Sender<Result<lq_pb::SubscribeResponse, Status>>,
tenant_id: String,
project_id: String,
cdc_topic: String,
user_filter: Option<LogicalFilter>,
_stream_slot: StreamSlot,
) {
loop {
let received = tokio::select! {
_ = tx.closed() => break,
received = rx.recv() => received,
};
match received {
Ok(envelope) => {
if !topic_matches_source(&envelope.topic, &cdc_topic) {
continue;
}
let payload =
match serde_json::from_str::<serde_json::Value>(&envelope.payload_json) {
Ok(value) => value,
Err(_) => continue,
};
if !event_matches_tenant_scope(&envelope.topic, &payload, &tenant_id, &project_id) {
continue;
}
let row = change_row(&payload);
if let Some(filter) = user_filter.as_ref() {
if !filter_matches_row(filter, &row) {
continue;
}
}
match tx.try_send(Ok(change_frame(&envelope, &payload))) {
Ok(()) => {}
Err(mpsc::error::TrySendError::Full(_)) => {
let _ = tx
.send(Err(livequery_backpressure_status(
"subscriber_channel",
"live query subscriber too slow; stream closed",
)))
.await;
break;
}
Err(mpsc::error::TrySendError::Closed(_)) => break,
}
}
Err(broadcast::error::RecvError::Lagged(_)) => {
let _ = tx
.send(Err(livequery_backpressure_status(
"delta feed lag",
"live query delta feed lagged; stream closed",
)))
.await;
break;
}
Err(broadcast::error::RecvError::Closed) => break,
}
}
}