Skip to content

Commit 629f082

Browse files
srperensPer Enstedtclaude
authored
fix: leak-fixes — stop pipeline on delete + overlay renderer cleanup (#530)
* fix: release pipeline resources on flow delete and overlay teardown - delete_flow now stops the active pipeline before removing the flow record. Previously the PipelineManager stayed in AppState.pipelines after the flow record was gone, with no API path left to reach it. - stop_flow now also unregisters the vision-mixer overlay renderer alongside the overlay state. Both share the same lifecycle and must be cleared together; otherwise the renderer entry keeps a background thread alive past pipeline shutdown. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * test: cover delete-while-running and overlay-registry teardown - test_delete_running_flow_releases_pipeline starts a flow, captures weak refs to the pipeline and its elements, calls delete_flow without a prior stop, and asserts the pipeline is removed from AppState.pipelines and that no element survives. - overlay_registries_round_trip exercises the register/unregister pair for both overlay_states and overlay_renderers, locking in the shared-lifecycle invariant so a future per-block registry isn't forgotten in stop_flow. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * fix: log delete-time stop_flow failure as error A failed stop_flow during delete_flow leaves the PipelineManager orphaned (the very leak the prior commit fixes), so it deserves error-level visibility, not warn. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Per Enstedt <per.enstedt@svt.se> Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
1 parent cb69484 commit 629f082

3 files changed

Lines changed: 174 additions & 0 deletions

File tree

backend/src/blocks/builtin/vision_mixer/tests.rs

Lines changed: 72 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,3 +111,75 @@ fn test_parse_initial_pgm_clamped() {
111111
props.insert("initial_pgm_input".to_string(), PropertyValue::UInt(99));
112112
assert_eq!(properties::parse_initial_pgm(&props, 4), 3); // max index = 3
113113
}
114+
115+
/// `overlay_states` and `overlay_renderers` share the same lifecycle:
116+
/// they are populated together in `build_overlay` and must be cleared
117+
/// together by the cleanup branch in `state.rs::stop_flow`. If you add
118+
/// another per-block registry in this module, mirror its unregister call
119+
/// in that branch and extend this test.
120+
#[test]
121+
fn overlay_registries_round_trip() {
122+
use super::overlay::{
123+
get_overlay_renderer, get_overlay_state, register_overlay_renderer, register_overlay_state,
124+
unregister_overlay_renderer, unregister_overlay_state, OverlayRenderer,
125+
VisionMixerOverlayState,
126+
};
127+
use gstreamer as gst;
128+
use gstreamer_app as gst_app;
129+
use std::sync::{Arc, Mutex};
130+
131+
gst::init().unwrap();
132+
133+
let block_id = "test-vm-overlay-cleanup-block-id";
134+
135+
let lo = layout::compute_layout(1280, 720, 4);
136+
let state = Arc::new(VisionMixerOverlayState::new(
137+
4,
138+
0,
139+
0,
140+
1,
141+
vec!["A".into(), "B".into(), "C".into(), "D".into()],
142+
lo,
143+
false,
144+
));
145+
146+
let caps = gst::Caps::builder("video/x-raw")
147+
.field("format", "BGRA")
148+
.field("width", 1280i32)
149+
.field("height", 720i32)
150+
.field("framerate", gst::Fraction::new(50, 1))
151+
.build();
152+
let appsrc = gst_app::AppSrc::builder().caps(&caps).build();
153+
154+
let renderer = Arc::new(Mutex::new(OverlayRenderer::new(
155+
appsrc,
156+
caps,
157+
Arc::clone(&state),
158+
1280,
159+
720,
160+
)));
161+
162+
register_overlay_state(block_id, Arc::clone(&state));
163+
register_overlay_renderer(block_id, Arc::clone(&renderer));
164+
165+
assert!(
166+
get_overlay_state(block_id).is_some(),
167+
"state should be registered"
168+
);
169+
assert!(
170+
get_overlay_renderer(block_id).is_some(),
171+
"renderer should be registered"
172+
);
173+
174+
unregister_overlay_state(block_id);
175+
unregister_overlay_renderer(block_id);
176+
177+
assert!(
178+
get_overlay_state(block_id).is_none(),
179+
"state must be cleaned (otherwise API still sees stale block)"
180+
);
181+
assert!(
182+
get_overlay_renderer(block_id).is_none(),
183+
"renderer must be cleaned (otherwise overlay-timer-* thread leaks)"
184+
);
185+
}

backend/src/state.rs

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -631,6 +631,22 @@ impl AppState {
631631
return Ok(false);
632632
}
633633

634+
// Stop the pipeline first if it is still running. Without this, the
635+
// PipelineManager stays in self.inner.pipelines after the flow record
636+
// is gone, with no API path left to reach it.
637+
let pipeline_active = {
638+
let pipelines = self.inner.pipelines.read().await;
639+
pipelines.contains_key(id)
640+
};
641+
if pipeline_active {
642+
if let Err(e) = self.stop_flow(id).await {
643+
error!(
644+
"Failed to stop flow {} before delete: {} — pipeline resources may leak",
645+
id, e
646+
);
647+
}
648+
}
649+
634650
// Delete from storage first (skip for ephemeral flows)
635651
let is_ephemeral = {
636652
let flows = self.inner.flows.read().await;
@@ -1365,6 +1381,12 @@ impl AppState {
13651381
crate::blocks::builtin::vision_mixer::overlay::unregister_overlay_state(
13661382
&block.id,
13671383
);
1384+
// Without this, the overlay-timer-* thread keeps polling the
1385+
// renderer registry, holds a strong AppSrc ref, and prevents
1386+
// the pipeline (and its NiceAgent) from finalizing.
1387+
crate::blocks::builtin::vision_mixer::overlay::unregister_overlay_renderer(
1388+
&block.id,
1389+
);
13681390
}
13691391
}
13701392
}

