Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

34 changes: 16 additions & 18 deletions binaries/coordinator/src/control_requests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,8 @@ use crate::{
ensure_remove_node_applied, ensure_replace_node_applied, ensure_set_param_forward_applied,
handle_get_trace_spans, handle_get_traces,
handlers::{
build_dataflow, dataflow_result, handle_destroy, parse_logs_node_id, reload_dataflow,
resolve_name, restart_node, retrieve_logs, send_log_message, start_dataflow, stop_dataflow,
attach_log_subscriber, build_dataflow, dataflow_result, handle_destroy, parse_logs_node_id,
reload_dataflow, resolve_name, restart_node, retrieve_logs, start_dataflow, stop_dataflow,
stop_node,
},
initiate_restart,
Expand Down Expand Up @@ -514,14 +514,13 @@ impl Coordinator {
found_tx,
} => {
if let Some(dataflow) = self.running_dataflows.get_mut(&dataflow_id) {
dataflow
.log_subscribers
.push(LogSubscriber::new(level, sender));
let buffered = std::mem::take(&mut dataflow.buffered_log_messages);
for message in buffered {
send_log_message(&mut dataflow.log_subscribers, &message).await;
}
let _ = found_tx.send(true);
attach_log_subscriber(
&mut dataflow.log_subscribers,
&mut dataflow.buffered_log_messages,
LogSubscriber::new(level, sender),
found_tx,
)
.await;
} else if self.archived_dataflows.contains_key(&dataflow_id) {
// Dataflow already finished before the CLI could subscribe.
// Acknowledge the subscription so the CLI doesn't error, then
Expand All @@ -538,14 +537,13 @@ impl Coordinator {
found_tx,
} => {
if let Some(build) = self.running_builds.get_mut(&build_id) {
build
.log_subscribers
.push(LogSubscriber::new(level, sender));
let buffered = std::mem::take(&mut build.buffered_log_messages);
for message in buffered {
send_log_message(&mut build.log_subscribers, &message).await;
}
let _ = found_tx.send(true);
attach_log_subscriber(
&mut build.log_subscribers,
&mut build.buffered_log_messages,
LogSubscriber::new(level, sender),
found_tx,
)
.await;
} else {
let _ = found_tx.send(false);
}
Expand Down
58 changes: 58 additions & 0 deletions binaries/coordinator/src/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,27 @@ pub(crate) fn resolve_name(
}
}

/// Register `subscriber` and replay the logs buffered before it attached.
///
/// `found_tx` is answered *before* the replay. The WS task that owns the
/// subscriber's receiver waits on it and only resumes draining that receiver
/// once it fires, so replaying first let the bounded channel fill up: every
/// further message then took `send_log_message`'s 100 ms timeout, which
/// stalled the coordinator's event loop and finally evicted the subscriber
/// with the rest of the backlog lost.
pub(crate) async fn attach_log_subscriber(
log_subscribers: &mut Vec<LogSubscriber>,
buffered_log_messages: &mut Vec<LogMessage>,
subscriber: LogSubscriber,
found_tx: tokio::sync::oneshot::Sender<bool>,
) {
log_subscribers.push(subscriber);
let _ = found_tx.send(true);
for message in std::mem::take(buffered_log_messages) {
send_log_message(log_subscribers, &message).await;
}
}

pub(crate) async fn send_log_message(
log_subscribers: &mut Vec<LogSubscriber>,
message: &LogMessage,
Expand Down Expand Up @@ -1221,6 +1242,43 @@ mod tests {
);
}

/// Attaching to a dataflow that buffered more logs than the subscriber's
/// channel holds must deliver the whole backlog. The WS task only drains
/// the channel after `found_tx` fires, so answering it after the replay
/// left the channel full: the replay hit 100 send timeouts, stalled the
/// event loop and evicted the subscriber with most of the backlog lost.
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn attach_log_subscriber_delivers_backlog_larger_than_channel() {
const BACKLOG: usize = 200;
let (tx, mut rx) = tokio::sync::mpsc::channel(4);
let (found_tx, found_rx) = tokio::sync::oneshot::channel::<bool>();

// Mirrors `ws_control`: wait for the subscribe answer, then drain.
let drain = tokio::spawn(async move {
assert!(found_rx.await.expect("found_tx dropped"));
let mut received = 0;
while rx.recv().await.is_some() {
received += 1;
}
received
});

let mut subscribers = Vec::new();
let mut buffered = vec![test_log_message(); BACKLOG];
attach_log_subscriber(
&mut subscribers,
&mut buffered,
LogSubscriber::new(log::LevelFilter::Info, tx),
found_tx,
)
.await;

assert!(buffered.is_empty(), "the backlog is handed over once");
assert_eq!(subscribers.len(), 1, "subscriber must not be evicted");
drop(subscribers); // close the channel so the drain task finishes
assert_eq!(drain.await.unwrap(), BACKLOG);
}

/// A topic subscriber that got closed (its CLI went away) is evicted when
/// the next frame arrives; its daemon streams must be stopped at the same
/// time. Evicting it silently left the daemons streaming until the
Expand Down
Loading