Files
bambuddy/backend/tests/unit/test_tls_proxy_teardown_2968.py
maziggy 0dfcff5925 Keep the RTSPS proxy's handler set off the server object (issue #3001)
asyncio's Server has a __dict__ and uvloop's, a Cython cdef class, does
not, so the attribute added in 1.2.5.4 raised AttributeError under uvloop.
Every RTSP camera failed before opening a socket, which is the
diagnostic's capture_exception at 0 ms.

Our own unit files all pin --loop asyncio for #1896 and were never
affected. The reports come from units we do not write: the Proxmox VE
Helper-Scripts LXC pins no loop, and installs predating that fix never
gained the flag because update.sh does not rewrite unit files. The loop
is not ours to assume, so fix the code rather than add another flag.

The set moves to a module-level WeakKeyDictionary, keyed weakly so an
abandoned proxy retires its own entry rather than leaking one and later
handing a new server a dead one's handlers.

Pinned on a real uvloop loop and, for hosts without uvloop, against a
__slots__ server; conftest builds its loop from the default policy, so
nothing in the suite had ever run the branch that broke.

Also routes the two external-camera teardowns through close_tls_proxy,
which #2968 introduced and left them out of.

-----

Say so at startup when running on uvloop (issue #3001)

An install on the wrong loop had no way to find out it was. #3001 was
loud enough to notice; the #1896 upload truncation it is also exposed to
is silent, and shows up as a print failing from a file that was corrupt
on arrival.

One WARNING in the lifespan naming the loop, the risk and the flag to
add. A warning and not a refusal: uvicorn has already chosen its loop by
the time any application code runs, and a server that answers requests
beats one that will not boot.

Asks the running loop what it is rather than whether uvloop imports --
uvicorn[standard] installs uvloop everywhere, so its presence says
nothing -- and matches on the module name so the question never imports
uvloop on a host without it.

-----

Repair a service file written before the --loop asyncio pin (issue #3001)

install.sh has pinned the loop since #1896, but nothing has ever
rewritten an existing service file, so every native install created
between 2025-11-28 (when uvicorn[standard] brought uvloop into the venv)
and 2026-07-05 still runs on uvloop no matter how often it is updated.

Both update scripts now add the flag themselves while the service is
stopped, so it takes effect on the same restart -- systemd via sed,
launchd via PlistBuddy, each backing the file up first and inserting
nothing but the flag.

Refuses to edit and explains instead when the shape is not a plain
single-line uvicorn unit: a wrapper script, a continued ExecStart,
several of them, a read-only file, or a service with drop-ins, since a
drop-in may be what defines ExecStart and editing the fragment would
change nothing while reporting success. A deliberate --loop uvloop is
left alone. Reads the effective ExecStart from systemd rather than the
file, so it is idempotent.
2026-08-30 08:01:42 +02:00

267 lines
10 KiB
Python

"""The RTSPS proxy must not leave a handler running past its server (#2968).
The reporter's log carries three of these, one per camera snapshot, at ERROR
with a traceback pointing into ``camera.py``:
ERROR [asyncio] Task was destroyed but it is pending!
task: <Task pending name='Task-1889625'
coro=<create_tls_proxy.<locals>._handle() done, defined at camera.py:243>
wait_for=<_GatheringFuture pending ...>>
``asyncio.start_server`` wraps the connection callback in a task and keeps only
a weak reference to it, so a handler still awaiting its two forwarders can be
collected while pending -- which is exactly what that message is. Nothing was
broken by it (the snapshot on either side of each one succeeded), but it reads
like a camera fault in a log people attach to bug reports, and the shape behind
it is real: teardown closed the listener and then waited on handlers that only
finish when the *peer* drops the socket.
Two things fix it. The handlers are strongly referenced for as long as they run,
and ``close_tls_proxy`` cancels them rather than hoping ffmpeg has already gone.
The upstream here is a real TLS listener rather than a bare socket, because the
proxy spends its first ten seconds inside ``open_connection``: a stand-in that
never completes a handshake never reaches the forwarding state these tests are
about. The proxy sets ``CERT_NONE`` (Bambu printers are self-signed), so a
throwaway certificate is all it takes.
"""
from __future__ import annotations
import asyncio
import datetime
import gc
import logging
import ssl
import pytest
from backend.app.services.camera import _proxy_handlers, close_tls_proxy, create_tls_proxy
@pytest.fixture(scope="module")
def self_signed_cert(tmp_path_factory):
"""Certificate and key for the stand-in printer, generated once."""
from cryptography import x509
from cryptography.hazmat.primitives import hashes, serialization
from cryptography.hazmat.primitives.asymmetric import rsa
from cryptography.x509.oid import NameOID
key = rsa.generate_private_key(public_exponent=65537, key_size=2048)
name = x509.Name([x509.NameAttribute(NameOID.COMMON_NAME, "127.0.0.1")])
now = datetime.datetime.now(datetime.timezone.utc)
cert = (
x509.CertificateBuilder()
.subject_name(name)
.issuer_name(name)
.public_key(key.public_key())
.serial_number(x509.random_serial_number())
.not_valid_before(now - datetime.timedelta(days=1))
.not_valid_after(now + datetime.timedelta(days=1))
.sign(key, hashes.SHA256())
)
directory = tmp_path_factory.mktemp("tls")
cert_file = directory / "cert.pem"
key_file = directory / "key.pem"
cert_file.write_bytes(cert.public_bytes(serialization.Encoding.PEM))
key_file.write_bytes(
key.private_bytes(
encoding=serialization.Encoding.PEM,
format=serialization.PrivateFormat.TraditionalOpenSSL,
encryption_algorithm=serialization.NoEncryption(),
)
)
return cert_file, key_file
async def _printer(self_signed_cert, on_data=None) -> tuple[asyncio.Server, int]:
"""A TLS listener standing in for the printer's RTSPS port.
Accepts, hands anything it receives to *on_data*, and otherwise waits --
which is the state the upstream is in while ffmpeg is being reaped.
"""
async def _accept(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
try:
while True:
data = await reader.read(4096)
if not data:
break
if on_data is not None:
on_data(data)
except (ConnectionError, OSError, asyncio.CancelledError):
pass
finally:
if not writer.is_closing():
writer.close()
cert_file, key_file = self_signed_cert
context = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
context.load_cert_chain(str(cert_file), str(key_file))
server = await asyncio.start_server(_accept, "127.0.0.1", 0, ssl=context)
return server, server.sockets[0].getsockname()[1]
async def _close(proxy) -> None:
"""Teardown, bounded.
Every close in this file goes through the timeout, including the ones in
``finally`` blocks that are only there to tidy up. Losing the cancellation
or the handler tracking makes ``close_tls_proxy`` wait on a peer that is
not going to drop, and an unbounded await turns that regression into a
hung suite instead of a failing test.
"""
await asyncio.wait_for(close_tls_proxy(proxy), timeout=5.0)
async def _shutdown(server: asyncio.Server) -> None:
"""Bounded teardown for the stand-in printer.
``wait_closed`` waits for the listener's own handlers, and one of those is
reading a socket the proxy still holds. Left unbounded it inherits any
regression in the proxy's teardown and hangs the suite in a second place.
"""
server.close()
try:
await asyncio.wait_for(server.wait_closed(), timeout=5.0)
except asyncio.TimeoutError:
pass
async def _wait_for(predicate, timeout: float = 5.0) -> bool:
"""Poll rather than sleep a fixed amount: these are real sockets."""
deadline = asyncio.get_running_loop().time() + timeout
while asyncio.get_running_loop().time() < deadline:
if predicate():
return True
await asyncio.sleep(0.02)
return predicate()
class TestTheHandlerIsHeldWhileItRuns:
@pytest.mark.asyncio
async def test_an_open_connection_is_tracked(self, self_signed_cert):
"""The set is the strong reference asyncio does not keep."""
upstream, upstream_port = await _printer(self_signed_cert)
try:
port, proxy = await create_tls_proxy("127.0.0.1", upstream_port)
try:
_, writer = await asyncio.open_connection("127.0.0.1", port)
assert await _wait_for(lambda: len(_proxy_handlers[proxy]) == 1)
assert not next(iter(_proxy_handlers[proxy])).done()
writer.close()
finally:
await _close(proxy)
finally:
await _shutdown(upstream)
@pytest.mark.asyncio
async def test_a_finished_handler_is_released(self, self_signed_cert):
"""Tracked for the connection's life, not the process's -- a long
stream must not accumulate one entry per reconnect."""
upstream, upstream_port = await _printer(self_signed_cert)
try:
port, proxy = await create_tls_proxy("127.0.0.1", upstream_port)
try:
_, writer = await asyncio.open_connection("127.0.0.1", port)
assert await _wait_for(lambda: len(_proxy_handlers[proxy]) == 1)
writer.close()
assert await _wait_for(lambda: _proxy_handlers[proxy] == set())
finally:
await _close(proxy)
finally:
await _shutdown(upstream)
class TestCloseDoesNotDependOnThePeer:
@pytest.mark.asyncio
async def test_a_live_connection_does_not_stall_the_close(self, self_signed_cert):
"""``server.close()`` leaves established connections running, so the
old close/wait pair finished only when the client happened to drop.
Here the client is still attached and close still returns."""
upstream, upstream_port = await _printer(self_signed_cert)
try:
port, proxy = await create_tls_proxy("127.0.0.1", upstream_port)
_, writer = await asyncio.open_connection("127.0.0.1", port)
assert await _wait_for(lambda: len(_proxy_handlers[proxy]) == 1)
await _close(proxy)
# Stronger than "the set is empty": since #3001 the handler set
# lives in a module-level registry rather than on the server, and
# close_tls_proxy retires the whole entry.
assert proxy not in _proxy_handlers
writer.close()
finally:
await _shutdown(upstream)
@pytest.mark.asyncio
async def test_no_handler_survives_the_close(self, self_signed_cert, caplog):
"""The actual complaint: nothing is left pending for the garbage
collector to shout about afterwards."""
upstream, upstream_port = await _printer(self_signed_cert)
try:
port, proxy = await create_tls_proxy("127.0.0.1", upstream_port)
_, writer = await asyncio.open_connection("127.0.0.1", port)
assert await _wait_for(lambda: len(_proxy_handlers[proxy]) == 1)
handler = next(iter(_proxy_handlers[proxy]))
with caplog.at_level(logging.ERROR, logger="asyncio"):
await _close(proxy)
writer.close()
await asyncio.sleep(0.05)
gc.collect()
await asyncio.sleep(0.05)
assert handler.done()
assert not [r for r in caplog.records if "Task was destroyed" in r.getMessage()]
finally:
await _shutdown(upstream)
@pytest.mark.asyncio
async def test_closing_twice_is_harmless(self, self_signed_cert):
"""Both callers reach their finally block on the error paths too."""
upstream, upstream_port = await _printer(self_signed_cert)
try:
_, proxy = await create_tls_proxy("127.0.0.1", upstream_port)
await _close(proxy)
await _close(proxy)
finally:
await _shutdown(upstream)
@pytest.mark.asyncio
async def test_it_works_on_a_server_it_did_not_create(self):
"""Degrades to the close/wait it replaces rather than raising."""
plain = await asyncio.start_server(lambda r, w: None, "127.0.0.1", 0)
await _close(plain)
assert not plain.is_serving()
@pytest.mark.asyncio
async def test_the_proxy_still_forwards(self_signed_cert):
"""The teardown changes must not cost the proxy its job: plain TCP in one
end, TLS to the printer out the other."""
received: list[bytes] = []
upstream, upstream_port = await _printer(self_signed_cert, on_data=received.append)
try:
port, proxy = await create_tls_proxy("127.0.0.1", upstream_port)
try:
_, writer = await asyncio.open_connection("127.0.0.1", port)
writer.write(b"OPTIONS rtsp://127.0.0.1/streaming/live/1 RTSP/1.0\r\n\r\n")
await writer.drain()
assert await _wait_for(lambda: bool(received))
assert b"OPTIONS" in received[0]
writer.close()
finally:
await _close(proxy)
finally:
await _shutdown(upstream)