Files
app/backend/tests/flow/test_senders.py
T
stroblmeandClaude Opus 5 640654bd66 Rename the import package app to fluksio
A wheel whose top-level module is `app` collides with anything else in a
user's venv, so the package that is about to be published takes the name
it is published under. Only the Python package moves; the repo, the
Docker WORKDIR and the compose project keep theirs.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-21 21:48:05 +02:00

92 lines
2.9 KiB
Python

"""The outbound nodes reuse one connection instead of opening one per message."""
import asyncio
from fluksio.flow.messages import DType, MessageSpec
from fluksio.flow.nodes import MqttNode
from fluksio.flow.nodes.http import close_shared_client, shared_client
def test_http_senders_share_one_pooled_client():
first = shared_client()
try:
assert shared_client() is first
finally:
close_shared_client()
# Closing lets the next request build a fresh one rather than reusing a
# closed pool.
assert shared_client() is not first
close_shared_client()
def _publisher() -> MqttNode:
node = MqttNode(
requires=[MessageSpec(name="setpoint", port="setpoint", dtype=DType.FLOAT)],
params={"topic": {"setpoint": "heating/setpoint"}},
)
node.assign_flow("heating", "out")
return node
def test_a_started_publisher_queues_instead_of_connecting():
"""The handler runs on a worker thread; it must not block on the broker."""
node = _publisher()
async def scenario() -> None:
node._publish_queue = asyncio.Queue(maxsize=4)
node._loop = asyncio.get_running_loop()
await asyncio.to_thread(node._publisher_handler, {}, setpoint=21.0)
# call_soon_threadsafe lands on the next loop pass.
await asyncio.sleep(0)
assert node._publish_queue.qsize() == 1
assert node._publish_queue.get_nowait() == {"setpoint": 21.0}
asyncio.run(scenario())
def test_a_full_publish_queue_drops_the_oldest():
"""A broker that cannot keep up must not grow the queue without bound."""
node = _publisher()
health: list[tuple[str, str | None]] = []
node._on_health = lambda _n, status, detail: health.append((status, detail))
async def scenario() -> None:
queue: asyncio.Queue[dict] = asyncio.Queue(maxsize=2)
for value in (1.0, 2.0, 3.0):
node._enqueue(queue, {"setpoint": value})
assert queue.qsize() == 2
assert queue.get_nowait() == {"setpoint": 2.0}
assert queue.get_nowait() == {"setpoint": 3.0}
assert health == [("degraded", "publish queue full")]
asyncio.run(scenario())
def test_a_string_goes_on_the_wire_bare():
"""Devices on a shared broker expect `ON`, not `"ON"`."""
class Recorder:
def __init__(self) -> None:
self.published: list[tuple[str, str]] = []
async def publish(self, topic, payload, **_):
self.published.append((topic, payload))
node = MqttNode(
requires=[
MessageSpec(name="plug", port="plug", dtype=DType.STR),
MessageSpec(name="level", port="level", dtype=DType.INT),
],
params={"topic": {"plug": "actor/plug", "level": "light/level"}},
)
node.assign_flow("house", "out")
client = Recorder()
asyncio.run(node._publish_with(client, {"plug": "ON", "level": 60}))
assert client.published == [("actor/plug", "ON"), ("light/level", "60")]