rust/bt-common: Add debug events and background tasks

Add event streaming, dynamic background stream management, and a debug
harness to bt-common.

Debug tools often need to continuously run background work and deliver
messages and/or structured events, instead of only doing work
only when commands are sent.

Add Event to the CommandRunner trait and a method to set persistent
streams which will be polled in the background and can deliver messages
or structured events using RunnerEvent.  Events are optional but each
CommandRunner implementation myst at least define type Event = ()

RunnerBackground<E> manages dynamic keyed or anonymous background streams,
display messages, and structured events.

Test: cargo test --workspace
Change-Id: I2200a0d21a679d0e2f6d380c5156213c6a6a6964
Reviewed-on: https://bluetooth-review.googlesource.com/c/bluetooth/+/4141
diff --git a/rust/bt-ascs/src/debug.rs b/rust/bt-ascs/src/debug.rs
index 0ce995f..ae0312a 100644
--- a/rust/bt-ascs/src/debug.rs
+++ b/rust/bt-ascs/src/debug.rs
@@ -149,6 +149,7 @@
 {
     type Set = AscsCmd;
     type Error = crate::types::Error;
+    type Event = ();
 
     fn run_command(
         &self,
diff --git a/rust/bt-broadcast-assistant/src/debug.rs b/rust/bt-broadcast-assistant/src/debug.rs
index a737f08..6afa9eb 100644
--- a/rust/bt-broadcast-assistant/src/debug.rs
+++ b/rust/bt-broadcast-assistant/src/debug.rs
@@ -209,6 +209,7 @@
 {
     type Set = AssistantCmd;
     type Error = crate::assistant::Error;
+    type Event = ();
 
     fn run_command(
         &self,
diff --git a/rust/bt-common/src/debug_command.rs b/rust/bt-common/src/debug_command.rs
index dd4ea63..5d55a9f 100644
--- a/rust/bt-common/src/debug_command.rs
+++ b/rust/bt-common/src/debug_command.rs
@@ -4,9 +4,15 @@
 
 //! Debug command traits and helpers for defining commands for integration
 //! into a debug tool.
+use futures::channel::mpsc::{unbounded, UnboundedReceiver, UnboundedSender};
+use futures::task::AtomicWaker;
+use futures::{FutureExt, Stream, StreamExt};
 use log::LevelFilter;
+use std::collections::HashMap;
+use std::pin::Pin;
 use std::str::FromStr;
-use std::sync::Arc;
+use std::sync::{Arc, Mutex};
+use std::task::{Context, Poll};
 
 #[must_use = "if unused the previous log level will immediately be restored"]
 pub struct ScopedVerbosityGuard {
@@ -265,10 +271,75 @@
     }
 }
 
+/// An event emitted by a [`CommandRunner`] through its
+/// [`CommandRunner::event_stream`].
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub enum RunnerEvent<E = ()> {
+    /// An arbitrary human-readable message from the runner or its background
+    /// tasks to be displayed by the harness.
+    Message(String),
+    /// A structured domain event specific to this runner that the harness or a
+    /// parent composite tool can react to programmatically.
+    Event(E),
+}
+
+impl<E> RunnerEvent<E> {
+    pub fn message(&self) -> Option<&str> {
+        match self {
+            Self::Message(s) => Some(s),
+            Self::Event(_) => None,
+        }
+    }
+
+    pub fn event(&self) -> Option<&E> {
+        match self {
+            Self::Message(_) => None,
+            Self::Event(e) => Some(e),
+        }
+    }
+
+    pub fn into_event(self) -> Option<E> {
+        match self {
+            Self::Message(_) => None,
+            Self::Event(e) => Some(e),
+        }
+    }
+
+    pub fn into_message(self) -> Option<String> {
+        match self {
+            Self::Message(s) => Some(s),
+            Self::Event(_) => None,
+        }
+    }
+
+    pub fn map_event<F, T>(self, f: F) -> RunnerEvent<T>
+    where
+        F: FnOnce(E) -> T,
+    {
+        match self {
+            Self::Message(m) => RunnerEvent::Message(m),
+            Self::Event(e) => RunnerEvent::Event(f(e)),
+        }
+    }
+}
+
+impl<E> From<String> for RunnerEvent<E> {
+    fn from(s: String) -> Self {
+        Self::Message(s)
+    }
+}
+
+impl<E> From<&str> for RunnerEvent<E> {
+    fn from(s: &str) -> Self {
+        Self::Message(s.to_string())
+    }
+}
+
 /// CommandRunner is used to perform a specific task based on the command set.
 pub trait CommandRunner {
     type Set: CommandSet;
     type Error: ::std::error::Error;
+    type Event;
 
     fn run_command(
         &self,
@@ -276,6 +347,15 @@
         args: Vec<String>,
     ) -> impl futures::Future<Output = Result<(), Self::Error>>;
 
+    /// Returns a stream that drives background processing for this runner and
+    /// yields display messages or structured events for the harness.
+    ///
+    /// Crates that do not need background tasks or event delivery can set
+    /// `type Event = ();` and rely on this default implementation.
+    fn event_stream(&self) -> impl Stream<Item = RunnerEvent<Self::Event>> + '_ {
+        futures::stream::pending::<RunnerEvent<Self::Event>>()
+    }
+
     fn run(
         &self,
         cmd: impl Into<CliCommand<Self::Set>>,
@@ -301,6 +381,7 @@
 impl<R: CommandRunner> CommandRunner for &R {
     type Set = R::Set;
     type Error = R::Error;
+    type Event = R::Event;
 
     fn run_command(
         &self,
@@ -309,11 +390,16 @@
     ) -> impl futures::Future<Output = Result<(), Self::Error>> {
         (**self).run_command(cmd, args)
     }
+
+    fn event_stream(&self) -> impl Stream<Item = RunnerEvent<Self::Event>> + '_ {
+        (**self).event_stream()
+    }
 }
 
 impl<R: CommandRunner> CommandRunner for Arc<R> {
     type Set = R::Set;
     type Error = R::Error;
+    type Event = R::Event;
 
     fn run_command(
         &self,
@@ -322,6 +408,309 @@
     ) -> impl futures::Future<Output = Result<(), Self::Error>> {
         (**self).run_command(cmd, args)
     }
+
+    fn event_stream(&self) -> impl Stream<Item = RunnerEvent<Self::Event>> + '_ {
+        (**self).event_stream()
+    }
+}
+
+type BoxedStream<E> = Pin<Box<dyn Stream<Item = RunnerEvent<E>> + Send + 'static>>;
+
+struct StreamsState<E> {
+    keyed: HashMap<String, BoxedStream<E>>,
+    anonymous: Vec<BoxedStream<E>>,
+}
+
+impl<E> Default for StreamsState<E> {
+    fn default() -> Self {
+        Self { keyed: HashMap::new(), anonymous: Vec::new() }
+    }
+}
+
+impl<E: Send + 'static> Stream for StreamsState<E> {
+    type Item = RunnerEvent<E>;
+
+    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
+        let mut ready = None;
+        self.keyed.retain(|_key, stream| {
+            if ready.is_some() {
+                return true;
+            }
+            match stream.as_mut().poll_next(cx) {
+                Poll::Ready(Some(item)) => {
+                    ready = Some(item);
+                    true
+                }
+                Poll::Ready(None) => false,
+                Poll::Pending => true,
+            }
+        });
+        if let Some(item) = ready {
+            return Poll::Ready(Some(item));
+        }
+
+        self.anonymous.retain_mut(|stream| {
+            if ready.is_some() {
+                return true;
+            }
+            match stream.as_mut().poll_next(cx) {
+                Poll::Ready(Some(item)) => {
+                    ready = Some(item);
+                    true
+                }
+                Poll::Ready(None) => false,
+                Poll::Pending => true,
+            }
+        });
+        if let Some(item) = ready {
+            return Poll::Ready(Some(item));
+        }
+
+        Poll::Pending
+    }
+}
+
+struct RunnerBackgroundInner<E: Send + 'static> {
+    msg_sender: UnboundedSender<RunnerEvent<E>>,
+    msg_receiver: Mutex<UnboundedReceiver<RunnerEvent<E>>>,
+    streams: Mutex<StreamsState<E>>,
+    waker: AtomicWaker,
+}
+
+/// Helper for managing background streams, display messages, and structured
+/// events for a [`CommandRunner`].
+#[derive(Clone)]
+pub struct RunnerBackground<E: Send + 'static = ()> {
+    inner: Arc<RunnerBackgroundInner<E>>,
+}
+
+impl<E: Send + 'static> RunnerBackground<E> {
+    pub fn new() -> Self {
+        let (msg_sender, msg_receiver) = unbounded();
+        Self {
+            inner: Arc::new(RunnerBackgroundInner {
+                msg_sender,
+                msg_receiver: Mutex::new(msg_receiver),
+                streams: Mutex::new(StreamsState::default()),
+                waker: AtomicWaker::new(),
+            }),
+        }
+    }
+
+    /// Queues a displayable text message to be emitted by
+    /// [`Self::event_stream`].
+    pub fn println(&self, msg: impl Into<String>) {
+        let _ = self.inner.msg_sender.unbounded_send(RunnerEvent::Message(msg.into()));
+        self.inner.waker.wake();
+    }
+
+    /// Queues a structured domain event to be emitted by
+    /// [`Self::event_stream`].
+    pub fn send(&self, event: E) {
+        let _ = self.inner.msg_sender.unbounded_send(RunnerEvent::Event(event));
+        self.inner.waker.wake();
+    }
+
+    /// Returns an [`UnboundedSender`] that converts sent items into
+    /// [`RunnerEvent::Event`].
+    pub fn event_sender(&self) -> UnboundedSender<E> {
+        let (tx, rx) = unbounded::<E>();
+        self.spawn_stream(rx.map(RunnerEvent::Event));
+        tx
+    }
+
+    /// Dynamically registers or replaces a background stream under a unique
+    /// key. If a stream was already registered under `key`, the previous
+    /// stream is cancelled.
+    pub fn set_stream<K: Into<String>>(
+        &self,
+        key: K,
+        stream: impl Stream<Item = RunnerEvent<E>> + Send + 'static,
+    ) {
+        let mut streams = self.inner.streams.lock().unwrap();
+        streams.keyed.insert(key.into(), Box::pin(stream));
+        drop(streams);
+        self.inner.waker.wake();
+    }
+
+    /// Cancels the stream registered under `key`, if present. Returns `true` if
+    /// a stream was cancelled.
+    pub fn cancel_stream(&self, key: &str) -> bool {
+        let mut streams = self.inner.streams.lock().unwrap();
+        let removed = streams.keyed.remove(key).is_some();
+        drop(streams);
+        if removed {
+            self.inner.waker.wake();
+        }
+        removed
+    }
+
+    /// Spawns an anonymous background stream that runs until completion.
+    pub fn spawn_stream(&self, stream: impl Stream<Item = RunnerEvent<E>> + Send + 'static) {
+        let mut streams = self.inner.streams.lock().unwrap();
+        streams.anonymous.push(Box::pin(stream));
+        drop(streams);
+        self.inner.waker.wake();
+    }
+
+    /// Returns `true` if a background stream is currently registered under
+    /// `key`.
+    pub fn has_stream(&self, key: &str) -> bool {
+        self.inner.streams.lock().unwrap().keyed.contains_key(key)
+    }
+
+    /// Returns the total number of active background streams (keyed and
+    /// anonymous).
+    pub fn active_stream_count(&self) -> usize {
+        let streams = self.inner.streams.lock().unwrap();
+        streams.keyed.len() + streams.anonymous.len()
+    }
+
+    /// Returns a stream yielding messages and events produced by this runner.
+    pub fn event_stream(&self) -> RunnerBackgroundStream<E> {
+        RunnerBackgroundStream { inner: self.inner.clone() }
+    }
+}
+
+impl<E: Send + 'static> RunnerBackgroundInner<E> {
+    fn poll_next_event(&self, cx: &mut Context<'_>) -> Poll<Option<RunnerEvent<E>>> {
+        self.waker.register(cx.waker());
+
+        // 1. Drain directly queued messages/events from msg_receiver first.
+        if let Poll::Ready(Some(item)) = self.msg_receiver.lock().unwrap().poll_next_unpin(cx) {
+            return Poll::Ready(Some(item));
+        }
+
+        // 2. Poll registered background streams.
+        if let Poll::Ready(Some(item)) = self.streams.lock().unwrap().poll_next_unpin(cx) {
+            return Poll::Ready(Some(item));
+        }
+
+        Poll::Pending
+    }
+}
+
+impl<E: Send + 'static> Default for RunnerBackground<E> {
+    fn default() -> Self {
+        Self::new()
+    }
+}
+
+/// Stream returned by [`RunnerBackground::event_stream`].
+pub struct RunnerBackgroundStream<E: Send + 'static = ()> {
+    inner: Arc<RunnerBackgroundInner<E>>,
+}
+
+impl<E: Send + 'static> Stream for RunnerBackgroundStream<E> {
+    type Item = RunnerEvent<E>;
+
+    fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
+        self.inner.poll_next_event(cx)
+    }
+}
+
+/// A harness for driving a [`CommandRunner`] alongside its background event
+/// stream.
+pub struct DebugHarness<R: CommandRunner> {
+    runner: R,
+}
+
+impl<R: CommandRunner> DebugHarness<R> {
+    pub fn new(runner: R) -> Self {
+        Self { runner }
+    }
+
+    pub fn runner(&self) -> &R {
+        &self.runner
+    }
+
+    pub fn runner_mut(&mut self) -> &mut R {
+        &mut self.runner
+    }
+
+    pub fn into_runner(self) -> R {
+        self.runner
+    }
+
+    pub fn event_stream(&self) -> impl Stream<Item = RunnerEvent<R::Event>> + '_ {
+        self.runner.event_stream()
+    }
+
+    pub async fn run(
+        &self,
+        cmd: impl Into<CliCommand<R::Set>>,
+        args: Vec<String>,
+    ) -> Result<(), R::Error> {
+        self.runner.run(cmd, args).await
+    }
+
+    async fn drive_events<Fut, F>(&self, future: Fut, mut on_event: F) -> Fut::Output
+    where
+        Fut: futures::Future,
+        F: FnMut(RunnerEvent<R::Event>),
+    {
+        let stream = self.runner.event_stream().fuse();
+        futures::pin_mut!(stream);
+        let future = future.fuse();
+        futures::pin_mut!(future);
+
+        loop {
+            futures::select! {
+                event = stream.select_next_some() => on_event(event),
+                output = future => {
+                    while let Some(Some(event)) = stream.next().now_or_never() {
+                        on_event(event);
+                    }
+                    return output;
+                }
+            }
+        }
+    }
+
+    /// Concurrently drives the runner's background `event_stream` while
+    /// executing `future`. Messages are passed to `on_message` and events
+    /// are passed to `on_event`.
+    pub async fn drive_with<Fut, FM, FE>(
+        &self,
+        future: Fut,
+        mut on_message: FM,
+        mut on_event: FE,
+    ) -> Fut::Output
+    where
+        Fut: futures::Future,
+        FM: FnMut(String),
+        FE: FnMut(R::Event),
+    {
+        self.drive_events(future, |event| match event {
+            RunnerEvent::Message(msg) => on_message(msg),
+            RunnerEvent::Event(evt) => on_event(evt),
+        })
+        .await
+    }
+
+    /// Concurrently executes `future` while collecting all `RunnerEvent`s
+    /// emitted during its execution.
+    pub async fn run_collecting_events<Fut>(
+        &self,
+        future: Fut,
+    ) -> (Fut::Output, Vec<RunnerEvent<R::Event>>)
+    where
+        Fut: futures::Future,
+    {
+        let mut events = Vec::new();
+        let output = self.drive_events(future, |event| events.push(event)).await;
+        (output, events)
+    }
+
+    /// Concurrently executes a command while collecting all `RunnerEvent`s
+    /// emitted during its execution.
+    pub async fn run_command_collecting(
+        &self,
+        cmd: impl Into<CliCommand<R::Set>>,
+        args: Vec<String>,
+    ) -> (Result<(), R::Error>, Vec<RunnerEvent<R::Event>>) {
+        self.run_collecting_events(self.run(cmd, args)).await
+    }
 }
 
 #[cfg(test)]
