diff --git a/solstone/convey/secure_listener/admission.py b/solstone/convey/secure_listener/admission.py index feea83860..58d58c96b 100644 --- a/solstone/convey/secure_listener/admission.py +++ b/solstone/convey/secure_listener/admission.py @@ -21,20 +21,31 @@ log = logging.getLogger("convey.secure_listener.admission") DEFAULT_SECURE_LISTENER_CAPACITY: Final[int] = 16 DEFAULT_SECURE_LISTENER_STREAMING_CAPACITY: Final[int] = 8 +DEFAULT_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS: Final[float] = 120.0 +SECURE_LISTENER_QUEUE_WARN_SECONDS: Final[float] = 60.0 MAX_SECURE_LISTENER_CAPACITY: Final[int] = 128 +MAX_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS: Final[float] = 600.0 _PermitKind = Literal["total", "streaming", "streaming_over_budget"] +_DepartureReason = Literal["granted", "cancelled", "timed_out"] class SecureListenerAdmissionRejected(Exception): """The listener refused admission before executor submission.""" +class SecureListenerQueueTimeout(SecureListenerAdmissionRejected): + """A queued listener request exceeded the admission wait deadline.""" + + reason_code = "secure_listener_queue_timeout" + + @dataclass(frozen=True) class SecureListenerAdmissionConfig: capacity: int = DEFAULT_SECURE_LISTENER_CAPACITY streaming_capacity: int = DEFAULT_SECURE_LISTENER_STREAMING_CAPACITY refuse_when_full: bool = False + queue_timeout_seconds: float = DEFAULT_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS @property def queue_limit(self) -> int: @@ -146,10 +157,19 @@ def resolve_admission_config() -> SecureListenerAdmissionConfig: default=False, warning_default="false", ) + queue_timeout_seconds = _resolve_float( + link_cfg, + "secure_listener_queue_timeout_seconds", + default=DEFAULT_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS, + valid_min=1.0, + valid_max=MAX_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS, + warning_default="120.0", + ) return SecureListenerAdmissionConfig( capacity=capacity, streaming_capacity=streaming_capacity, refuse_when_full=refuse_when_full, + queue_timeout_seconds=queue_timeout_seconds, ) @@ -198,6 +218,47 @@ def _resolve_bool( return raw +def _resolve_float( + link_cfg: dict[str, Any], + key: str, + *, + default: float, + valid_min: float, + valid_max: float, + warning_default: str, +) -> float: + raw = link_cfg.get(key, default) + if isinstance(raw, bool): + log.warning( + "Invalid link.%s in journal config: %r \u2014 defaulting to %s", + key, + raw, + warning_default, + ) + return default + if raw == 0: + log.info("link.%s is 0; secure listener queue timeout disabled", key) + return 0.0 + if not isinstance(raw, (int, float)): + log.warning( + "Invalid link.%s in journal config: %r \u2014 defaulting to %s", + key, + raw, + warning_default, + ) + return default + value = float(raw) + if not (valid_min <= value <= valid_max): + log.warning( + "Invalid link.%s in journal config: %r \u2014 defaulting to %s", + key, + raw, + warning_default, + ) + return default + return value + + class SecureListenerAdmission: """Admission, queueing, and content-free telemetry for listener work.""" @@ -220,6 +281,7 @@ class SecureListenerAdmission: self._active_streaming_over_budget = 0 self._rejected_total = 0 self._rejected_streaming = 0 + self._rejected_queue_timeout = 0 self._admitted_streaming_over_budget = 0 async def acquire(self) -> SecureListenerPermit: @@ -241,15 +303,35 @@ class SecureListenerAdmission: queued_at=started, ) self._waiters.append(waiter) + queue_timeout_seconds = self.config.queue_timeout_seconds + departure_reason: _DepartureReason = "cancelled" try: - return await waiter.future + if queue_timeout_seconds > 0.0: + permit = await asyncio.wait_for( + waiter.future, + timeout=queue_timeout_seconds, + ) + else: + permit = await waiter.future + except TimeoutError as exc: + departure_reason = "timed_out" + self._reclaim_abandoned_waiter(waiter) + self._record_queue_timeout_rejection() + raise SecureListenerQueueTimeout from exc except BaseException: - if self._cancel_waiter(waiter): - raise - if waiter.permit is not None: - waiter.permit.release() + departure_reason = "cancelled" + self._reclaim_abandoned_waiter(waiter) raise + else: + departure_reason = "granted" + return permit + finally: + self._warn_if_slow_waiter_departure( + waiter, + departure_reason, + queue_timeout_seconds, + ) async def submit( self, @@ -340,6 +422,7 @@ class SecureListenerAdmission: "rejected": { "total": self._rejected_total, "streaming": self._rejected_streaming, + "queue_timeout": self._rejected_queue_timeout, }, "admitted_over_budget": { "streaming": self._admitted_streaming_over_budget, @@ -386,6 +469,40 @@ class SecureListenerAdmission: return False return True + def _reclaim_abandoned_waiter(self, waiter: _QueuedWaiter) -> None: + if self._cancel_waiter(waiter): + return + if waiter.permit is not None: + waiter.permit.release() + + def _record_queue_timeout_rejection(self) -> None: + with self._lock: + self._rejected_total += 1 + self._rejected_queue_timeout += 1 + + def _warn_if_slow_waiter_departure( + self, + waiter: _QueuedWaiter, + reason: _DepartureReason, + queue_timeout_seconds: float, + ) -> None: + waiter_age_s = time.monotonic() - waiter.queued_at + if waiter_age_s <= SECURE_LISTENER_QUEUE_WARN_SECONDS: + return + with self._lock: + active_total = self._active_total + queue_depth = len(self._waiters) + log.warning( + "Secure listener admission waiter departed departure_reason=%s " + "waiter_age_s=%.3f active_total=%d queue_depth=%d " + "queue_timeout_seconds=%.3f", + reason, + waiter_age_s, + active_total, + queue_depth, + queue_timeout_seconds, + ) + def _wake_waiters_locked(self) -> None: while self._waiters and self._active_total < self.config.capacity: waiter = self._waiters.popleft() diff --git a/solstone/convey/secure_listener/wsgi.py b/solstone/convey/secure_listener/wsgi.py index 21db76208..01ada25bb 100644 --- a/solstone/convey/secure_listener/wsgi.py +++ b/solstone/convey/secure_listener/wsgi.py @@ -20,7 +20,11 @@ from werkzeug.exceptions import HTTPException from solstone.think.link.window import window_open -from .admission import SecureListenerAdmission, SecureListenerAdmissionRejected +from .admission import ( + SecureListenerAdmission, + SecureListenerAdmissionRejected, + SecureListenerQueueTimeout, +) from .identity import ConveyIdentity from .mux import ( RESET_CTX_APP_CANCELLATION, @@ -430,19 +434,27 @@ async def dispatch_stream( wsgi_input = environ["wsgi.input"] if wsgi_input.remaining > 0: stream_writer.begin_drain(RESET_CTX_APP_CANCELLATION) - except SecureListenerAdmissionRejected: - await write_json_response( - stream_writer, - 503, - "Service Unavailable", - {"error": "secure listener capacity is full"}, - extra_headers=( - ( - "Retry-After", - str(SECURE_LISTENER_REFUSAL_RETRY_AFTER_SECONDS), - ), - ), + except SecureListenerAdmissionRejected as exc: + body = ( + {"error": "secure listener queue timeout"} + if isinstance(exc, SecureListenerQueueTimeout) + else {"error": "secure listener capacity is full"} ) + try: + await write_json_response( + stream_writer, + 503, + "Service Unavailable", + body, + extra_headers=( + ( + "Retry-After", + str(SECURE_LISTENER_REFUSAL_RETRY_AFTER_SECONDS), + ), + ), + ) + except ConnectionError: + pass stream_writer.begin_drain(RESET_CTX_BODY_DISCARD_CANCELLATION) return DispatchResult(endpoint=endpoint, status=503) except asyncio.CancelledError: diff --git a/solstone/think/link/README.md b/solstone/think/link/README.md index e30c99da2..91a2cab52 100644 --- a/solstone/think/link/README.md +++ b/solstone/think/link/README.md @@ -32,7 +32,7 @@ The `spl` repo's `home/` continues as the open-source reference implementation o TLS termination, multiplexing, and inline WSGI dispatch now live in `solstone/convey/secure_listener/`, because Convey owns both listening ports: the DL web port and the PL secure-listener port 7657. -Secure-listener capacity is configured with `link.secure_listener_capacity`; `link.secure_listener_streaming_capacity = 0` disables the streaming lane split. +Secure-listener capacity is configured with `link.secure_listener_capacity`; `link.secure_listener_streaming_capacity = 0` disables the streaming lane split; `link.secure_listener_queue_timeout_seconds` bounds queued admission waits, accepts 1.0-600.0 seconds, and `0` disables queue-timeout refusal. ## naming diff --git a/tests/link/secure_listener_harness.py b/tests/link/secure_listener_harness.py index 64742cc24..842d1d404 100644 --- a/tests/link/secure_listener_harness.py +++ b/tests/link/secure_listener_harness.py @@ -4,6 +4,7 @@ from __future__ import annotations import asyncio +import contextlib from dataclasses import dataclass from pathlib import Path from typing import Any @@ -21,10 +22,22 @@ from solstone.convey.secure_listener.tls import ( issue_server_cert, ) from solstone.think.link.auth import AuthorizedClients -from solstone.think.link.ca import LoadedCa, load_or_generate_ca +from solstone.think.link.ca import LoadedCa, cert_fingerprint, load_or_generate_ca +from solstone.think.link.client import ( + ClientIdentity, + TunnelSession, + _open_tunnel_session, + _TcpEncryptedTransport, +) from solstone.think.link.nonces import NonceStore from solstone.think.link.paths import authorized_clients_path, ca_dir, nonces_path -from tests.link.certless_helpers import make_convey_app +from tests.link.certless_helpers import ( + DirectPairCandidate, + DirectPairRequest, + build_csr, + make_convey_app, + post_pair_framed, +) @dataclass @@ -45,6 +58,7 @@ class SecureListenerHarness: monkeypatch: pytest.MonkeyPatch, *, link: dict[str, Any] | None = None, + admission_config: SecureListenerAdmissionConfig | None = None, ) -> SecureListenerHarness: app, journal = make_convey_app( tmp_path, @@ -67,7 +81,8 @@ class SecureListenerHarness: authorized, ) admission = SecureListenerAdmission( - SecureListenerAdmissionConfig( + admission_config + or SecureListenerAdmissionConfig( capacity=4, streaming_capacity=4, refuse_when_full=False, @@ -120,3 +135,44 @@ class SecureListenerHarness: def pair_url(self, nonce: str, *, path: str = "/app/network/pair") -> str: return f"https://{self.host}:{self.port}{path}?token={nonce}" + + +async def pair_and_open_session( + harness: SecureListenerHarness, + *, + nonce: str, + label: str, +) -> TunnelSession: + private_key, private_key_pem, csr_pem = build_csr(label) + harness.seed_nonce(nonce, label) + response = await asyncio.to_thread( + post_pair_framed, + DirectPairRequest( + candidates=(DirectPairCandidate(harness.host, harness.port),), + path=f"/app/network/pair?token={nonce}", + ca_fingerprint_pin=harness.ca.fingerprint_sha256(), + ), + {"csr": csr_pem, "device_label": label}, + private_key, + ) + identity = ClientIdentity( + private_key_pem=private_key_pem.decode("ascii"), + client_cert_pem=response.client_cert, + ca_chain_pem="".join(response.ca_chain), + fingerprint=cert_fingerprint(response.client_cert), + home_instance_id=response.instance_id, + home_label=response.home_label, + home_attestation=response.home_attestation, + local_endpoints=tuple(response.local_endpoints), + ) + reader, writer = await asyncio.open_connection(harness.host, harness.port) + try: + return await _open_tunnel_session( + _TcpEncryptedTransport(reader, writer), + identity, + ) + except BaseException: + writer.close() + with contextlib.suppress(Exception): + await writer.wait_closed() + raise diff --git a/tests/link/test_mux.py b/tests/link/test_mux.py index b1c7e0237..43aadbcce 100644 --- a/tests/link/test_mux.py +++ b/tests/link/test_mux.py @@ -2301,7 +2301,7 @@ async def test_capacity_snapshot_is_readable_while_workers_are_occupied( } assert set(payload["limit"]) == {"total", "streaming", "queue"} assert set(payload["queued"]) == {"total"} - assert set(payload["rejected"]) == {"total", "streaming"} + assert set(payload["rejected"]) == {"total", "streaming", "queue_timeout"} assert set(payload["admitted_over_budget"]) == {"streaming"} assert payload["active"]["total"] == admission.config.capacity assert payload["limit"]["total"] == admission.config.capacity diff --git a/tests/link/test_secure_listener_queue_bound.py b/tests/link/test_secure_listener_queue_bound.py new file mode 100644 index 000000000..0bb1be140 --- /dev/null +++ b/tests/link/test_secure_listener_queue_bound.py @@ -0,0 +1,997 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import asyncio +import contextlib +import logging +import re +import threading +from collections.abc import Callable +from pathlib import Path +from typing import Any + +import pytest +from flask import Response + +from solstone.convey import root as root_module +from solstone.convey.secure_listener import admission as admission_module +from solstone.convey.secure_listener import mux as mux_module +from solstone.convey.secure_listener import wsgi as wsgi_module +from solstone.convey.secure_listener.admission import ( + DEFAULT_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS, + SECURE_LISTENER_QUEUE_WARN_SECONDS, + SecureListenerAdmission, + SecureListenerAdmissionConfig, +) +from solstone.convey.secure_listener.framing import ( + FLAG_DATA, + FLAG_RESET, + INITIAL_WINDOW, + RESET_CANCEL, + Frame, + FrameDecoder, + build_close, + build_data, + build_open, + build_reset, + parse_reset_reason, +) +from solstone.convey.secure_listener.mux import ( + RESET_CTX_BODY_DISCARD_CANCELLATION, + RESET_CTX_HANDLER_EXCEPTION, + Multiplexer, +) +from solstone.convey.secure_listener.wsgi import DispatchResult, dispatch_stream +from solstone.think.link.auth import AuthorizedClients +from solstone.think.link.client import _http_head_bytes, _parse_http_response +from solstone.think.link.dialer import _REQUEST_TIMEOUT_SECONDS +from solstone.think.link.paths import authorized_clients_path +from tests.link.certless_helpers import FakeStreamWriter, make_convey_app, pl_identity +from tests.link.secure_listener_harness import ( + SecureListenerHarness, + pair_and_open_session, +) + +STATE_TIMEOUT_S = 3.0 +RESPONSE_TIMEOUT_S = 3.0 +POLL_INTERVAL_S = 0.005 +BLOCK_GUARD_TIMEOUT_S = 10.0 + +FINGERPRINT = "sha256:" + ("7" * 64) +QUEUE_TIMEOUT_BODY = b'{"error":"secure listener queue timeout"}' +CAPACITY_REFUSAL_BODY = b'{"error":"secure listener capacity is full"}' +STREAMING_REFUSAL_BODY = b'{"error":"secure listener streaming capacity is full"}\n' +OK_BODY = b"queue-ok" + + +def _decode_frames(chunks: list[bytes]) -> list[Frame]: + decoder = FrameDecoder() + for chunk in chunks: + decoder.feed(chunk) + return decoder.drain() + + +def _authorize_fingerprint( + monkeypatch: pytest.MonkeyPatch, + fingerprint: str = FINGERPRINT, +) -> None: + authorized = AuthorizedClients(authorized_clients_path()) + authorized.add(fingerprint, "pytest phone", "inst-1") + monkeypatch.setattr(root_module, "get_authorized_clients", lambda: authorized) + + +def _admission( + *, + capacity: int = 1, + streaming_capacity: int | None = None, + queue_timeout_seconds: float = 0.05, + refuse_when_full: bool = False, +) -> SecureListenerAdmission: + return SecureListenerAdmission( + SecureListenerAdmissionConfig( + capacity=capacity, + streaming_capacity=capacity + if streaming_capacity is None + else streaming_capacity, + refuse_when_full=refuse_when_full, + queue_timeout_seconds=queue_timeout_seconds, + ) + ) + + +async def _shutdown_admission(admission: SecureListenerAdmission) -> None: + await asyncio.to_thread(admission.shutdown, wait=True, cancel_futures=True) + + +async def _wait_for_snapshot( + admission: SecureListenerAdmission, + predicate: Callable[[dict[str, Any]], bool], +) -> dict[str, Any]: + loop = asyncio.get_running_loop() + deadline = loop.time() + STATE_TIMEOUT_S + while True: + snapshot = admission.snapshot() + if predicate(snapshot): + return snapshot + if loop.time() >= deadline: + raise AssertionError(f"snapshot predicate not met: {snapshot!r}") + await asyncio.sleep(POLL_INTERVAL_S) + + +async def _wait_for_mux_state(mux: Multiplexer, stream_id: int) -> Any: + loop = asyncio.get_running_loop() + deadline = loop.time() + STATE_TIMEOUT_S + while True: + state = mux._streams.get(stream_id) + if state is not None: + return state + if loop.time() >= deadline: + raise AssertionError(f"stream {stream_id} did not open") + await asyncio.sleep(POLL_INTERVAL_S) + + +async def _wait_for_stream_payload(sent: list[bytes], stream_id: int) -> bytes: + loop = asyncio.get_running_loop() + deadline = loop.time() + RESPONSE_TIMEOUT_S + while True: + payload = b"".join( + frame.payload + for frame in _decode_frames(sent) + if frame.stream_id == stream_id and frame.flags & FLAG_DATA + ) + if payload: + return payload + if loop.time() >= deadline: + raise AssertionError(f"stream {stream_id} did not emit DATA") + await asyncio.sleep(POLL_INTERVAL_S) + + +async def _dispatch_raw_request( + app: Any, + admission: SecureListenerAdmission, + method: str, + path: str, + *, + body: bytes = b"", + headers: dict[str, str] | None = None, + writer: FakeStreamWriter | None = None, +) -> tuple[DispatchResult, int, dict[str, str], bytes, FakeStreamWriter]: + reader = asyncio.StreamReader() + reader.feed_data( + _http_head_bytes( + method, + path, + headers=headers, + content_length=len(body), + ) + + body + ) + reader.feed_eof() + stream_writer = writer or FakeStreamWriter() + result = await dispatch_stream( + app, + pl_identity(FINGERPRINT), + reader, + stream_writer, + asyncio.get_running_loop(), + admission, + ) + status, response_headers, response_body = _parse_http_response( + bytes(stream_writer.data) + ) + assert result.status == status + return result, status, response_headers, response_body, stream_writer + + +def _register_basic_endpoints( + app: Any, + release_hold: threading.Event | None = None, + *, + invoked: dict[str, int] | None = None, +) -> None: + @app.get("/_queue_bound/hold") + def queue_bound_hold() -> Response: + if release_hold is not None: + release_hold.wait(timeout=BLOCK_GUARD_TIMEOUT_S) + return Response(b"held", content_type="text/plain") + + @app.get("/_queue_bound/ok") + def queue_bound_ok() -> Response: + if invoked is not None: + invoked["count"] = invoked.get("count", 0) + 1 + return Response(OK_BODY, content_type="text/plain") + + @app.post("/_queue_bound/upload") + def queue_bound_upload() -> Response: + if invoked is not None: + invoked["count"] = invoked.get("count", 0) + 1 + return Response(b"uploaded", content_type="text/plain") + + +def _register_streaming_endpoint(app: Any) -> None: + @app.get("/_queue_bound/events") + def queue_bound_events() -> Response: + return Response(iter((b"data: queue-bound\n\n",)), mimetype="text/event-stream") + + +async def _open_mux_request( + mux: Multiplexer, + stream_id: int, + method: str, + path: str, + *, + headers: dict[str, str] | None = None, + content_length: int = 0, +) -> None: + head = _http_head_bytes( + method, + path, + headers=headers, + content_length=content_length, + ) + suffix = build_close(stream_id).encode() if content_length == 0 else b"" + await mux.feed(build_open(stream_id, head).encode() + suffix) + + +async def _open_mux_get(mux: Multiplexer, stream_id: int, path: str) -> None: + await _open_mux_request(mux, stream_id, "GET", path) + + +def _assert_content_free(value: bytes | str) -> None: + text = value.decode("utf-8", "replace") if isinstance(value, bytes) else value + assert "sha256" not in text + assert re.search(r"\b\d{1,3}(?:\.\d{1,3}){3}\b", text) is None + + +@pytest.mark.asyncio +async def test_queue_timeout_returns_503_retry_after_and_distinct_body( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + release_hold = threading.Event() + _register_basic_endpoints(app, release_hold) + admission = _admission() + hold_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/hold") + ) + try: + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["active"]["total"] == 1, + ) + + _result, status, headers, body, _writer = await asyncio.wait_for( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok"), + timeout=RESPONSE_TIMEOUT_S, + ) + + assert status == 503 + assert headers["retry-after"] == str( + wsgi_module.SECURE_LISTENER_REFUSAL_RETRY_AFTER_SECONDS + ) + assert body == QUEUE_TIMEOUT_BODY + assert body != CAPACITY_REFUSAL_BODY + assert body != STREAMING_REFUSAL_BODY + _assert_content_free(body) + finally: + release_hold.set() + with contextlib.suppress(Exception): + await hold_task + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_queue_timeout_clears_queue_counts_and_keeps_streaming_separate( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + release_hold = threading.Event() + _register_basic_endpoints(app, release_hold) + _register_streaming_endpoint(app) + admission = _admission(streaming_capacity=1) + held_streaming = admission.try_acquire_streaming() + assert held_streaming is not None + hold_task: ( + asyncio.Task[ + tuple[DispatchResult, int, dict[str, str], bytes, FakeStreamWriter] + ] + | None + ) = None + try: + _result, streaming_status, _headers, _body, _writer = await asyncio.wait_for( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/events"), + timeout=RESPONSE_TIMEOUT_S, + ) + assert streaming_status == 503 + + hold_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/hold") + ) + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["active"]["total"] == 1, + ) + _result, status, _headers, _body, _writer = await asyncio.wait_for( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok"), + timeout=RESPONSE_TIMEOUT_S, + ) + assert status == 503 + + snapshot = admission.snapshot() + assert snapshot["queued"]["total"] == 0 + assert snapshot["longest_wait_ms"] == 0 + assert snapshot["rejected"]["queue_timeout"] == 1 + assert snapshot["rejected"]["total"] == 1 + assert snapshot["rejected"]["streaming"] == 1 + + release_hold.set() + await asyncio.wait_for(hold_task, timeout=RESPONSE_TIMEOUT_S) + snapshot = await _wait_for_snapshot( + admission, + lambda item: item["active"]["total"] == 0, + ) + assert snapshot["active"]["total"] == 0 + + _result, later_status, _headers, later_body, _writer = await asyncio.wait_for( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok"), + timeout=RESPONSE_TIMEOUT_S, + ) + assert later_status == 200 + assert later_body == OK_BODY + finally: + held_streaming.release() + release_hold.set() + if hold_task is not None: + with contextlib.suppress(Exception): + await hold_task + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_cancelled_waiter_warns_without_timeout_count_and_quick_grant_is_quiet( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + release_hold = threading.Event() + _register_basic_endpoints(app, release_hold) + admission = _admission(queue_timeout_seconds=10.0) + sent: list[bytes] = [] + + async def send(data: bytes, *, urgent: bool = False) -> None: + sent.append(data) + + async def handler(reader: asyncio.StreamReader, writer: Any) -> None: + await dispatch_stream( + app, + pl_identity(FINGERPRINT), + reader, + writer, + asyncio.get_running_loop(), + admission, + ) + + monkeypatch.setattr(admission_module, "SECURE_LISTENER_QUEUE_WARN_SECONDS", 0.0) + mux = Multiplexer(send, handler, is_listener=True) + try: + await _open_mux_get(mux, 1, "/_queue_bound/hold") + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["active"]["total"] == 1, + ) + await _open_mux_get(mux, 3, "/_queue_bound/ok") + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["queued"]["total"] == 1, + ) + + with caplog.at_level( + logging.WARNING, + logger="convey.secure_listener.admission", + ): + await mux.feed(build_reset(3, RESET_CANCEL).encode()) + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["queued"]["total"] == 0, + ) + + warnings = [ + record + for record in caplog.records + if record.levelno == logging.WARNING + and "Secure listener admission waiter departed" in record.message + ] + assert len(warnings) == 1 + message = warnings[0].message + assert "departure_reason=cancelled" in message + assert "waiter_age_s=" in message + assert "active_total=1" in message + assert "queue_depth=0" in message + assert "queue_timeout_seconds=10.000" in message + _assert_content_free(message) + assert admission.snapshot()["rejected"]["queue_timeout"] == 0 + finally: + release_hold.set() + await mux.close() + await _shutdown_admission(admission) + + caplog.clear() + monkeypatch.setattr(admission_module, "SECURE_LISTENER_QUEUE_WARN_SECONDS", 60.0) + admission = _admission(queue_timeout_seconds=10.0) + held = await admission.acquire() + try: + queued_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok") + ) + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["queued"]["total"] == 1, + ) + with caplog.at_level( + logging.WARNING, + logger="convey.secure_listener.admission", + ): + held.release() + _result, status, _headers, _body, _writer = await asyncio.wait_for( + queued_task, + timeout=RESPONSE_TIMEOUT_S, + ) + assert status == 200 + assert not [ + record + for record in caplog.records + if "Secure listener admission waiter departed" in record.message + ] + finally: + if not held._released: + held.release() + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_deadline_boundary_refuses_short_and_serves_before_generous_deadline( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + release_short = threading.Event() + _register_basic_endpoints(app, release_short) + admission = _admission(queue_timeout_seconds=0.05) + hold_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/hold") + ) + try: + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["active"]["total"] == 1, + ) + _result, status, _headers, _body, _writer = await asyncio.wait_for( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok"), + timeout=RESPONSE_TIMEOUT_S, + ) + assert status == 503 + assert admission.snapshot()["rejected"]["queue_timeout"] == 1 + finally: + release_short.set() + with contextlib.suppress(Exception): + await hold_task + await _shutdown_admission(admission) + + admission = _admission(queue_timeout_seconds=10.0) + held = await admission.acquire() + try: + queued_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok") + ) + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["queued"]["total"] == 1, + ) + held.release() + _result, status, _headers, body, _writer = await asyncio.wait_for( + queued_task, + timeout=RESPONSE_TIMEOUT_S, + ) + assert status == 200 + assert body == OK_BODY + assert admission.snapshot()["rejected"]["queue_timeout"] == 0 + finally: + if not held._released: + held.release() + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_timeout_reclaims_raced_permit_after_delivery_intercept( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + invoked = {"count": 0} + _register_basic_endpoints(app, invoked=invoked) + admission = _admission(queue_timeout_seconds=0.05) + first_permit = await admission.acquire() + captured: dict[str, Any] = {} + captured_event = asyncio.Event() + real_deliver = admission._deliver_waiter + + def capture(waiter: Any, permit: Any) -> None: + captured["waiter"] = waiter + captured["permit"] = permit + captured["active_total"] = admission.snapshot()["active"]["total"] + captured["waiter_has_permit"] = waiter.permit is not None + captured_event.set() + + admission._deliver_waiter = capture + try: + queued_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok") + ) + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["queued"]["total"] == 1, + ) + first_permit.release() + await asyncio.wait_for(captured_event.wait(), timeout=STATE_TIMEOUT_S) + + assert captured["waiter_has_permit"] is True + assert captured["active_total"] == admission.config.capacity + + _result, status, _headers, _body, _writer = await asyncio.wait_for( + queued_task, + timeout=RESPONSE_TIMEOUT_S, + ) + assert status == 503 + assert invoked["count"] == 0 + + real_deliver(captured["waiter"], captured["permit"]) + snapshot = await _wait_for_snapshot( + admission, + lambda item: item["active"]["total"] == 0, + ) + assert snapshot["active"]["total"] == 0 + + _result, later_status, _headers, later_body, _writer = await asyncio.wait_for( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok"), + timeout=RESPONSE_TIMEOUT_S, + ) + assert later_status == 200 + assert later_body == OK_BODY + finally: + admission._deliver_waiter = real_deliver + if not first_permit._released: + first_permit.release() + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_raced_permit_release_guard_falsification_pins_active_total( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + _register_basic_endpoints(app) + admission = _admission(queue_timeout_seconds=0.05) + first_permit = await admission.acquire() + captured: dict[str, Any] = {} + captured_event = asyncio.Event() + real_deliver = admission._deliver_waiter + + def capture(waiter: Any, permit: Any) -> None: + captured["waiter"] = waiter + captured["permit"] = permit + captured["active_total"] = admission.snapshot()["active"]["total"] + captured["waiter_has_permit"] = waiter.permit is not None + captured_event.set() + + admission._deliver_waiter = capture + try: + queued_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok") + ) + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["queued"]["total"] == 1, + ) + first_permit.release() + await asyncio.wait_for(captured_event.wait(), timeout=STATE_TIMEOUT_S) + + assert captured["waiter_has_permit"] is True + assert captured["active_total"] == admission.config.capacity + captured["permit"].release = lambda: None + + _result, status, _headers, _body, _writer = await asyncio.wait_for( + queued_task, + timeout=RESPONSE_TIMEOUT_S, + ) + assert status == 503 + real_deliver(captured["waiter"], captured["permit"]) + assert admission.snapshot()["active"]["total"] == admission.config.capacity + finally: + admission._deliver_waiter = real_deliver + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_concurrent_queue_timeouts_account_and_warn_individually( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + _register_basic_endpoints(app) + monkeypatch.setattr(admission_module, "SECURE_LISTENER_QUEUE_WARN_SECONDS", 0.0) + admission = _admission(queue_timeout_seconds=0.05) + held = await admission.acquire() + count = 3 + try: + with caplog.at_level( + logging.WARNING, + logger="convey.secure_listener.admission", + ): + tasks = [ + asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok") + ) + for _index in range(count) + ] + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["queued"]["total"] == count, + ) + results = await asyncio.wait_for( + asyncio.gather(*tasks), + timeout=RESPONSE_TIMEOUT_S, + ) + + assert [status for _result, status, *_rest in results] == [503] * count + snapshot = admission.snapshot() + assert snapshot["rejected"]["queue_timeout"] == count + assert snapshot["queued"]["total"] == 0 + assert snapshot["longest_wait_ms"] == 0 + warnings = [ + record + for record in caplog.records + if record.levelno == logging.WARNING + and "Secure listener admission waiter departed" in record.message + ] + assert len(warnings) == count + finally: + held.release() + snapshot = await _wait_for_snapshot( + admission, + lambda item: item["active"]["total"] == 0, + ) + assert snapshot["active"]["total"] == 0 + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_zero_recv_credit_timeout_does_not_block_later_stream( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + _register_basic_endpoints(app) + sent: list[bytes] = [] + admission = _admission(queue_timeout_seconds=0.05) + held = await admission.acquire() + + async def send(data: bytes, *, urgent: bool = False) -> None: + sent.append(data) + + async def handler(reader: asyncio.StreamReader, writer: Any) -> None: + await dispatch_stream( + app, + pl_identity(FINGERPRINT), + reader, + writer, + asyncio.get_running_loop(), + admission, + ) + + mux = Multiplexer(send, handler, is_listener=True) + try: + await _open_mux_request( + mux, + 1, + "POST", + "/_queue_bound/upload", + headers={"content-type": "application/octet-stream"}, + content_length=INITIAL_WINDOW + 100, + ) + state = await _wait_for_mux_state(mux, 1) + await mux.feed(build_data(1, b"x" * state.recv_credit).encode()) + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["rejected"]["queue_timeout"] == 1, + ) + + held.release() + await _open_mux_get(mux, 3, "/_queue_bound/ok") + payload = await _wait_for_stream_payload(sent, 3) + + assert b"HTTP/1.1 200" in payload + assert OK_BODY in payload + finally: + if not held._released: + held.release() + await mux.close() + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_zero_queue_timeout_disables_refusal_but_keeps_slow_warning( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, +) -> None: + app, _journal = make_convey_app( + tmp_path, + monkeypatch, + link={"posture": "spl", "secure_listener_queue_timeout_seconds": 0}, + ) + _authorize_fingerprint(monkeypatch) + release_hold = threading.Event() + _register_basic_endpoints(app, release_hold) + monkeypatch.setattr(admission_module, "SECURE_LISTENER_QUEUE_WARN_SECONDS", 0.0) + admission = _admission(queue_timeout_seconds=0.0) + hold_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/hold") + ) + try: + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["active"]["total"] == 1, + ) + queued_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok") + ) + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["queued"]["total"] == 1, + ) + with caplog.at_level( + logging.WARNING, + logger="convey.secure_listener.admission", + ): + release_hold.set() + _result, status, _headers, body, _writer = await asyncio.wait_for( + queued_task, + timeout=RESPONSE_TIMEOUT_S, + ) + + assert status == 200 + assert body == OK_BODY + assert admission.snapshot()["rejected"]["queue_timeout"] == 0 + warnings = [ + record + for record in caplog.records + if "Secure listener admission waiter departed" in record.message + ] + assert len(warnings) == 1 + assert "departure_reason=granted" in warnings[0].message + finally: + release_hold.set() + with contextlib.suppress(Exception): + await hold_task + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_refuse_when_full_false_still_allows_deadline_timeout( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + release_hold = threading.Event() + _register_basic_endpoints(app, release_hold) + admission = _admission(queue_timeout_seconds=0.05, refuse_when_full=False) + assert admission.config.refuse_when_full is False + hold_task = asyncio.create_task( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/hold") + ) + try: + await _wait_for_snapshot( + admission, + lambda snapshot: snapshot["active"]["total"] == 1, + ) + _result, status, _headers, body, _writer = await asyncio.wait_for( + _dispatch_raw_request(app, admission, "GET", "/_queue_bound/ok"), + timeout=RESPONSE_TIMEOUT_S, + ) + + assert status == 503 + assert body == QUEUE_TIMEOUT_BODY + assert admission.snapshot()["rejected"]["queue_timeout"] == 1 + finally: + release_hold.set() + with contextlib.suppress(Exception): + await hold_task + await _shutdown_admission(admission) + + +@pytest.mark.asyncio +async def test_paired_mtls_client_parses_queue_timeout_and_cancel_clears_queue( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + forward_root = tmp_path / "forward" + forward_root.mkdir() + harness = await SecureListenerHarness.start( + forward_root, + monkeypatch, + admission_config=SecureListenerAdmissionConfig( + capacity=1, + streaming_capacity=1, + refuse_when_full=False, + queue_timeout_seconds=0.05, + ), + ) + session = None + held = None + try: + session = await pair_and_open_session( + harness, + nonce="10000000000000000000000000000041", + label="queue-timeout-phone", + ) + held = await harness.admission.acquire() + status, headers, body = await asyncio.wait_for( + session.request("GET", "/app/network/api/status"), + timeout=RESPONSE_TIMEOUT_S, + ) + + assert status == 503 + assert headers["retry-after"] == str( + wsgi_module.SECURE_LISTENER_REFUSAL_RETRY_AFTER_SECONDS + ) + assert body == QUEUE_TIMEOUT_BODY + finally: + if held is not None: + held.release() + if session is not None: + with contextlib.suppress(Exception): + await session.close() + await harness.close() + + inverse_root = tmp_path / "inverse" + inverse_root.mkdir() + harness = await SecureListenerHarness.start( + inverse_root, + monkeypatch, + admission_config=SecureListenerAdmissionConfig( + capacity=1, + streaming_capacity=1, + refuse_when_full=False, + queue_timeout_seconds=10.0, + ), + ) + session = None + held = None + reset_payloads: list[dict[str, Any]] = [] + reset_cancel_seen = asyncio.Event() + real_dispatch = mux_module.Multiplexer._dispatch + + async def observed_dispatch(self: Multiplexer, frame: Frame) -> None: + if frame.flags & FLAG_RESET: + with contextlib.suppress(ValueError): + if parse_reset_reason(frame) == RESET_CANCEL: + reset_cancel_seen.set() + await real_dispatch(self, frame) + + def capture_reset(event: str, fields: dict[str, Any]) -> None: + if event == "stream_reset": + reset_payloads.append(fields) + + monkeypatch.setattr(mux_module.Multiplexer, "_dispatch", observed_dispatch) + harness.listener._emit = capture_reset + try: + session = await pair_and_open_session( + harness, + nonce="10000000000000000000000000000042", + label="queue-cancel-phone", + ) + held = await harness.admission.acquire() + request_task = asyncio.create_task( + session.request("GET", "/app/network/api/status") + ) + await _wait_for_snapshot( + harness.admission, + lambda snapshot: snapshot["queued"]["total"] == 1, + ) + request_task.cancel() + with contextlib.suppress(asyncio.CancelledError): + await request_task + await asyncio.wait_for(reset_cancel_seen.wait(), timeout=STATE_TIMEOUT_S) + snapshot = await _wait_for_snapshot( + harness.admission, + lambda item: item["queued"]["total"] == 0, + ) + + assert snapshot["queued"]["total"] == 0 + assert snapshot["rejected"]["queue_timeout"] == 0 + assert not any( + payload.get("context") == RESET_CTX_HANDLER_EXCEPTION + for payload in reset_payloads + ) + finally: + if held is not None: + held.release() + if session is not None: + with contextlib.suppress(Exception): + await session.close() + await harness.close() + + +class _ClosedWriter(FakeStreamWriter): + def __init__(self) -> None: + super().__init__() + self.begin_drain_called = False + + async def write(self, data: bytes) -> None: + raise ConnectionError("stream closed") + + def begin_drain(self, context: str) -> None: + self.begin_drain_called = True + super().begin_drain(context) + + +@pytest.mark.asyncio +async def test_queue_timeout_dead_writer_returns_503_without_handler_exception( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + app, _journal = make_convey_app(tmp_path, monkeypatch, link={"posture": "spl"}) + _authorize_fingerprint(monkeypatch) + _register_basic_endpoints(app) + admission = _admission(queue_timeout_seconds=0.05) + held = await admission.acquire() + writer = _ClosedWriter() + reader = asyncio.StreamReader() + reader.feed_data( + _http_head_bytes("GET", "/_queue_bound/ok", headers=None, content_length=0) + ) + reader.feed_eof() + try: + result = await asyncio.wait_for( + dispatch_stream( + app, + pl_identity(FINGERPRINT), + reader, + writer, + asyncio.get_running_loop(), + admission, + ), + timeout=RESPONSE_TIMEOUT_S, + ) + + assert result.status == 503 + assert writer.begin_drain_called is True + assert writer.drain_context == RESET_CTX_BODY_DISCARD_CANCELLATION + finally: + held.release() + await _shutdown_admission(admission) + + +def test_secure_listener_queue_bounds_stay_below_link_request_timeout() -> None: + assert ( + SECURE_LISTENER_QUEUE_WARN_SECONDS + < DEFAULT_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS + < _REQUEST_TIMEOUT_SECONDS + == 180 + ) diff --git a/tests/test_secure_listener_runtime.py b/tests/test_secure_listener_runtime.py index d5c23a4f1..23024a127 100644 --- a/tests/test_secure_listener_runtime.py +++ b/tests/test_secure_listener_runtime.py @@ -24,6 +24,7 @@ from solstone.convey.secure_listener.accept import ( ) from solstone.convey.secure_listener.admission import ( DEFAULT_SECURE_LISTENER_CAPACITY, + DEFAULT_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS, DEFAULT_SECURE_LISTENER_STREAMING_CAPACITY, SecureListenerAdmission, SecureListenerAdmissionConfig, @@ -134,6 +135,7 @@ def test_stop_all_after_loop_closed_does_not_raise(): def test_secure_listener_admission_config_defaults_to_current_capacity( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, ) -> None: _runtime_journal( tmp_path, @@ -141,12 +143,18 @@ def test_secure_listener_admission_config_defaults_to_current_capacity( {"setup": {"completed_at": 1700000000000}}, ) - config = resolve_admission_config() + with caplog.at_level( + logging.INFO, + logger="convey.secure_listener.admission", + ): + config = resolve_admission_config() assert config.capacity == DEFAULT_SECURE_LISTENER_CAPACITY == 16 assert config.streaming_capacity == DEFAULT_SECURE_LISTENER_STREAMING_CAPACITY == 8 assert config.refuse_when_full is False + assert config.queue_timeout_seconds == DEFAULT_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS assert config.queue_limit == 32 + assert not caplog.records def test_secure_listener_admission_config_reads_link_namespace( @@ -162,6 +170,7 @@ def test_secure_listener_admission_config_reads_link_namespace( "secure_listener_capacity": 24, "secure_listener_streaming_capacity": 6, "secure_listener_refuse_when_full": True, + "secure_listener_queue_timeout_seconds": 30.5, }, }, ) @@ -171,9 +180,112 @@ def test_secure_listener_admission_config_reads_link_namespace( assert config.capacity == 24 assert config.streaming_capacity == 6 assert config.refuse_when_full is True + assert config.queue_timeout_seconds == 30.5 assert config.queue_limit == 48 +@pytest.mark.parametrize( + ("raw", "expected"), + [ + (120, 120.0), + (1.0, 1.0), + (600, 600.0), + ], +) +def test_secure_listener_admission_config_queue_timeout_accepts_numbers_silently( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, + raw: object, + expected: float, +) -> None: + _runtime_journal( + tmp_path, + monkeypatch, + { + "setup": {"completed_at": 1700000000000}, + "link": {"secure_listener_queue_timeout_seconds": raw}, + }, + ) + + with caplog.at_level( + logging.INFO, + logger="convey.secure_listener.admission", + ): + config = resolve_admission_config() + + assert config.queue_timeout_seconds == expected + assert not caplog.records + + +def test_secure_listener_admission_config_queue_timeout_zero_disables_at_info( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, +) -> None: + _runtime_journal( + tmp_path, + monkeypatch, + { + "setup": {"completed_at": 1700000000000}, + "link": {"secure_listener_queue_timeout_seconds": 0}, + }, + ) + + with caplog.at_level( + logging.INFO, + logger="convey.secure_listener.admission", + ): + config = resolve_admission_config() + + assert config.queue_timeout_seconds == 0.0 + assert [record.levelno for record in caplog.records] == [logging.INFO] + assert ( + "link.secure_listener_queue_timeout_seconds is 0; " + "secure listener queue timeout disabled" + ) in caplog.text + + +@pytest.mark.parametrize( + "raw", + [ + 601, + False, + True, + -1, + 0.5, + "abc", + ], +) +def test_secure_listener_admission_config_queue_timeout_warns_and_falls_back( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, + raw: object, +) -> None: + _runtime_journal( + tmp_path, + monkeypatch, + { + "setup": {"completed_at": 1700000000000}, + "link": {"secure_listener_queue_timeout_seconds": raw}, + }, + ) + + with caplog.at_level( + logging.WARNING, + logger="convey.secure_listener.admission", + ): + config = resolve_admission_config() + + assert config.queue_timeout_seconds == DEFAULT_SECURE_LISTENER_QUEUE_TIMEOUT_SECONDS + assert [record.levelno for record in caplog.records] == [logging.WARNING] + assert ( + "Invalid link.secure_listener_queue_timeout_seconds in journal config: " + f"{raw!r} \u2014 defaulting to 120.0" + ) in caplog.text + + def test_secure_listener_admission_config_warns_and_falls_back( tmp_path: Path, monkeypatch: pytest.MonkeyPatch,