backend/tests/pipeline_lifecycle_test.rs

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@ use std::collections::HashMap;
88
use strom::blocks::BlockRegistry;
99
use strom::events::EventBroadcaster;
1010
use strom::gst::pipeline::PipelineManager;
11+
use strom::state::AppState;
12+
use strom::storage::JsonFileStorage;
1113
use strom_types::{Flow, Link};
1214
use tempfile::NamedTempFile;
1315

@@ -162,3 +164,81 @@ async fn test_leak_detection_catches_circular_reference() {
162164
leak detection would miss real leaks!"
163165
);
164166
}
167+
168+
/// Regression test for the delete-without-stop leak: deleting a flow that has
169+
/// an active pipeline must fully release the pipeline (encoders, sockets,
170+
/// element refs). Before the fix in delete_flow, the PipelineManager stayed
171+
/// in AppState.pipelines after deletion and orphaned NVENC sessions / sockets
172+
/// would accumulate until process exit.
173+
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
174+
async fn test_delete_running_flow_releases_pipeline() {
175+
gstreamer::init().unwrap();
176+
177+
let storage_file = NamedTempFile::new().unwrap();
178+
let blocks_file = NamedTempFile::new().unwrap();
179+
let storage = JsonFileStorage::new(storage_file.path());
180+
181+
let state = AppState::new(
182+
storage,
183+
blocks_file.path(),
184+
std::env::temp_dir(),
185+
vec![],
186+
"all".to_string(),
187+
vec![],
188+
);
189+
190+
let flow = build_test_flow("delete_running_flow_test");
191+
let flow_id = flow.id;
192+
state.upsert_flow(flow).await.expect("upsert_flow failed");
193+
194+
state.start_flow(&flow_id).await.expect("start_flow failed");
195+
196+
// Capture weak refs to the running pipeline before deletion
197+
let (pipeline_weak, element_weak_refs) = {
198+
let pipelines = state.pipelines_read().await;
199+
let manager = pipelines
200+
.get(&flow_id)
201+
.expect("Pipeline should be active after start_flow");
202+
(manager.pipeline_weak(), manager.element_weak_refs())
203+
};
204+
assert!(
205+
pipeline_weak.upgrade().is_some(),
206+
"Pipeline should be alive before delete"
207+
);
208+
assert!(
209+
!element_weak_refs.is_empty(),
210+
"Pipeline should have elements"
211+
);
212+
213+
// Delete WITHOUT calling stop_flow first — this is the path the UI takes
214+
let deleted = state
215+
.delete_flow(&flow_id)
216+
.await
217+
.expect("delete_flow failed");
218+
assert!(deleted, "delete_flow should report the flow was deleted");
219+
220+
// Pipeline must no longer be tracked in AppState
221+
{
222+
let pipelines = state.pipelines_read().await;
223+
assert!(
224+
!pipelines.contains_key(&flow_id),
225+
"Pipeline still in AppState.pipelines after delete — orphaned manager"
226+
);
227+
}
228+
229+
// And every GStreamer object must be finalized — surviving refs would
230+
// mean a leaked encoder / socket / thread that no API path can release.
231+
assert!(
232+
pipeline_weak.upgrade().is_none(),
233+
"Pipeline still alive after delete_flow — resources leaked"
234+
);
235+
let leaked: Vec<_> = element_weak_refs
236+
.iter()
237+
.filter_map(|(name, weak)| weak.upgrade().map(|_| name.clone()))
238+
.collect();
239+
assert!(
240+
leaked.is_empty(),
241+
"Elements still alive after delete_flow: {:?}",
242+
leaked
243+
);
244+
}

0 commit comments

Comments
 (0)