Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 33 additions & 1 deletion src/labthings_fastapi/message_broker.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@

thing: str
affordance: str
message_type: Literal["property", "action"]
message_type: Literal["property", "action", "stream"]
payload: Any


Expand Down Expand Up @@ -83,6 +83,19 @@
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:
Expand Down Expand Up @@ -148,3 +161,22 @@
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.

Check warning on line 177 in src/labthings_fastapi/message_broker.py

View workflow job for this annotation

GitHub Actions / coverage

176-177 lines are not covered with tests
async with anyio.create_task_group() as tg:
n = len(streams)
for stream in streams:
tg.start_soon(stream.aclose)
return n
Loading
Loading