Skip to content

Commit 9c07c83

Browse files
authored
Set src element on inference bus messages (#54)
The video inference element posted its Element messages with gst::message::Element::new, which leaves the message src unset. Bus consumers rely on src to attribute each message to the element that produced it, so unattributed messages can't be matched back to this element on shared buses with multiple inference variants. Build the messages with Element::builder(s).src(&*self.obj()).build() at every post site in the video element (inference results, error, visual-anomaly, object-tracking, and classification fallbacks) so each message carries its producing element. Add an e2e regression test asserting that inference element messages expose a non-empty src.
1 parent cd771c3 commit 9c07c83

2 files changed

Lines changed: 45 additions & 6 deletions

File tree

src/video/imp.rs

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -898,7 +898,9 @@ impl BaseTransformImpl for EdgeImpulseVideoInfer {
898898
inbuf.pts().unwrap_or(gst::ClockTime::ZERO),
899899
e.to_string(),
900900
);
901-
let _ = self.obj().post_message(gst::message::Element::new(s));
901+
let _ = self
902+
.obj()
903+
.post_message(gst::message::Element::builder(s).src(&*self.obj()).build());
902904
// Put the model back in the state
903905
let mut state = self.state.lock().unwrap();
904906
state.model = Some(model);
@@ -1061,7 +1063,9 @@ impl BaseTransformImpl for EdgeImpulseVideoInfer {
10611063
elapsed.as_millis() as u32,
10621064
resize_time_ms,
10631065
);
1064-
let _ = self.obj().post_message(gst::message::Element::new(s));
1066+
let _ = self
1067+
.obj()
1068+
.post_message(gst::message::Element::builder(s).src(&*self.obj()).build());
10651069
}
10661070
}
10671071

@@ -1148,7 +1152,9 @@ impl BaseTransformImpl for EdgeImpulseVideoInfer {
11481152
elapsed.as_millis() as u32,
11491153
resize_time_ms,
11501154
);
1151-
let _ = self.obj().post_message(gst::message::Element::new(s));
1155+
let _ = self
1156+
.obj()
1157+
.post_message(gst::message::Element::builder(s).src(&*self.obj()).build());
11521158
}
11531159
}
11541160

@@ -1229,7 +1235,9 @@ impl BaseTransformImpl for EdgeImpulseVideoInfer {
12291235
elapsed.as_millis() as u32,
12301236
resize_time_ms,
12311237
);
1232-
let _ = self.obj().post_message(gst::message::Element::new(s));
1238+
let _ = self
1239+
.obj()
1240+
.post_message(gst::message::Element::builder(s).src(&*self.obj()).build());
12331241
} else {
12341242
// Classification fallback if bounding_boxes is empty
12351243
if let Some(classification) = result_value
@@ -1270,7 +1278,9 @@ impl BaseTransformImpl for EdgeImpulseVideoInfer {
12701278
elapsed.as_millis() as u32,
12711279
resize_time_ms,
12721280
);
1273-
let _ = self.obj().post_message(gst::message::Element::new(s));
1281+
let _ = self
1282+
.obj()
1283+
.post_message(gst::message::Element::builder(s).src(&*self.obj()).build());
12741284
}
12751285
} else {
12761286
// No bounding_boxes field, treat as classification
@@ -1312,7 +1322,9 @@ impl BaseTransformImpl for EdgeImpulseVideoInfer {
13121322
elapsed.as_millis() as u32,
13131323
resize_time_ms,
13141324
);
1315-
let _ = self.obj().post_message(gst::message::Element::new(s));
1325+
let _ = self
1326+
.obj()
1327+
.post_message(gst::message::Element::builder(s).src(&*self.obj()).build());
13161328
}
13171329

13181330
// Put the model back in the state

tests/e2e.rs

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,7 @@ fn run_pipeline_for(pipeline_str: &str, duration: Duration) -> TestResults {
8282
"result": serde_json::from_str::<serde_json::Value>(&result_json)
8383
.unwrap_or(serde_json::Value::Null),
8484
"timing_ms": timing_ms,
85+
"src": msg.src().map(|s| s.name().to_string()),
8586
});
8687
res.inference_messages.push(entry);
8788
res.buffer_count += 1;
@@ -200,6 +201,32 @@ fn test_video_inference_message_structure() {
200201
);
201202
}
202203

204+
/// Regression: inference element messages must carry their `src` element.
205+
///
206+
/// GStreamer element messages should identify the element that posted them so
207+
/// bus consumers can attribute each message to its producing element. The
208+
/// inference result messages are posted via `gst::message::Element` and must be
209+
/// built with `src` set to this element; posting without a `src` leaves the
210+
/// message unattributed on the bus.
211+
#[test]
212+
fn test_inference_messages_carry_src_element() {
213+
let src = video_test_source(5);
214+
let pipeline = format!("{src} ! edgeimpulsevideoinfer ! fakesink");
215+
216+
let results = run_pipeline_for(&pipeline, Duration::from_secs(10));
217+
218+
assert!(
219+
!results.inference_messages.is_empty(),
220+
"Expected at least one inference bus message"
221+
);
222+
for msg in &results.inference_messages {
223+
assert!(
224+
msg.get("src").and_then(|v| v.as_str()).is_some(),
225+
"Inference element message must carry its src element. Message: {msg}"
226+
);
227+
}
228+
}
229+
203230
#[test]
204231
fn test_overlay_does_not_crash() {
205232
let src = video_test_source(5);

0 commit comments

Comments
 (0)