From 1b7bee11602cd8c70d8fc109bc95f11fd681bb10 Mon Sep 17 00:00:00 2001 From: Richard Bowman Date: Fri, 9 Oct 2026 00:09:32 +0100 Subject: [PATCH 1/7] Extend MessageBroker I've added "stream" as a message type, and provided "next message" and "close streams by affordance" methods. This should prepare the way for using MessageBroker for MJPEG streams. --- src/labthings_fastapi/message_broker.py | 34 ++++++++++++++++++++++++- 1 file changed, 33 insertions(+), 1 deletion(-) diff --git a/src/labthings_fastapi/message_broker.py b/src/labthings_fastapi/message_broker.py index d7ce98c7..1621170a 100644 --- a/src/labthings_fastapi/message_broker.py +++ b/src/labthings_fastapi/message_broker.py @@ -36,7 +36,7 @@ class Message: thing: str affordance: str - message_type: Literal["property", "action"] + message_type: Literal["property", "action", "stream"] payload: Any @@ -83,6 +83,19 @@ async def subscribe( streams = affordances.setdefault(affordance, WeakSet()) streams.add(stream) + async def next_message(self, thing: str, affordance: str) -> Message: + """Get the next message from a particular affordance. + + Note that there's no timeout: standard `anyio` commands may be used to do that. + + :param thing: The name of the `.Thing` being subscribed to. + :param affordance: The name of the affordance being subscribed to. + :return: the next message from the specified affordance. + """ + send, recv = anyio.create_memory_object_stream[Message](max_buffer_size=1) + await self.subscribe(thing, affordance, send) + return await recv.receive() + async def unsubscribe( self, thing: str, affordance: str, stream: MemoryObjectSendStream[Message] ) -> None: @@ -148,3 +161,22 @@ async def close_streams(self) -> None: for subs in thing_subs.values(): for stream in subs: tg.start_soon(stream.aclose) + + async def close_streams_for_affordance(self, thing: str, affordance: str) -> int: + """Close all streams subscribed to a particular affordance. + + This may be called to signal that a stream has stopped. + + :param thing: the Thing name. + :param affordance: the affordance name. + :return: the number of streams that were closed. + """ + try: + streams = self._subscriptions[thing][affordance] + except KeyError: + return 0 # If there are no streams, ignore it. + async with anyio.create_task_group() as tg: + n = len(streams) + for stream in streams: + tg.start_soon(stream.aclose) + return n From 88ea8ff85d5c28887afac72ee1dc58b500654492 Mon Sep 17 00:00:00 2001 From: Richard Bowman Date: Fri, 9 Oct 2026 00:10:20 +0100 Subject: [PATCH 2/7] Expose the message broker. The thing server interface now exposes the message broker. --- src/labthings_fastapi/thing_server_interface.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/src/labthings_fastapi/thing_server_interface.py b/src/labthings_fastapi/thing_server_interface.py index 3acbf011..a9d0bb8c 100644 --- a/src/labthings_fastapi/thing_server_interface.py +++ b/src/labthings_fastapi/thing_server_interface.py @@ -20,7 +20,7 @@ from labthings_fastapi.exceptions import FeatureNotEnabledError, ServerNotRunningError from labthings_fastapi.global_lock import GlobalLock -from labthings_fastapi.message_broker import Message +from labthings_fastapi.message_broker import Message, MessageBroker if TYPE_CHECKING: from labthings_fastapi.actions import ActionManager @@ -156,6 +156,11 @@ def publish(self, message: Message) -> None: except ServerNotRunningError: pass # If the server isn't running yet, we can't publish events. + @property + def message_broker(self) -> MessageBroker: + """The message broker associated with the current server.""" + return self._get_server().message_broker + @property def settings_folder(self) -> str: """The path to a folder where persistent files may be saved.""" From 451fa2aee8e678427e7bed69bc71d7488560ee5e Mon Sep 17 00:00:00 2001 From: Richard Bowman Date: Fri, 9 Oct 2026 00:18:29 +0100 Subject: [PATCH 3/7] Refactor MJPEG stream This strips out the pub/sub logic from MJPEGStream in favour of using MessageBroker, as is already done for property notifications. The resulting class should be much simpler. I've also inherited from BaseDescriptor to save duplication in MJPEGStreamDescriptor. --- src/labthings_fastapi/outputs/mjpeg_stream.py | 259 +++++------------- tests/test_mjpeg_stream.py | 1 - 2 files changed, 76 insertions(+), 184 deletions(-) diff --git a/src/labthings_fastapi/outputs/mjpeg_stream.py b/src/labthings_fastapi/outputs/mjpeg_stream.py index 39f4d7e2..5bc60bd9 100644 --- a/src/labthings_fastapi/outputs/mjpeg_stream.py +++ b/src/labthings_fastapi/outputs/mjpeg_stream.py @@ -6,33 +6,28 @@ from __future__ import annotations -import logging import threading from dataclasses import dataclass from datetime import datetime from typing import ( - TYPE_CHECKING, Any, AsyncGenerator, - Literal, Optional, - Union, - overload, ) import anyio +from anyio.streams.memory import MemoryObjectReceiveStream, MemoryObjectSendStream from fastapi import FastAPI from fastapi.responses import HTMLResponse, StreamingResponse -from typing_extensions import Self -if TYPE_CHECKING: - from labthings_fastapi.thing import Thing - from labthings_fastapi.thing_server_interface import ThingServerInterface +from labthings_fastapi.base_descriptor import BaseDescriptor +from labthings_fastapi.message_broker import Message +from labthings_fastapi.thing import Thing @dataclass -class RingbufferEntry: - """A single entry in a ringbuffer. +class Frame: + """A single frame in the stream. This structure comprises one frame as a JPEG, plus a timestamp and a buffer index. Each time a frame is added to the stream, it is @@ -51,6 +46,24 @@ class RingbufferEntry: """The index of the frame within the stream.""" +def frame_payload(message: Message) -> Frame: + """Extract a Frame object from a Message. + + This checks that the message is of "stream" type, and checks the + type of its payload. + + :param message: the message containing the Frame. + :return: the Frame object. + :raises TypeError: if the message isn't a stream message containing + a `Frame` object. + """ + if message.message_type != "stream": + raise TypeError("The message {message} was not part of a stream.") + if not isinstance(message.payload, Frame): + raise TypeError("The message {message} didn't contain a Frame.") + return message.payload + + class MJPEGStreamResponse(StreamingResponse): """A StreamingResponse that streams an MJPEG stream. @@ -64,7 +77,10 @@ class MJPEGStreamResponse(StreamingResponse): """The media_type used to describe the endpoint in FastAPI.""" def __init__( - self, gen: AsyncGenerator[bytes, None], status_code: int = 200 + self, + send: MemoryObjectSendStream[Message], + recv: MemoryObjectReceiveStream[Message], + status_code: int = 200, ) -> None: """Set up StreamingResponse that streams an MJPEG stream. @@ -73,18 +89,20 @@ def __init__( types that mark it as an MJPEG stream. This is sufficient to enable it to work in an `img` tag, with the `src` set to the MJPEG stream's endpoint. - It expects an async generator that supplies individual JPEGs to be streamed, - such as the one provided by `.MJPEGStream`. + It expects to get a stream that receives `Message` objects with `Frame` + payloads. Both the send and receive streams are retained, because the send + stream is only weakly referenced by the message broker. NB the ``status_code`` argument is used by FastAPI to set the status code of the response in OpenAPI. - :param gen: an async generator, yielding `bytes` objects each of which is - one image, in JPEG format. + :param send: the send stream subscribed to the MJPEG stream. + :param recv: the receive stream subscribed to the MJPEG stream. :param status_code: The status code associated with the response, by default a 200 code is returned. """ - self.frame_async_generator = gen + self._send_stream = send + self._receive_stream = recv StreamingResponse.__init__( self, self.mjpeg_async_generator(), @@ -101,9 +119,10 @@ async def mjpeg_async_generator(self) -> AsyncGenerator[bytes, None]: :yield: JPEG frames, each with a ``--frame`` marker prepended. """ - async for frame in self.frame_async_generator: + async for message in self._receive_stream: + frame = frame_payload(message) yield b"--frame\r\nContent-Type: image/jpeg\r\n\r\n" - yield frame + yield frame.frame yield b"\r\n" @@ -128,26 +147,19 @@ class MJPEGStream: of new frames, and then retrieving the frame (shortly) afterwards. """ - def __init__( - self, thing_server_interface: ThingServerInterface, ringbuffer_size: int = 10 - ) -> None: + def __init__(self, thing: Thing, name: str) -> None: """Initialise an MJPEG stream. See the class docstring for `.MJPEGStream`. Note that it will often be initialised by `.MJPEGStreamDescriptor`. - :param thing_server_interface: the `~lt.ThingServerInterface` of the - `~lt.Thing` associated with this stream. It's used to run the async - code that relays frames to open connections. - :param ringbuffer_size: The number of frames to retain in - memory, to allow retrieval after the frame has been sent. + :param thing: the `~lt.Thing` on which this stream is defined. + :param name: The attribute name of this stream. """ self._lock = threading.Lock() - self.condition = anyio.Condition() - self._streaming = False - self._ringbuffer: list[RingbufferEntry] = [] - self._thing_server_interface = thing_server_interface - self.reset(ringbuffer_size=ringbuffer_size) + self._thing = thing + self._name = name + self.reset() def reset(self, ringbuffer_size: Optional[int] = None) -> None: """Reset the stream and optionally change the ringbuffer size. @@ -157,16 +169,6 @@ def reset(self, ringbuffer_size: Optional[int] = None) -> None: :param ringbuffer_size: the number of frames to keep in memory. """ with self._lock: - self._streaming = True - n = ringbuffer_size or len(self._ringbuffer) - self._ringbuffer = [ - RingbufferEntry( - frame=b"", - index=-1, - timestamp=datetime.min, - ) - for i in range(n) - ] self.last_frame_i = -1 def stop(self) -> None: @@ -174,53 +176,12 @@ def stop(self) -> None: Stop the stream and cause all clients to disconnect. """ - with self._lock: - self._streaming = False - self._thing_server_interface.start_async_task_soon( - self.notify_stream_stopped - ) - - async def ringbuffer_entry(self, i: int) -> RingbufferEntry: - """Return the ith frame acquired by the camera. - - The ringbuffer means we can retrieve frames even if they are not - the latest frame. Specifying ``i`` also makes it simple to ensure - that every frame in a stream is acquired. - - :param i: The index of the frame to read. - - :return: the frame, together with a timestamp and its index. - - :raise ValueError: if the frame is not available. - """ - if i < 0: - raise ValueError("i must be >= 0") - if i < self.last_frame_i - len(self._ringbuffer) + 2: - raise ValueError("the ith frame has been overwritten") - if i > self.last_frame_i: - # TODO: await the ith frame - raise ValueError("the ith frame has not yet been acquired") - entry = self._ringbuffer[i % len(self._ringbuffer)] - if entry.index != i: - raise ValueError("the ith frame has been overwritten") - return entry - - async def next_frame(self) -> int: - """Wait for the next frame, and return its index. - - This async function will yield until a new frame arrives, then return - its index. The index may then be used to retrieve the new frame - with `.MJPEGStream.ringbuffer_entry`. - - :return: the index of the next frame to arrive. - - :raise StopAsyncIteration: if the stream has stopped. - """ - async with self.condition: - await self.condition.wait() - if not self._streaming: - raise StopAsyncIteration() - return self.last_frame_i + tsi = self._thing._thing_server_interface + tsi.call_async_task( + tsi.message_broker.close_streams_for_affordance, + self._thing.name, + self._name, + ) async def grab_frame(self) -> bytes: """Wait for the next frame, and return it. @@ -230,9 +191,11 @@ async def grab_frame(self) -> bytes: :return: The next JPEG frame, as a `bytes` object. """ - i = await self.next_frame() - entry = await self.ringbuffer_entry(i) - return entry.frame + message = await self._thing._thing_server_interface.message_broker.next_message( + self._thing.name, self._name + ) + payload = frame_payload(message) + return payload.frame async def next_frame_size(self) -> int: """Wait for the next frame and return its size. @@ -241,34 +204,7 @@ async def next_frame_size(self) -> int: :return: The size of the next JPEG frame, in bytes. """ - i = await self.next_frame() - entry = await self.ringbuffer_entry(i) - return len(entry.frame) - - async def frame_async_generator(self) -> AsyncGenerator[bytes, None]: - """Yield frames as bytes objects. - - This generator will return frames from the MJPEG stream. - - Note that this will wait for a new frame each time. There is no - guarantee that we won't skip frames. - - :yield: the frames in sequence, as a `bytes` object containing - JPEG data. - """ - while self._streaming: - try: - i = await self.next_frame() - entry = await self.ringbuffer_entry(i) - yield entry.frame - except StopAsyncIteration: - break - except Exception as e: # noqa: BLE001 - # It's important that errors in the stream don't crash the server. - # This may be something we can remove in the future, now streams stop - # more elegantly. However, it will require careful testing.f - logging.exception(f"Error in stream: {e}, stream stopped") - return + return len(await self.grab_frame()) async def mjpeg_stream_response(self) -> MJPEGStreamResponse: """Return a StreamingResponse that streams an MJPEG stream. @@ -281,7 +217,11 @@ async def mjpeg_stream_response(self) -> MJPEGStreamResponse: :return: a streaming response in MJPEG format. """ - return MJPEGStreamResponse(self.frame_async_generator()) + send, recv = anyio.create_memory_object_stream[Message](max_buffer_size=1) + await self._thing._thing_server_interface.message_broker.subscribe( + self._thing.name, self._name, send + ) + return MJPEGStreamResponse(send, recv) def add_frame(self, frame: bytes, timestamp: Optional[datetime] = None) -> None: """Add a JPEG to the MJPEG stream. @@ -307,40 +247,18 @@ def add_frame(self, frame: bytes, timestamp: Optional[datetime] = None) -> None: ): raise ValueError("Invalid JPEG") with self._lock: - entry = self._ringbuffer[(self.last_frame_i + 1) % len(self._ringbuffer)] - entry.timestamp = timestamp if timestamp is not None else datetime.now() - entry.frame = frame - entry.index = self.last_frame_i + 1 - self._thing_server_interface.start_async_task_soon( - self.notify_new_frame, entry.index - ) - - async def notify_new_frame(self, i: int) -> None: - """Notify any waiting tasks that a new frame is available. - - :param i: The number of the frame (which counts up since the server starts) - """ - async with self.condition: - self.last_frame_i = i - self.condition.notify_all() - - async def notify_stream_stopped(self) -> None: - """Raise an exception in any waiting tasks to signal the stream has stopped. - - This should be run only when streaming has stopped, i.e. ``self._streaming`` - is ``False`` and an error will be raised if this isn't the case. - - :raises RuntimeError: if the stream is still streaming. - """ - if self._streaming is True: - raise RuntimeError( - "This function should only be called when the stream is stopped." + self.last_frame_i += 1 + payload = Frame( + frame=frame, + timestamp=timestamp if timestamp is not None else datetime.now(), + index=self.last_frame_i, ) - async with self.condition: - self.condition.notify_all() + self._thing._thing_server_interface.publish( + Message(self._thing.name, self._name, "stream", payload) + ) -class MJPEGStreamDescriptor: +class MJPEGStreamDescriptor(BaseDescriptor[Thing, MJPEGStream]): """A descriptor that returns a MJPEGStream object when accessed. If this descriptor is added to a `~lt.Thing`, it will create an `.MJPEGStream` @@ -357,48 +275,23 @@ def __init__(self, **kwargs: Any) -> None: :param \**kwargs: keyword arguments are passed to the initialiser of `.MJPEGStream`. """ + super().__init__() self._kwargs: Any = kwargs - def __set_name__(self, _owner: Thing, name: str) -> None: - """Remember the name to which we are assigned. - - The name is important, as it will set the URL of the HTTP endpoint used - to access the stream. - - :param _owner: the `~lt.Thing` to which we are attached. - :param name: the name to which this descriptor is assigned. - """ - self.name = name - - @overload - def __get__(self, obj: Literal[None], type: type | None = None) -> Self: ... # noqa: D105 - - @overload - def __get__(self, obj: Thing, type: type | None = None) -> MJPEGStream: ... # noqa: D105 - - def __get__( - self, obj: Optional[Thing], type: type[Thing] | None = None - ) -> Union[MJPEGStream, Self]: - """Return the MJPEG Stream, or the descriptor object. - - When accessed on the class, this ``__get__`` method will return the descriptor - object. This allows LabThings to add it to the HTTP API. - - When accessed on the object, an `.MJPEGStream` is returned. + def instance_get(self, obj: Thing) -> MJPEGStream: + """Return the MJPEG Stream. - :param obj: the host `~lt.Thing`, or ``None`` if accessed on the class. - :param type: the class on which we are defined. + :param obj: the host `~lt.Thing`. :return: an `.MJPEGStream`, or this descriptor. """ - if obj is None: - return self try: return obj.__dict__[self.name] except KeyError: obj.__dict__[self.name] = MJPEGStream( + thing=obj, + name=self.name, **self._kwargs, - thing_server_interface=obj._thing_server_interface, ) return obj.__dict__[self.name] diff --git a/tests/test_mjpeg_stream.py b/tests/test_mjpeg_stream.py index 1aefa2b8..f8b30cb8 100644 --- a/tests/test_mjpeg_stream.py +++ b/tests/test_mjpeg_stream.py @@ -41,7 +41,6 @@ def _make_images(self): time.sleep(1 / self.framerate) i = i + 1 self.stream.stop() - self._streaming = False @pytest.fixture From 59474393ba2cbd7f416d59364208e8c1eb1fbbf8 Mon Sep 17 00:00:00 2001 From: Richard Bowman Date: Fri, 9 Oct 2026 00:32:42 +0100 Subject: [PATCH 4/7] Pull over improved tests from #372 This uses an event to speed up tests by eliminating `time.sleep` and adds some more checks on the data received. --- tests/test_mjpeg_stream.py | 171 ++++++++++++++++++++++++++++++++----- 1 file changed, 148 insertions(+), 23 deletions(-) diff --git a/tests/test_mjpeg_stream.py b/tests/test_mjpeg_stream.py index f8b30cb8..76198abc 100644 --- a/tests/test_mjpeg_stream.py +++ b/tests/test_mjpeg_stream.py @@ -12,7 +12,9 @@ class Telly(lt.Thing): _stream_thread: threading.Thread _streaming: bool = False framerate: float = 1000 - frame_limit: int = 3 + frame_limit: int = 999 + frame_event: threading.Event | None = None + initial_delay: float = 0 stream = lt.outputs.MJPEGStreamDescriptor() @@ -23,6 +25,10 @@ def __enter__(self): def __exit__(self, exc_t, exc_v, exc_tb): self._streaming = False + if self.frame_event: + # Trigger an iteration of the loop, so that it + # will terminate rather than hang forever + self.frame_event.set() self._stream_thread.join() def _make_images(self): @@ -35,47 +41,166 @@ def _make_images(self): image.save(dest, "jpeg") jpegs.append(dest.getvalue()) + if self.initial_delay > 0: + time.sleep(self.initial_delay) + i = 0 while self._streaming and (i < self.frame_limit or self.frame_limit < 0): self.stream.add_frame(jpegs[i % len(jpegs)]) - time.sleep(1 / self.framerate) i = i + 1 + if self.frame_event: + self.frame_event.wait() + self.frame_event.clear() + else: + time.sleep(1 / self.framerate) self.stream.stop() + self._streaming = False + + +def assert_magic_bytes(frame: bytes) -> None: + """Check that a `bytes` object starts and ends with the JPEG markers.""" + assert frame[0] == 0xFF + assert frame[1] == 0xD8 + assert frame[-2] == 0xFF + assert frame[-1] == 0xD9 @pytest.fixture -def client(): - """Yield a test client connected to a ThingServer""" - server = lt.ThingServer.from_things({"telly": Telly}) - with server.test_client() as client: - yield client +def server(): + """Yield a ThingServer with the `Telly` thing.""" + return lt.ThingServer.from_things({"telly": Telly}) + + +@pytest.fixture +def telly(server): + """Yield a Telly thing from the server.""" + telly = server.things["telly"] + assert isinstance(telly, Telly) + return telly + +def test_grab_and_shutdown(server: lt.ThingServer, telly: Telly): + """Check we can grab frames, and shut down cleanly. -def test_mjpeg_stream(client): - """Verify the MJPEG stream contains at least one frame marker. + This test uses an Event to synchronise new frames with the various + methods we're calling to retrieve them. This is intended to make the + test suite faster and less reliant on `time.sleep`. + + We check that `grab_frame` and `next_frame_size` both work when the + camera emits frames, and also that they raise an error if they're called + once the camera has stopped. + + We verify that the stream can be stopped, though we don't shut down any + long-running listeners, because limitations in TestClient mean this + isn't possible. It would be possible if we spun up an actual HTTP server, + but that's quite high-effort. + + The async grab functions need to run in the event loop, which happens in + a background thread during the `with server.test_client()` block. We use + the thing server interface to run tasks in the event loop for convenience. + """ + telly.frame_event = threading.Event() # Make timings more deterministic. + with server.test_client(): + # this `with` block starts an event loop and runs the server. + + # Grab some frames and check we get a JPEG back + for _ in range(3): + assert telly._stream_thread.is_alive() # Catch premature termination + # Start the coroutine to grab a frame + future = telly._thing_server_interface.start_async_task_soon( + telly.stream.grab_frame + ) + # then use the event to trigger the next frame. + telly.frame_event.set() + # then wait for the coroutine to finish + frame = future.result() + # this should return a valid JPEG, which we can (sort of) check below + assert_magic_bytes(frame) + + # Repeat the process for grabbing a frame to test `next_frame_size` + future = telly._thing_server_interface.start_async_task_soon( + telly.stream.next_frame_size + ) + telly.frame_event.set() + size = future.result() # wait for the frame to be grabbed + assert isinstance(size, int) + assert size > 0 + + # Close all streams + telly._thing_server_interface.call_async_task(telly.stream.close_streams) + + # We shouldn't be able to get any more frames now + # This means that the stream won't generate any new stream pairs + with pytest.raises(StopAsyncIteration): + telly.stream.connected_stream() + # The grab functions depend on the function above, so they should also fail. + with pytest.raises(StopAsyncIteration): + telly._thing_server_interface.call_async_task(telly.stream.grab_frame) + with pytest.raises(StopAsyncIteration): + telly._thing_server_interface.call_async_task(telly.stream.next_frame_size) + # The background thread gets shut down by `Telly.__exit__`. + + +def test_mjpeg_stream_http(server: lt.ThingServer, telly: Telly): + """Verify the MJPEG stream works, and is shut down cleanly. A limitation of the TestClient is that it can't actually stream. This means that all of the frames sent by our test Thing will arrive in a single packet. - For now, we just check it starts with the frame separator, - but it might be possible in the future to check there are three - images there. + For now, we download all the data and then chop it up afterwards. + + The `Telly` will send exactly 3 JPEGs, with a short delay at + the start to make sure the `StreamingResponse` is created and + doesn't miss any frames. + + This test also verifies the stream is shut down by the server + when it's stopped - if it wasn't terminated by the server, the + `client.stream()` call would hang indefinitely. """ - with client.stream("GET", "/telly/stream") as stream: - stream.raise_for_status() - received = 0 - for b in stream.iter_bytes(): - received += 1 - assert b.startswith(b"--frame") + telly.frame_limit = 3 # stream 3 frames and then stop + telly.initial_delay = 0.05 # give the client time to start listening + with server.test_client() as client: + with client.stream("GET", "/telly/stream") as response: + # Note: we don't actually enter this `with` block until after + # the camera background thread has finished and the stream + # has been closed. + response.raise_for_status() + parts = 0 + mjpeg_data = b"" + for b in response.iter_bytes(): + parts += 1 + mjpeg_data += b + # Due to a quirk in TestClient, we get all the data in a single + # chunk - this limits our ability to test the stream. + # If that's fixed in the future, the assertion below will fail, + # which should prompt us to improve these tests. + assert parts == 1 + + # Check the received data contained the expected number of frames + # We split chunks using the known header, and remove extra white + # space before checking for JPEG start/end bytes. + chunks = mjpeg_data.split(b"--frame\r\nContent-Type: image/jpeg") + n = 0 + for chunk in chunks: + # Check each chunk is a JPEG + stripped = chunk.strip(b"\r\n") + if len(stripped) > 10: # If the chunk doesn't look empty + assert_magic_bytes(stripped) + n += 1 + assert telly.stream.last_frame_i == 2 + assert n == telly.frame_limit if __name__ == "__main__": - import uvicorn - - server = lt.ThingServer.from_things({"telly": Telly}) - telly = server.things["telly"] + """This block allows you to connect manually with a web browser. + + That's helpful, because the tests above don't actually stream anything + as noted in `test_mjpeg_stream_http`. + """ + thing_server = lt.ThingServer.from_things({"telly": Telly}) + telly = thing_server.things["telly"] assert isinstance(telly, Telly) telly.framerate = 6 telly.frame_limit = -1 - uvicorn.run(server.app, port=5000, ws="websockets-sansio") + thing_server.serve() From 6a479af1566ac49d9da08adace479c16f726caca Mon Sep 17 00:00:00 2001 From: Richard Bowman Date: Fri, 9 Oct 2026 00:45:13 +0100 Subject: [PATCH 5/7] Remove unused ringbuffer argument --- src/labthings_fastapi/outputs/mjpeg_stream.py | 9 ++------- 1 file changed, 2 insertions(+), 7 deletions(-) diff --git a/src/labthings_fastapi/outputs/mjpeg_stream.py b/src/labthings_fastapi/outputs/mjpeg_stream.py index 5bc60bd9..4318007a 100644 --- a/src/labthings_fastapi/outputs/mjpeg_stream.py +++ b/src/labthings_fastapi/outputs/mjpeg_stream.py @@ -161,13 +161,8 @@ def __init__(self, thing: Thing, name: str) -> None: self._name = name self.reset() - def reset(self, ringbuffer_size: Optional[int] = None) -> None: - """Reset the stream and optionally change the ringbuffer size. - - Discard all frames from the ringbuffer and reset the frame index. - - :param ringbuffer_size: the number of frames to keep in memory. - """ + def reset(self) -> None: + """Reset the frame index.""" with self._lock: self.last_frame_i = -1 From e16243024dfcd6709999bc22a4298faa383f0c30 Mon Sep 17 00:00:00 2001 From: Richard Bowman Date: Fri, 9 Oct 2026 00:46:21 +0100 Subject: [PATCH 6/7] Test the type checking function and remove broken tests. The stream no longer has the notion of being "active" so we can't check for errors there. This might be a feature we want to revive - I'm not totally sure. --- tests/test_mjpeg_stream.py | 31 ++++++++++++++++++++----------- 1 file changed, 20 insertions(+), 11 deletions(-) diff --git a/tests/test_mjpeg_stream.py b/tests/test_mjpeg_stream.py index 76198abc..5028e8da 100644 --- a/tests/test_mjpeg_stream.py +++ b/tests/test_mjpeg_stream.py @@ -1,11 +1,14 @@ import io import threading import time +from datetime import datetime import pytest from PIL import Image import labthings_fastapi as lt +from labthings_fastapi.message_broker import Message +from labthings_fastapi.outputs.mjpeg_stream import Frame, frame_payload class Telly(lt.Thing): @@ -79,6 +82,21 @@ def telly(server): return telly +def test_frame_payload(): + """Check that a frame is returned, or we get an error.""" + frame = Frame(b"payload", datetime.now(), 1) + good_message = Message("thing", "affordance", "stream", frame) + assert frame_payload(good_message) is frame + + bad_message = Message("thing", "affordance", "action", frame) + with pytest.raises(TypeError): + frame_payload(bad_message) + + bad_message2 = Message("thing", "affordance", "stream", b"payload") + with pytest.raises(TypeError): + frame_payload(bad_message2) + + def test_grab_and_shutdown(server: lt.ThingServer, telly: Telly): """Check we can grab frames, and shut down cleanly. @@ -127,17 +145,8 @@ def test_grab_and_shutdown(server: lt.ThingServer, telly: Telly): assert size > 0 # Close all streams - telly._thing_server_interface.call_async_task(telly.stream.close_streams) - - # We shouldn't be able to get any more frames now - # This means that the stream won't generate any new stream pairs - with pytest.raises(StopAsyncIteration): - telly.stream.connected_stream() - # The grab functions depend on the function above, so they should also fail. - with pytest.raises(StopAsyncIteration): - telly._thing_server_interface.call_async_task(telly.stream.grab_frame) - with pytest.raises(StopAsyncIteration): - telly._thing_server_interface.call_async_task(telly.stream.next_frame_size) + telly.stream.stop() + # The background thread gets shut down by `Telly.__exit__`. From fc5c0287ace1bcb4a4f3c1a37c7a9ff95c592fad Mon Sep 17 00:00:00 2001 From: Richard Bowman Date: Fri, 9 Oct 2026 01:28:31 +0100 Subject: [PATCH 7/7] Add a function to get the next frame + metadata. This is tested by a better-isolated test that mocks the server and uses a real MessageBroker. --- src/labthings_fastapi/outputs/mjpeg_stream.py | 16 +++-- tests/test_mjpeg_stream.py | 70 ++++++++++++++++--- 2 files changed, 73 insertions(+), 13 deletions(-) diff --git a/src/labthings_fastapi/outputs/mjpeg_stream.py b/src/labthings_fastapi/outputs/mjpeg_stream.py index 4318007a..7752499f 100644 --- a/src/labthings_fastapi/outputs/mjpeg_stream.py +++ b/src/labthings_fastapi/outputs/mjpeg_stream.py @@ -178,18 +178,24 @@ def stop(self) -> None: self._name, ) - async def grab_frame(self) -> bytes: + async def grab_frame_with_metadata(self) -> Frame: """Wait for the next frame, and return it. - This copies the frame for safety, so there is no need to release - or return the buffer. + This includes metadata such as the timestamp. - :return: The next JPEG frame, as a `bytes` object. + :return: The next JPEG frame, in a dataclass with metadata. """ message = await self._thing._thing_server_interface.message_broker.next_message( self._thing.name, self._name ) - payload = frame_payload(message) + return frame_payload(message) + + async def grab_frame(self) -> bytes: + """Wait for the next frame, and return it. + + :return: The next JPEG frame, as a `bytes` object. + """ + payload = await self.grab_frame_with_metadata() return payload.frame async def next_frame_size(self) -> int: diff --git a/tests/test_mjpeg_stream.py b/tests/test_mjpeg_stream.py index 5028e8da..94b31762 100644 --- a/tests/test_mjpeg_stream.py +++ b/tests/test_mjpeg_stream.py @@ -3,12 +3,22 @@ import time from datetime import datetime +import anyio import pytest +from anyio import from_thread, to_thread from PIL import Image import labthings_fastapi as lt -from labthings_fastapi.message_broker import Message -from labthings_fastapi.outputs.mjpeg_stream import Frame, frame_payload +from labthings_fastapi.message_broker import Message, MessageBroker +from labthings_fastapi.outputs.mjpeg_stream import Frame, MJPEGStream, frame_payload + + +def make_jpeg(colour: str) -> bytes: + """Return a 10x10 JPEG of the given colour.""" + image = Image.new("RGB", (10, 10), colour) + dest = io.BytesIO() + image.save(dest, "jpeg") + return dest.getvalue() class Telly(lt.Thing): @@ -37,12 +47,7 @@ def __exit__(self, exc_t, exc_v, exc_tb): def _make_images(self): """Stream a series of solid colours""" colours = ["#F00", "#0F0", "#00F"] - jpegs = [] - for c in colours: - image = Image.new("RGB", (10, 10), c) - dest = io.BytesIO() - image.save(dest, "jpeg") - jpegs.append(dest.getvalue()) + jpegs = [make_jpeg(c) for c in colours] if self.initial_delay > 0: time.sleep(self.initial_delay) @@ -150,6 +155,55 @@ def test_grab_and_shutdown(server: lt.ThingServer, telly: Telly): # The background thread gets shut down by `Telly.__exit__`. +@pytest.fixture +def mjpeg_stream(mocker): + """An MJPEGStream object with a mocked Thing.""" + thing = mocker.Mock() + broker = MessageBroker() + tsi = thing._thing_server_interface + tsi.message_broker = broker + tsi.start_async_task_soon.side_effect = from_thread.run + + def publish(message: Message) -> None: + from_thread.run(broker.publish, message) + + tsi.publish.side_effect = publish + thing.name = "thing" + yield MJPEGStream(thing, "stream") + + +@pytest.fixture +def jpeg_bytes(): + return make_jpeg("#F0F") + + +async def test_grab_frame(mjpeg_stream, jpeg_bytes): + """Test we can add a frame in a thread and grab its bytes.""" + # Listen for the next frame (grab_frame) then add a frame. + async with anyio.create_task_group() as tg: + handle = tg.start_soon(mjpeg_stream.grab_frame) + await to_thread.run_sync(mjpeg_stream.add_frame, jpeg_bytes) + assert handle.return_value is jpeg_bytes + + +async def test_grab_frame_with_metadata(mjpeg_stream, jpeg_bytes): + """Test we can add a frame in a thread and grab it.""" + # Listen for the next frame (grab_frame) then add a frame. + async with anyio.create_task_group() as tg: + handle = tg.start_soon(mjpeg_stream.grab_frame_with_metadata) + await to_thread.run_sync(mjpeg_stream.add_frame, jpeg_bytes) + assert handle.return_value.frame is jpeg_bytes + + +async def test_next_frame_size(mjpeg_stream, jpeg_bytes): + """Test we can add a frame in a thread and grab its size.""" + # Listen for the next frame (grab_frame) then add a frame. + async with anyio.create_task_group() as tg: + handle = tg.start_soon(mjpeg_stream.next_frame_size) + await to_thread.run_sync(mjpeg_stream.add_frame, jpeg_bytes) + assert handle.return_value == len(jpeg_bytes) + + def test_mjpeg_stream_http(server: lt.ThingServer, telly: Telly): """Verify the MJPEG stream works, and is shut down cleanly.