@@ -494,6 +883,7 @@
         impl CommandRunner for TestRunner {
             type Set = RunnerCmd;
             type Error = std::io::Error;
+            type Event = ();
 
             fn run_command(
                 &self,
@@ -530,4 +920,350 @@
         // Reset persistent verbosity back to Info
         log::set_max_level(LevelFilter::Info);
     }
+
+    #[test]
+    fn test_command_runner_default_event_stream() {
+        gen_commandset! {
+            DefaultCmd {
+                Action = ("action", [], [], "Default action"),
+            }
+        }
+
+        struct SimpleRunner;
+        impl CommandRunner for SimpleRunner {
+            type Set = DefaultCmd;
+            type Error = std::io::Error;
+            type Event = ();
+
+            fn run_command(
+                &self,
+                _cmd: Self::Set,
+                _args: Vec<String>,
+            ) -> impl futures::Future<Output = Result<(), Self::Error>> {
+                futures::future::ready(Ok(()))
+            }
+        }
+
+        let runner = SimpleRunner;
+        let stream = runner.event_stream();
+        futures::pin_mut!(stream);
+        let mut cx = std::task::Context::from_waker(futures::task::noop_waker_ref());
+        assert!(stream.as_mut().poll_next(&mut cx).is_pending());
+    }
+
+    #[test]
+    fn test_runner_event() {
+        let msg: RunnerEvent<u32> = RunnerEvent::Message("hello".to_string());
+        assert_eq!(msg.message(), Some("hello"));
+        assert_eq!(msg.event(), None);
+        assert_eq!(msg.clone().into_event(), None);
+        assert_eq!(msg.clone().into_message(), Some("hello".to_string()));
+
+        let evt: RunnerEvent<u32> = RunnerEvent::Event(42);
+        assert_eq!(evt.message(), None);
+        assert_eq!(evt.event(), Some(&42));
+        assert_eq!(evt.clone().into_event(), Some(42));
+        assert_eq!(evt.into_message(), None);
+
+        let str_msg: RunnerEvent<()> = "test msg".into();
+        assert_eq!(str_msg, RunnerEvent::Message("test msg".to_string()));
+
+        let string_msg: RunnerEvent<()> = String::from("string msg").into();
+        assert_eq!(string_msg, RunnerEvent::Message("string msg".to_string()));
+
+        let mapped = RunnerEvent::Event(10).map_event(|x| x * 2);
+        assert_eq!(mapped, RunnerEvent::Event(20));
+
+        let mapped_msg: RunnerEvent<i32> =
+            RunnerEvent::Message("hi".to_string()).map_event(|x: u32| x as i32);
+        assert_eq!(mapped_msg, RunnerEvent::Message("hi".to_string()));
+    }
+
+    #[test]
+    fn test_runner_background_println_and_send() {
+        futures::executor::block_on(async {
+            let bg: RunnerBackground<u32> = RunnerBackground::new();
+            bg.println("log message 1");
+            bg.send(100);
+            bg.println("log message 2");
+
+            let stream = bg.event_stream();
+            futures::pin_mut!(stream);
+
+            assert_eq!(
+                stream.next().await,
+                Some(RunnerEvent::Message("log message 1".to_string()))
+            );
+            assert_eq!(stream.next().await, Some(RunnerEvent::Event(100)));
+            assert_eq!(
+                stream.next().await,
+                Some(RunnerEvent::Message("log message 2".to_string()))
+            );
+
+            // When no more events, polling should return Pending.
+            let mut cx = std::task::Context::from_waker(futures::task::noop_waker_ref());
+            assert!(stream.as_mut().poll_next(&mut cx).is_pending());
+        });
+    }
+
+    #[test]
+    fn test_runner_background_event_sender() {
+        futures::executor::block_on(async {
+            let bg: RunnerBackground<String> = RunnerBackground::new();
+            let sender = bg.event_sender();
+
+            sender.unbounded_send("hello".to_string()).unwrap();
+            sender.unbounded_send("world".to_string()).unwrap();
+
+            let stream = bg.event_stream();
+            futures::pin_mut!(stream);
+
+            assert_eq!(stream.next().await, Some(RunnerEvent::Event("hello".to_string())));
+            assert_eq!(stream.next().await, Some(RunnerEvent::Event("world".to_string())));
+
+            drop(sender);
+            // Polling after dropping sender should clean up the stream and
+            // return Pending.
+            let mut cx = std::task::Context::from_waker(futures::task::noop_waker_ref());
+            assert!(stream.as_mut().poll_next(&mut cx).is_pending());
+            assert_eq!(bg.active_stream_count(), 0);
+        });
+    }
+
+    #[test]
+    fn test_runner_background_keyed_streams_and_cancellation() {
+        futures::executor::block_on(async {
+            let bg: RunnerBackground<u32> = RunnerBackground::new();
+            assert_eq!(bg.has_stream("discovery"), false);
+
+            let (tx1, rx1) = futures::channel::mpsc::unbounded::<RunnerEvent<u32>>();
+            bg.set_stream("discovery", rx1);
+            assert_eq!(bg.has_stream("discovery"), true);
+            assert_eq!(bg.active_stream_count(), 1);
+
+            tx1.unbounded_send(RunnerEvent::Event(1)).unwrap();
+            let stream = bg.event_stream();
+            futures::pin_mut!(stream);
+            assert_eq!(stream.next().await, Some(RunnerEvent::Event(1)));
+
+            // Replacing the stream cancels the old one.
+            let (tx2, rx2) = futures::channel::mpsc::unbounded::<RunnerEvent<u32>>();
+            bg.set_stream("discovery", rx2);
+            assert_eq!(bg.has_stream("discovery"), true);
+            // Old tx1 should have receiver dropped / closed.
+            assert!(tx1.is_closed());
+
+            tx2.unbounded_send(RunnerEvent::Event(2)).unwrap();
+            assert_eq!(stream.next().await, Some(RunnerEvent::Event(2)));
+
+            // Cancelling the stream.
+            assert!(bg.cancel_stream("discovery"));
+            assert_eq!(bg.has_stream("discovery"), false);
+            assert!(tx2.is_closed());
+            assert!(!bg.cancel_stream("discovery"));
+            assert_eq!(bg.active_stream_count(), 0);
+        });
+    }
+
+    #[test]
+    fn test_runner_background_spawn_stream_and_cleanup() {
+        futures::executor::block_on(async {
+            let bg: RunnerBackground<u32> = RunnerBackground::new();
+            let stream_data = futures::stream::iter(vec![
+                RunnerEvent::Message("from stream".to_string()),
+                RunnerEvent::Event(99),
+            ]);
+            bg.spawn_stream(stream_data);
+            assert_eq!(bg.active_stream_count(), 1);
+
+            let stream = bg.event_stream();
+            futures::pin_mut!(stream);
+
+            assert_eq!(stream.next().await, Some(RunnerEvent::Message("from stream".to_string())));
+            assert_eq!(stream.next().await, Some(RunnerEvent::Event(99)));
+
+            // Polling again will observe the stream returned None and clean it
+            // up.
+            let mut cx = std::task::Context::from_waker(futures::task::noop_waker_ref());
+            assert!(stream.as_mut().poll_next(&mut cx).is_pending());
+            assert_eq!(bg.active_stream_count(), 0);
+        });
+    }
+
+    #[test]
+    fn test_runner_background_waker_waking() {
+        let bg: RunnerBackground<u32> = RunnerBackground::new();
+        let bg_clone = bg.clone();
+        let (polled_tx, polled_rx) = futures::channel::oneshot::channel();
+
+        let handle = std::thread::spawn(move || {
+            futures::executor::block_on(polled_rx).unwrap();
+            bg_clone.println("woken up");
+            bg_clone.send(42);
+        });
+
+        futures::executor::block_on(async {
+            let stream = bg.event_stream();
+            futures::pin_mut!(stream);
+            let mut polled_tx = Some(polled_tx);
+            let first = futures::future::poll_fn(|cx| {
+                let res = stream.as_mut().poll_next(cx);
+                if let Some(tx) = polled_tx.take() {
+                    assert!(res.is_pending());
+                    tx.send(()).unwrap();
+                }
+                res
+            })
+            .await;
+            assert_eq!(first, Some(RunnerEvent::Message("woken up".to_string())));
+            assert_eq!(stream.next().await, Some(RunnerEvent::Event(42)));
+        });
+
+        handle.join().unwrap();
+    }
+
+    #[test]
+    fn test_debug_harness() {
+        gen_commandset! {
+            HarnessCmd {
+                Greet = ("greet", [], ["name"], "Greet a user"),
+                Trigger = ("trigger", [], [], "Trigger background event"),
+            }
+        }
+
+        struct HarnessRunner {
+            bg: RunnerBackground<String>,
+        }
+
+        impl CommandRunner for HarnessRunner {
+            type Set = HarnessCmd;
+            type Error = std::io::Error;
+            type Event = String;
+
+            fn run_command(
+                &self,
+                cmd: Self::Set,
+                args: Vec<String>,
+            ) -> impl futures::Future<Output = Result<(), Self::Error>> {
+                match cmd {
+                    HarnessCmd::Greet => {
+                        let name = args.first().cloned().unwrap_or_else(|| "world".to_string());
+                        self.bg.println(format!("Hello, {}!", name));
+                    }
+                    HarnessCmd::Trigger => {
+                        self.bg.send("event_fired".to_string());
+                    }
+                }
+                futures::future::ready(Ok(()))
+            }
+
+            fn event_stream(&self) -> impl Stream<Item = RunnerEvent<Self::Event>> + '_ {
+                self.bg.event_stream()
+            }
+        }
+
+        let runner = HarnessRunner { bg: RunnerBackground::new() };
+        let harness = DebugHarness::new(runner);
+
+        futures::executor::block_on(async {
+            // Test run_command_collecting
+            let (res, events) =
+                harness.run_command_collecting(HarnessCmd::Greet, vec!["Alice".to_string()]).await;
+            assert!(res.is_ok());
+            assert_eq!(events, vec![RunnerEvent::Message("Hello, Alice!".to_string())]);
+
+            // Test drive_with with Trigger
+            let mut messages = Vec::new();
+            let mut domain_events = Vec::new();
+            harness
+                .drive_with(
+                    harness.run(HarnessCmd::Trigger, vec![]),
+                    |msg| messages.push(msg),
+                    |evt| domain_events.push(evt),
+                )
+                .await
+                .unwrap();
+
+            assert!(messages.is_empty());
+            assert_eq!(domain_events, vec!["event_fired".to_string()]);
+        });
+    }
+
+    #[test]
+    fn test_debug_harness_terminating_stream() {
+        gen_commandset! {
+            FiniteCmd {
+                Action = ("action", [], [], "Finite action"),
+            }
+        }
+
+        struct FiniteStreamRunner {
+            stream_done_tx: std::sync::Mutex<Option<futures::channel::oneshot::Sender<()>>>,
+        }
+        impl CommandRunner for FiniteStreamRunner {
+            type Set = FiniteCmd;
+            type Error = std::io::Error;
+            type Event = u32;
+
+            fn run_command(
+                &self,
+                _cmd: Self::Set,
+                _args: Vec<String>,
+            ) -> impl futures::Future<Output = Result<(), Self::Error>> {
+                futures::future::ready(Ok(()))
+            }
+
+            fn event_stream(&self) -> impl Stream<Item = RunnerEvent<Self::Event>> + '_ {
+                let mut done_tx = self.stream_done_tx.lock().unwrap().take();
+                futures::stream::iter(vec![
+                    RunnerEvent::Message("finite msg".to_string()),
+                    RunnerEvent::Event(42),
+                ])
+                .chain(futures::stream::poll_fn(move |_| {
+                    if let Some(tx) = done_tx.take() {
+                        let _ = tx.send(());
+                    }
+                    Poll::Ready(None)
+                }))
+            }
+        }
+
+        let (stream_done_tx, stream_done_rx) = futures::channel::oneshot::channel::<()>();
+        let harness = DebugHarness::new(FiniteStreamRunner {
+            stream_done_tx: std::sync::Mutex::new(Some(stream_done_tx)),
+        });
+        futures::executor::block_on(async {
+            let (tx, rx) = futures::channel::oneshot::channel::<()>();
+            let handle = std::thread::spawn(move || {
+                futures::executor::block_on(stream_done_rx).unwrap();
+                let _ = tx.send(());
+            });
+
+            // The stream terminates after 2 items while rx is still pending.
+            // run_collecting_events should not busy-loop and should cleanly
+            // complete when rx resolves.
+            let (output, events) = harness.run_collecting_events(rx).await;
+            assert!(output.is_ok());
+            assert_eq!(
+                events,
+                vec![RunnerEvent::Message("finite msg".to_string()), RunnerEvent::Event(42),]
+            );
+            handle.join().unwrap();
+        });
+    }
+
+    #[test]
+    fn test_runner_background_stream_static() {
+        let bg: RunnerBackground<u32> = RunnerBackground::new();
+        bg.send(77);
+        let stream = bg.event_stream();
+
+        fn assert_static<T: 'static>(_: &T) {}
+        assert_static(&stream);
+
+        futures::executor::block_on(async {
+            futures::pin_mut!(stream);
+            assert_eq!(stream.next().await, Some(RunnerEvent::Event(77)));
+        });
+    }
 }
diff --git a/rust/bt-pacs/src/debug.rs b/rust/bt-pacs/src/debug.rs
index 342d29b..9b815d7 100644
--- a/rust/bt-pacs/src/debug.rs
+++ b/rust/bt-pacs/src/debug.rs
@@ -32,6 +32,7 @@
 impl<T: bt_gatt::GattTypes> CommandRunner for PacsDebug<T> {
     type Set = PacsCmd;
     type Error = bt_gatt::types::Error;
+    type Event = ();
 
     fn run_command(
         &self,
diff --git a/rust/bt-vcs/src/debug.rs b/rust/bt-vcs/src/debug.rs
index fd3cdc9..780f045 100644
--- a/rust/bt-vcs/src/debug.rs
+++ b/rust/bt-vcs/src/debug.rs
@@ -59,6 +59,7 @@
 impl<T: bt_gatt::GattTypes> CommandRunner for VcsDebug<T> {
     type Set = VcsCmd;
     type Error = Error;
+    type Event = ();
 
     fn run_command(
         &self,