From 63903e23cf0c06b3c9dd2bce6df0e4d886b9ec05 Mon Sep 17 00:00:00 2001 From: Leonid Ryzhyk Date: Tue, 28 Jul 2026 11:42:46 -0700 Subject: [PATCH 1/2] adapters: republish completion counters when an output endpoint goes away `total_completed_records` and `total_completed_steps` are the minimum over the registered output endpoints. An endpoint that has not delivered a step holds both back, and `disconnect_output` drops an endpoint along with everything still queued for it, so a departing endpoint never reports the progress the others already have. `remove_output` only dropped the map entry: the only other refresh happens on a step or an output batch, and an idle pipeline produces neither, so the minimum stayed pinned at the departed endpoint's last step forever. `/completion_status` then reports `inprogress` for a pipeline that has delivered all of its output, and a checkpoint waiting on `total_completed_records` in `CheckpointThread::run` never unblocks. `server::test_with_kafka::test_server` hit this: it opens 200 `/egress` streams and drops each client instantly, and ~39 of them were still torn down in the window where the circuit ran the steps carrying the final output. Failure: https://github.com/feldera/feldera/actions/runs/30380682951/job/90350576953 Signed-off-by: Leonid Ryzhyk --- crates/adapters/src/controller/stats.rs | 9 +++ crates/adapters/src/controller/test.rs | 84 +++++++++++++++++++++++++ 2 files changed, 93 insertions(+) diff --git a/crates/adapters/src/controller/stats.rs b/crates/adapters/src/controller/stats.rs index f2208974978..9138364acef 100644 --- a/crates/adapters/src/controller/stats.rs +++ b/crates/adapters/src/controller/stats.rs @@ -710,6 +710,15 @@ impl ControllerStatus { pub fn remove_output(&self, endpoint_id: &EndpointId) { self.outputs.write().remove(endpoint_id); + + // `total_completed_records` and `total_completed_steps` are the minimum + // over the registered output endpoints, and this endpoint may have been + // the one holding them back: it is dropped along with everything still + // queued for it, so it will never report the progress the others + // already have. Republish the counters now, since the only other + // refresh happens on a step or an output batch, and an idle pipeline + // produces neither. + self.update_total_completed_records(None); } /// Initialize stats for a new output endpoint. diff --git a/crates/adapters/src/controller/test.rs b/crates/adapters/src/controller/test.rs index ecb570fb2f3..24745d020a4 100644 --- a/crates/adapters/src/controller/test.rs +++ b/crates/adapters/src/controller/test.rs @@ -5026,6 +5026,90 @@ fn test_output_progress_counter_waits_for_owed_output() { assert_eq!(processed(1), RESTORED_RECORDS); } +/// Dropping an output endpoint that is behind must republish the pipeline's +/// completion counters. +/// +/// `total_completed_records` and `total_completed_steps` are the minimum over the +/// registered output endpoints, so an endpoint that has not delivered a step holds +/// both back. An abandoned `/egress` stream is disconnected with its queue still +/// unread, which means the endpoint disappears without ever reaching that step. If +/// removal does not recompute the minimum, nothing else does until the next step, +/// and an idle pipeline runs no further step: `/completion_status` then reports +/// `inprogress` forever and a checkpoint waiting on `total_completed_records` +/// never unblocks. +#[test] +fn test_removing_a_lagging_output_endpoint_republishes_completion() { + use super::stats::ProcessedRecords; + use crate::{ControllerStatus, OutputEndpointConfig}; + use uuid::Uuid; + + // Records the circuit processes in the one step of this test. + const STEP_RECORDS: u64 = 10; + + let config = serde_json::from_value(json!({ + "name": "test_remove_lagging_output", + "workers": 1, + })) + .unwrap(); + let status = ControllerStatus::new(config, 0, None, Uuid::nil()); + + let output_config: OutputEndpointConfig = serde_json::from_value(json!({ + "stream": "v1", + "transport": { "name": "http_output", "config": {} }, + "format": { "name": "json", "config": {} } + })) + .unwrap(); + + // Two `/egress` streams of the same view: one client reads, the other walks + // away. + status.add_output(&0, "reader", &output_config, None, true); + status.add_output(&1, "abandoned", &output_config, None, true); + + // The circuit initiates and evaluates one step, and its output is queued for + // both endpoints. + status + .global_metrics + .total_initiated_steps + .store(1, Ordering::Release); + status.processed_data(BufferSize { + records: STEP_RECORDS as usize, + bytes: 0, + }); + status.enqueue_batch(0, STEP_RECORDS as usize); + status.enqueue_batch(1, STEP_RECORDS as usize); + + // Only `reader` transmits the batch. + let parker = Parker::new(); + let unparker = parker.unparker().clone(); + let step_processed = ProcessedRecords { + total_processed_input_records: STEP_RECORDS, + total_processed_steps: 1, + }; + status.output_batch(0, Some(step_processed), STEP_RECORDS as usize, &unparker); + + assert_eq!( + status.global_metrics.total_completed_steps(), + 0, + "the step is not complete while `abandoned` still owes its output" + ); + assert_eq!(status.num_total_completed_records(), 0); + + // The abandoned client disconnects, so the endpoint goes away with the batch + // still queued. The step is now complete as far as any endpoint is concerned. + status.remove_output(&1); + + assert_eq!( + status.global_metrics.total_completed_steps(), + 1, + "removing the endpoint that owed the output must complete the step" + ); + assert_eq!( + status.num_total_completed_records(), + STEP_RECORDS, + "removing the endpoint that owed the output must complete its records" + ); +} + /// Only a changed relation makes an output endpoint fall behind. A connector whose /// own definition changed still emits nothing for inputs already processed, so it /// stays caught up and keeps its seeded progress counter. From 824b05aaef79adf4ac14e6749e9b7f70f1cef002 Mon Sep 17 00:00:00 2001 From: Leonid Ryzhyk Date: Tue, 28 Jul 2026 11:42:51 -0700 Subject: [PATCH 2/2] adapters: report which connector stalled a completion token in tests `wait_for_completion` unwrapped the timeout, so a stall said only that 20 seconds had passed. Dump the completion status and `/stats` first: the step the token waits for, next to each connector's `total_processed_steps`, names the endpoint that is holding `total_completed_steps` back. Signed-off-by: Leonid Ryzhyk --- crates/adapters/src/server.rs | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/crates/adapters/src/server.rs b/crates/adapters/src/server.rs index 73ce742e417..c2d957c46d0 100644 --- a/crates/adapters/src/server.rs +++ b/crates/adapters/src/server.rs @@ -4021,7 +4021,7 @@ outputs: } pub(super) async fn wait_for_completion(server: &TestServer, token: &str) { - async_wait( + let completed = async_wait( || async { let CompletionStatusResponse { status, .. } = get_completion_status(server, token).await; @@ -4030,7 +4030,16 @@ outputs: 20_000, ) .await - .unwrap(); + .is_ok(); + + if !completed { + // The step the token waits for, next to each connector's + // `total_processed_steps`, names the endpoint that is holding + // `total_completed_steps` back. + let status = get_completion_status(server, token).await; + print_stats(server).await; + panic!("completion token '{token}' is still {status:?} after 20 seconds"); + } } pub(super) fn test_batches(start: u32, len: u32) -> Vec> {