From 8c6de47ea6b86ceb1f5d18596438855c11bbe16d Mon Sep 17 00:00:00 2001 From: Andre van Zuydam Date: Wed, 30 Sep 2026 04:43:42 +0200 Subject: [PATCH 1/8] refactor(queue): revive a persisted job through Job.from_data() MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The file backends store each job as a JSON file and rebuild the Job from that dict on pop. The seven-field mapping (id, topic-or-default, data, priority, attempts, error) was spelt out at three pop sites — FileQueue's claim-by-id, LiteBackend.pop() and LiteBackend.pop_batch(). It now lives once as Job.from_data(queue, job_data, default_topic). Behaviour-preserving: same fields, same defaults, same topic fallback. The reclaim path in LiteBackend keeps its own construction — it overrides attempts and error and defaults data, so it is not the same mapping. Co-Authored-By: Claude Opus 4.8 Co-Authored-By: Tina4 <82961293+tina4stack@users.noreply.github.com> --- tina4_python/queue/__init__.py | 9 +-------- tina4_python/queue/job.py | 20 ++++++++++++++++++++ tina4_python/queue/lite_backend.py | 20 ++------------------ 3 files changed, 23 insertions(+), 26 deletions(-) diff --git a/tina4_python/queue/__init__.py b/tina4_python/queue/__init__.py index f410d5e2..beb4553b 100644 --- a/tina4_python/queue/__init__.py +++ b/tina4_python/queue/__init__.py @@ -431,14 +431,7 @@ def pop_by_id(self, id: str) -> Job | None: # claim the pending file — mirrors LiteBackend.pop(). self._backend._write_reserved(job_data) os.unlink(filepath) - return Job( - queue=self, job_id=job_data["id"], - topic=job_data.get("topic", self.topic), - data=job_data["data"], - priority=job_data.get("priority", 0), - attempts=job_data.get("attempts", 0), - error=job_data.get("error"), - ) + return Job.from_data(self, job_data, self.topic) except (json.JSONDecodeError, FileNotFoundError): continue except FileNotFoundError: diff --git a/tina4_python/queue/job.py b/tina4_python/queue/job.py index fbf0658c..018dbd80 100644 --- a/tina4_python/queue/job.py +++ b/tina4_python/queue/job.py @@ -20,6 +20,26 @@ def __init__(self, queue, job_id, topic: str, data: dict, # without trawling logs. self.error: str | None = error + @classmethod + def from_data(cls, queue, job_data: dict, default_topic: str) -> "Job": + """Build a Job from a persisted job_data dict. + + The file backends (LiteBackend and the FileQueue) store each job as a + JSON file and rebuild the Job from that dict on pop. The field mapping + is identical everywhere a stored job is revived — id, topic (falling + back to the queue's own topic), data, priority, attempts and error — + so it lives here once rather than being spelt out at each pop site. + """ + return cls( + queue=queue, + job_id=job_data["id"], + topic=job_data.get("topic", default_topic), + data=job_data["data"], + priority=job_data.get("priority", 0), + attempts=job_data.get("attempts", 0), + error=job_data.get("error"), + ) + @property def data(self): """Alias for payload — deprecated, use .payload instead.""" diff --git a/tina4_python/queue/lite_backend.py b/tina4_python/queue/lite_backend.py index e54f7ede..3d573a71 100644 --- a/tina4_python/queue/lite_backend.py +++ b/tina4_python/queue/lite_backend.py @@ -245,15 +245,7 @@ def pop(self, queue_ref) -> Job | None: except FileNotFoundError: continue # Already consumed by another worker - return Job( - queue=queue_ref, - job_id=job_data["id"], - topic=job_data.get("topic", self._topic), - data=job_data["data"], - priority=job_data.get("priority", 0), - attempts=job_data.get("attempts", 0), - error=job_data.get("error"), - ) + return Job.from_data(queue_ref, job_data, self._topic) return None @@ -278,15 +270,7 @@ def pop_batch(self, count: int, queue_ref) -> list: except FileNotFoundError: continue - results.append(Job( - queue=queue_ref, - job_id=job_data["id"], - topic=job_data.get("topic", self._topic), - data=job_data["data"], - priority=job_data.get("priority", 0), - attempts=job_data.get("attempts", 0), - error=job_data.get("error"), - )) + results.append(Job.from_data(queue_ref, job_data, self._topic)) return results From 1f2e73a75b33909c8037fcfb19924b1b10bee6a6 Mon Sep 17 00:00:00 2001 From: Andre van Zuydam Date: Wed, 30 Sep 2026 04:44:33 +0200 Subject: [PATCH 2/8] refactor(rate-limit): enforce once in RateLimiter.apply() MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit RateLimiter.apply() already holds the enforcement path — resolve the IP, check it, set the x-ratelimit headers, and on refusal add retry-after and the 429. That same body was copied into two staticmethod entry points: RateLimiter.before_rate_limit (shared instance) and RateLimiterMiddleware.before_rate_limit (the middleware's own instance). Both now just resolve their limiter and delegate to apply(). The response shape, header set and 429 wording are unchanged — the enforcement lives in one place. Co-Authored-By: Claude Opus 4.8 Co-Authored-By: Tina4 <82961293+tina4stack@users.noreply.github.com> --- tina4_python/core/middleware.py | 13 +------------ tina4_python/core/rate_limiter.py | 13 +------------ 2 files changed, 2 insertions(+), 24 deletions(-) diff --git a/tina4_python/core/middleware.py b/tina4_python/core/middleware.py index ed820679..0b9bf173 100644 --- a/tina4_python/core/middleware.py +++ b/tina4_python/core/middleware.py @@ -528,18 +528,7 @@ def _get_limiter(cls): @staticmethod def before_rate_limit(request, response): """Middleware hook — enforces rate limiting before the route handler.""" - limiter = RateLimiterMiddleware._get_limiter() - ip = getattr(request, "ip", None) or "unknown" - allowed, info = limiter.check(ip) - limiter.apply_headers(response, info) - if not allowed: - retry_after = max(1, int(info.get("reset", limiter.window))) - response.header("retry-after", str(retry_after)) - if hasattr(response, "error"): - response.error("Too Many Requests", f"Rate limit exceeded. Retry in {retry_after}s.", 429) - else: - setattr(response, "status_code", 429) - return request, response + return RateLimiterMiddleware._get_limiter().apply(request, response) @staticmethod def check(ip: str): diff --git a/tina4_python/core/rate_limiter.py b/tina4_python/core/rate_limiter.py index e39661aa..c4489487 100644 --- a/tina4_python/core/rate_limiter.py +++ b/tina4_python/core/rate_limiter.py @@ -37,18 +37,7 @@ def _shared(cls) -> "RateLimiter": @staticmethod def before_rate_limit(request, response): """Class-based middleware entry point — enforces the shared rate limit.""" - limiter = RateLimiter._shared() - ip = getattr(request, "ip", None) or "unknown" - allowed, info = limiter.check(ip) - limiter.apply_headers(response, info) - if not allowed: - retry_after = max(1, int(info.get("reset", limiter.window))) - response.header("retry-after", str(retry_after)) - if hasattr(response, "error"): - response.error("Too Many Requests", f"Rate limit exceeded. Retry in {retry_after}s.", 429) - else: - setattr(response, "status_code", 429) - return request, response + return RateLimiter._shared().apply(request, response) def __init__(self): self._env_snapshot: tuple = () From 0c047a01235732d52cc0797b76de75fedd6646e3 Mon Sep 17 00:00:00 2001 From: Andre van Zuydam Date: Wed, 30 Sep 2026 04:45:51 +0200 Subject: [PATCH 3/8] refactor(queue): share dead_letters()/retry_job() across topic backends MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The Kafka and RabbitMQ backends both keep dead jobs on a .dead_letter topic and adapt an external connector to the same dequeue/enqueue interface. Their dead_letters() and retry_job() were byte-identical drain-filter-requeue loops — pure over that interface, carrying no broker specifics. They now live once in TopicDeadLetterMixin, which both backends inherit. The broker-specific methods (push, pop, fail, retry, reject, purge, clear) stay in each backend unchanged — those differ genuinely (Kafka re-produces past a committed offset; RabbitMQ tracks jobs and re-publishes with an attempt count) and are not touched. Co-Authored-By: Claude Opus 4.8 Co-Authored-By: Tina4 <82961293+tina4stack@users.noreply.github.com> --- tina4_python/queue/kafka_backend.py | 51 ++---------------- tina4_python/queue/rabbitmq_backend.py | 53 ++---------------- tina4_python/queue/topic_dead_letter.py | 71 +++++++++++++++++++++++++ 3 files changed, 79 insertions(+), 96 deletions(-) create mode 100644 tina4_python/queue/topic_dead_letter.py diff --git a/tina4_python/queue/kafka_backend.py b/tina4_python/queue/kafka_backend.py index f94eac70..09e1430a 100644 --- a/tina4_python/queue/kafka_backend.py +++ b/tina4_python/queue/kafka_backend.py @@ -12,9 +12,10 @@ import os from tina4_python.queue.job import Job +from tina4_python.queue.topic_dead_letter import TopicDeadLetterMixin -class KafkaBackend: +class KafkaBackend(TopicDeadLetterMixin): """Backend adapter wrapping KafkaBackend for the unified Queue API.""" def __init__(self, topic: str, max_retries: int): @@ -132,52 +133,8 @@ def failed(self, max_retries: int = None) -> list[dict]: "mongodb backend to enumerate retryable failures." ) - def dead_letters(self, max_retries: int = None) -> list[dict]: - """Consume dead_letter topic, republish, return jobs at/over max_retries. - - Accepts max_retries to match the LiteBackend contract — Queue.dead_letters() - passes it as a kwarg, so without this signature the call raised TypeError. - """ - mr = max_retries if max_retries is not None else self._max_retries - dl_topic = f"{self._topic}.dead_letter" - results = [] - requeue = [] - while True: - msg = self._backend.dequeue(dl_topic) - if msg is None: - break - payload = msg.get("payload", msg) - attempts = msg.get("attempts", 0) - if attempts >= mr: - results.append({"id": msg.get("id"), "data": payload, - "attempts": attempts, "error": msg.get("error")}) - requeue.append(msg) - for msg in requeue: - self._backend.enqueue(dl_topic, msg) - return results - - def retry_job(self, job_id: str, delay_seconds: int = 0) -> bool: - """Move job from dead_letter topic back to main topic.""" - dl_topic = f"{self._topic}.dead_letter" - found = None - requeue = [] - while True: - msg = self._backend.dequeue(dl_topic) - if msg is None: - break - if msg.get("id") == job_id and found is None: - found = msg - else: - requeue.append(msg) - for msg in requeue: - self._backend.enqueue(dl_topic, msg) - if found is None: - return False - found["attempts"] = found.get("attempts", 0) + 1 - found["status"] = "pending" - found.pop("error", None) - self._backend.enqueue(self._topic, found) - return True + # dead_letters() and retry_job() come from TopicDeadLetterMixin — both are + # pure over the .dead_letter topic and identical to the rabbitmq backend. def clear(self) -> int: """Not performable on Kafka — raises naming the backend and the operation. diff --git a/tina4_python/queue/rabbitmq_backend.py b/tina4_python/queue/rabbitmq_backend.py index 59158c7f..b9002bae 100644 --- a/tina4_python/queue/rabbitmq_backend.py +++ b/tina4_python/queue/rabbitmq_backend.py @@ -13,11 +13,10 @@ from tina4_python.queue.amqp_url import parse_amqp_url as _parse_amqp_url from tina4_python.queue.job import Job +from tina4_python.queue.topic_dead_letter import TopicDeadLetterMixin - - -class RabbitMQBackend: +class RabbitMQBackend(TopicDeadLetterMixin): """Backend adapter wrapping RabbitMQBackend for the unified Queue API.""" def __init__(self, topic: str, max_retries: int): @@ -151,52 +150,8 @@ def failed(self, max_retries: int = None) -> list[dict]: "file or mongodb backend to enumerate retryable failures." ) - def dead_letters(self, max_retries: int = None) -> list[dict]: - """Drain the dead_letter queue, re-enqueue, and return jobs at/over max_retries. - - Accepts max_retries to match the LiteBackend contract — Queue.dead_letters() - passes it as a kwarg, so without this signature the call raised TypeError. - """ - mr = max_retries if max_retries is not None else self._max_retries - dl_topic = f"{self._topic}.dead_letter" - results = [] - requeue = [] - while True: - msg = self._backend.dequeue(dl_topic) - if msg is None: - break - payload = msg.get("payload", msg) - attempts = msg.get("attempts", 0) - if attempts >= mr: - results.append({"id": msg.get("id"), "data": payload, - "attempts": attempts, "error": msg.get("error")}) - requeue.append(msg) - for msg in requeue: - self._backend.enqueue(dl_topic, msg) - return results - - def retry_job(self, job_id: str, delay_seconds: int = 0) -> bool: - """Move a job from the dead_letter queue back to the main topic.""" - dl_topic = f"{self._topic}.dead_letter" - found = None - requeue = [] - while True: - msg = self._backend.dequeue(dl_topic) - if msg is None: - break - if msg.get("id") == job_id and found is None: - found = msg - else: - requeue.append(msg) - for msg in requeue: - self._backend.enqueue(dl_topic, msg) - if found is None: - return False - found["attempts"] = found.get("attempts", 0) + 1 - found["status"] = "pending" - found.pop("error", None) - self._backend.enqueue(self._topic, found) - return True + # dead_letters() and retry_job() come from TopicDeadLetterMixin — both are + # pure over the .dead_letter topic and identical to the kafka backend. def clear(self) -> int: """Not performable on RabbitMQ — raises naming the backend and the operation. diff --git a/tina4_python/queue/topic_dead_letter.py b/tina4_python/queue/topic_dead_letter.py new file mode 100644 index 00000000..f8d36a3a --- /dev/null +++ b/tina4_python/queue/topic_dead_letter.py @@ -0,0 +1,71 @@ +# Copyright (c) 2026 Code Infinity +# SPDX-License-Identifier: MPL-2.0 +# This Source Code Form is subject to the terms of the Mozilla Public +# License, v. 2.0. If a copy of the MPL was not distributed with this +# file, You can obtain one at https://mozilla.org/MPL/2.0/. + +# Tina4 Queue — dead-letter helpers shared by the topic-log backends. +""" +Dead-letter enumeration and single-job retry, shared by the broker backends +that keep dead jobs on a ``.dead_letter`` topic (Kafka and RabbitMQ). + +Both backends adapt an external connector to the same ``dequeue``/``enqueue`` +interface, so these two operations are pure over that interface — they carry +no broker specifics at all. The broker differences (push, pop, fail, retry, +reject, purge, clear) stay in each backend; only these two identical +drain-filter-requeue loops live here. + +A host class must expose ``self._backend`` (a connector with ``dequeue`` and +``enqueue``), ``self._topic`` (str), and ``self._max_retries`` (int). +""" + + +class TopicDeadLetterMixin: + """dead_letters() and retry_job() for a ``.dead_letter`` backend.""" + + def dead_letters(self, max_retries: int = None) -> list[dict]: + """Drain the dead_letter topic, re-enqueue it, and return jobs at/over max_retries. + + Accepts max_retries to match the LiteBackend contract — Queue.dead_letters() + passes it as a kwarg, so without this signature the call raised TypeError. + """ + mr = max_retries if max_retries is not None else self._max_retries + dl_topic = f"{self._topic}.dead_letter" + results = [] + requeue = [] + while True: + msg = self._backend.dequeue(dl_topic) + if msg is None: + break + payload = msg.get("payload", msg) + attempts = msg.get("attempts", 0) + if attempts >= mr: + results.append({"id": msg.get("id"), "data": payload, + "attempts": attempts, "error": msg.get("error")}) + requeue.append(msg) + for msg in requeue: + self._backend.enqueue(dl_topic, msg) + return results + + def retry_job(self, job_id: str, delay_seconds: int = 0) -> bool: + """Move a job from the dead_letter topic back to the main topic.""" + dl_topic = f"{self._topic}.dead_letter" + found = None + requeue = [] + while True: + msg = self._backend.dequeue(dl_topic) + if msg is None: + break + if msg.get("id") == job_id and found is None: + found = msg + else: + requeue.append(msg) + for msg in requeue: + self._backend.enqueue(dl_topic, msg) + if found is None: + return False + found["attempts"] = found.get("attempts", 0) + 1 + found["status"] = "pending" + found.pop("error", None) + self._backend.enqueue(self._topic, found) + return True From c7390d3d301db69d906c3dc5ff62aabe3a7b1348 Mon Sep 17 00:00:00 2001 From: Andre van Zuydam Date: Wed, 30 Sep 2026 04:46:40 +0200 Subject: [PATCH 4/8] refactor(orm): share the uncapped has_many paging loop The lazy relationship descriptor and the imperative has_many() both page through ALL children in _LAZY_PAGE_SIZE blocks so a parent with more than one page never loses the tail (REL-EAGER-UNBOUNDED). The page loop was written out in both places; it now lives once as fields.fetch_all_pages(db, sql, params, start_offset). Behaviour-preserving: same page size, same stop-on-short-block, same start offset. model.py already imported from fields, so the helper adds no new coupling. Both callers still wrap the returned rows in their own model. Co-Authored-By: Claude Opus 4.8 Co-Authored-By: Tina4 <82961293+tina4stack@users.noreply.github.com> --- tina4_python/orm/fields.py | 31 ++++++++++++++++++++++--------- tina4_python/orm/model.py | 12 ++---------- 2 files changed, 24 insertions(+), 19 deletions(-) diff --git a/tina4_python/orm/fields.py b/tina4_python/orm/fields.py index cce71850..565f2f77 100644 --- a/tina4_python/orm/fields.py +++ b/tina4_python/orm/fields.py @@ -622,6 +622,27 @@ def __init__(self, to=None, related_name: str = None, **kwargs): _LAZY_PAGE_SIZE = 1000 +def fetch_all_pages(db, sql: str, params: list, start_offset: int = 0) -> list: + """Return every row of ``sql``, paging in blocks of ``_LAZY_PAGE_SIZE``. + + Uncapped by design (REL-EAGER-UNBOUNDED): a relationship load that passed a + silent ``limit=1000`` truncated the tail with no signal, so this pages until + a short block instead. Shared by the descriptor lazy-load and the imperative + ``has_many`` so both revive the full set the same way. Returns raw rows; the + caller wraps them in the related model. + """ + records = [] + offset = start_offset + while True: + result = db.fetch(sql, params, limit=_LAZY_PAGE_SIZE, offset=offset) + batch = result.records + records.extend(batch) + if len(batch) < _LAZY_PAGE_SIZE: + break + offset += _LAZY_PAGE_SIZE + return records + + class RelationshipDescriptor: """Base descriptor for ORM relationships. Lazy-loads on first access.""" @@ -698,15 +719,7 @@ def _load(self, obj): sql = f"SELECT * FROM {table} WHERE {where} ORDER BY {order_col}" # REL-EAGER-UNBOUNDED: page through ALL children rather than silently # truncating at a fixed cap. - records = [] - offset = 0 - while True: - result = db.fetch(sql, [pk_value], limit=_LAZY_PAGE_SIZE, offset=offset) - batch = result.records - records.extend(batch) - if len(batch) < _LAZY_PAGE_SIZE: - break - offset += _LAZY_PAGE_SIZE + records = fetch_all_pages(db, sql, [pk_value]) return [related_cls(row) for row in records] diff --git a/tina4_python/orm/model.py b/tina4_python/orm/model.py index 53ec7a96..6c64777b 100644 --- a/tina4_python/orm/model.py +++ b/tina4_python/orm/model.py @@ -1506,7 +1506,7 @@ def has_many(self, related_class, foreign_key: str = None, limit: int = None, of row count whether accessed imperatively or lazily. An explicit ``limit`` still pages (explicit, never silent). """ - from tina4_python.orm.fields import _LAZY_PAGE_SIZE + from tina4_python.orm.fields import fetch_all_pages pk = self._get_pk() pk_value = getattr(self, pk) fk = foreign_key or f"{self.__class__.__name__.lower()}_id" @@ -1522,15 +1522,7 @@ def has_many(self, related_class, foreign_key: str = None, limit: int = None, of return [related_class(row) for row in result.records] # No explicit limit -> page through ALL rows (uncapped, parity with lazy). - records = [] - page_offset = offset - while True: - result = db.fetch(sql, [pk_value], limit=_LAZY_PAGE_SIZE, offset=page_offset) - batch = result.records - records.extend(batch) - if len(batch) < _LAZY_PAGE_SIZE: - break - page_offset += _LAZY_PAGE_SIZE + records = fetch_all_pages(db, sql, [pk_value], offset) return [related_class(row) for row in records] def belongs_to(self, related_class, foreign_key: str = None) -> Self | None: From 2ffbccce9514f7773afcb0f7f8c80303e14f1ec7 Mon Sep 17 00:00:00 2001 From: Andre van Zuydam Date: Wed, 30 Sep 2026 04:47:54 +0200 Subject: [PATCH 5/8] refactor(websocket): share room membership across the two connections MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both connection types carry room membership — the raw-socket WebSocketConnection and the ASGI _AsgiWebSocketConnection. join_room, leave_room, broadcast_to_room and the rooms property only touch this connection's own _rooms set and delegate to _manager, so they do not depend on the transport. They now live once in RoomMemberMixin, which both connections inherit; server.py already imports from the websocket package. Transport-specific code (framing, send, ping, close, the ASGI vs stream-reader wiring) stays in each class unchanged. Co-Authored-By: Claude Opus 4.8 Co-Authored-By: Tina4 <82961293+tina4stack@users.noreply.github.com> --- tina4_python/core/server.py | 29 ++----------- tina4_python/websocket/__init__.py | 66 ++++++++++++++++++------------ 2 files changed, 43 insertions(+), 52 deletions(-) diff --git a/tina4_python/core/server.py b/tina4_python/core/server.py index a5abe763..6782f801 100644 --- a/tina4_python/core/server.py +++ b/tina4_python/core/server.py @@ -1181,7 +1181,7 @@ async def dash(req, res): # ── WebSocket support ────────────────────────────────────────── -from tina4_python.websocket import CLOSE_GOING_AWAY, WebSocketConnection, WebSocketManager +from tina4_python.websocket import CLOSE_GOING_AWAY, RoomMemberMixin, WebSocketConnection, WebSocketManager _ws_manager = WebSocketManager() @@ -1309,7 +1309,7 @@ async def _handle_asgi_websocket(scope: dict, receive, send): _ws_manager.remove(conn) -class _AsgiWebSocketConnection: +class _AsgiWebSocketConnection(RoomMemberMixin): """WebSocket connection wrapper for ASGI servers (uvicorn, etc.). Supports both Router's (conn, event, data) style and WebSocketServer's @@ -1344,10 +1344,8 @@ def __init__(self, scope, receive, send, path, params, manager): def closed(self) -> bool: return self._closed - @property - def rooms(self) -> set: - """Return the set of room names this connection has joined.""" - return self._rooms + # rooms / join_room / leave_room / broadcast_to_room come from + # RoomMemberMixin — transport-agnostic, shared with WebSocketConnection. def on_message(self, handler): """Register a message handler (decorator style).""" @@ -1361,25 +1359,6 @@ def on_error(self, handler): """Register an error handler (decorator style).""" self._on_error = handler - def join_room(self, room_name: str) -> None: - """Join a named room.""" - self._rooms.add(room_name) - if self._manager: - self._manager._join_room(self.id, room_name) - - def leave_room(self, room_name: str) -> None: - """Leave a named room.""" - self._rooms.discard(room_name) - if self._manager: - self._manager._leave_room(self.id, room_name) - - async def broadcast_to_room(self, room_name: str, message: str | bytes, - exclude_self: bool = False) -> None: - """Broadcast a message to all connections in a room.""" - if self._manager: - exclude = self.id if exclude_self else None - await self._manager.broadcast_to_room(room_name, message, exclude=exclude) - async def send(self, message: str | bytes): """Send a text or binary message.""" if self._closed: diff --git a/tina4_python/websocket/__init__.py b/tina4_python/websocket/__init__.py index 03ded7d6..8ad2cd4a 100644 --- a/tina4_python/websocket/__init__.py +++ b/tina4_python/websocket/__init__.py @@ -186,7 +186,42 @@ async def _read_frame(reader: asyncio.StreamReader, max_size: int = 1048576) -> return bool(fin), opcode, payload -class WebSocketConnection: +class RoomMemberMixin: + """Room membership for a WebSocket connection — shared across transports. + + Joining, leaving and room broadcasts only ever touch this connection's own + ``_rooms`` set and delegate to ``_manager``; none of it depends on how the + bytes reach the socket. The raw-socket connection and the ASGI connection + therefore share this one copy. A host must expose ``self.id``, + ``self._rooms`` (a set) and ``self._manager``. + """ + + @property + def rooms(self) -> set[str]: + """Return the set of room names this connection has joined.""" + return self._rooms + + def join_room(self, room_name: str) -> None: + """Join a named room.""" + self._rooms.add(room_name) + if self._manager: + self._manager._join_room(self.id, room_name) + + def leave_room(self, room_name: str) -> None: + """Leave a named room.""" + self._rooms.discard(room_name) + if self._manager: + self._manager._leave_room(self.id, room_name) + + async def broadcast_to_room(self, room_name: str, message: str | bytes, + exclude_self: bool = False) -> None: + """Broadcast a message to all connections in a room.""" + if self._manager: + exclude = self.id if exclude_self else None + await self._manager.broadcast_to_room(room_name, message, exclude=exclude) + + +class WebSocketConnection(RoomMemberMixin): """Represents a single WebSocket connection.""" def __init__(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter, @@ -266,31 +301,8 @@ async def broadcast_to(self, path: str, message: str | bytes): if self._manager: await self._manager.broadcast(message, path=path) - # ── Rooms ────────────────────────────────────────────────── - - @property - def rooms(self) -> set[str]: - """Return the set of room names this connection has joined.""" - return self._rooms - - def join_room(self, room_name: str) -> None: - """Join a named room.""" - self._rooms.add(room_name) - if self._manager: - self._manager._join_room(self.id, room_name) - - def leave_room(self, room_name: str) -> None: - """Leave a named room.""" - self._rooms.discard(room_name) - if self._manager: - self._manager._leave_room(self.id, room_name) - - async def broadcast_to_room(self, room_name: str, message: str | bytes, - exclude_self: bool = False) -> None: - """Broadcast a message to all connections in a room.""" - if self._manager: - exclude = self.id if exclude_self else None - await self._manager.broadcast_to_room(room_name, message, exclude=exclude) + # ── Rooms ── join_room / leave_room / broadcast_to_room / rooms come from + # RoomMemberMixin — transport-agnostic, shared with the ASGI connection. async def ping(self, data: bytes = b""): """Send a ping frame.""" @@ -968,7 +980,7 @@ def _handle_upgrade(self, reader: asyncio.StreamReader, __all__ = [ - "WebSocketServer", "WebSocketConnection", "WebSocketManager", + "WebSocketServer", "WebSocketConnection", "WebSocketManager", "RoomMemberMixin", "compute_accept_key", "build_frame", "origin_allowed", "OP_TEXT", "OP_BINARY", "OP_CLOSE", "OP_PING", "OP_PONG", "CLOSE_NORMAL", "CLOSE_GOING_AWAY", "CLOSE_PROTOCOL_ERROR", "CLOSE_TOO_LARGE", From 38fd1691bc476889c16f6854676c8fb825f5b3b4 Mon Sep 17 00:00:00 2001 From: Andre van Zuydam Date: Wed, 30 Sep 2026 04:48:54 +0200 Subject: [PATCH 6/8] refactor(database): build the LIMIT/OFFSET clause once on the base adapter MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit SQLite, MySQL and PostgreSQL each appended their own pagination clause with the same if/else: unpaginated returns the SQL and params untouched; paginated puts LIMIT/OFFSET on a new line (so a trailing -- comment cannot swallow it) with two bound values. Only the placeholder differed — ? on SQLite, %s on the servers — which each adapter already declares as PARAM_MARKER. The append now lives once as DatabaseAdapter._paginate_sql(). Each adapter still decides `paginated` for itself (SQLite and PostgreSQL exclude writes; MySQL does not), so no engine's pagination rule changes. Co-Authored-By: Claude Opus 4.8 Co-Authored-By: Tina4 <82961293+tina4stack@users.noreply.github.com> --- tina4_python/database/adapter.py | 16 ++++++++++++++++ tina4_python/database/mysql.py | 7 +------ tina4_python/database/postgres.py | 7 +------ tina4_python/database/sqlite.py | 10 +--------- 4 files changed, 19 insertions(+), 21 deletions(-) diff --git a/tina4_python/database/adapter.py b/tina4_python/database/adapter.py index 816a51a9..5490e240 100644 --- a/tina4_python/database/adapter.py +++ b/tina4_python/database/adapter.py @@ -977,6 +977,22 @@ def _total_from_page(row_count: int, limit, offset, paginated: bool) -> int | No return 0 return None + def _paginate_sql(self, sql: str, params, limit, offset, paginated: bool): + """Append this engine's LIMIT/OFFSET clause when ``paginated``. + + Returns ``(sql, params)``. Each adapter decides ``paginated`` for + itself (SQLite/PostgreSQL also exclude writes), but the append is the + same everywhere: this engine's placeholder (``?`` or ``%s``) for the + two bound values. The clause goes on a NEW LINE — appended inline it + can land inside a trailing ``-- comment`` and be swallowed, the same + bug at the append site rather than the detector. + """ + if not paginated: + return sql, params or [] + marker = self.PARAM_MARKER + return (f"{sql}\nLIMIT {marker} OFFSET {marker}", + (params or []) + [limit, offset]) + @staticmethod def _strip_trailing_order_by(sql: str) -> str: """Strip a trailing top-level ``ORDER BY`` so the SQL can be safely diff --git a/tina4_python/database/mysql.py b/tina4_python/database/mysql.py index 0cc705bb..af85dad3 100644 --- a/tina4_python/database/mysql.py +++ b/tina4_python/database/mysql.py @@ -166,12 +166,7 @@ def fetch(self, sql: str, params: list = None, # a syntax error MEASURED on a live PostgreSQL. It worked on sqlite and # crashed on the server, which is the swap ADR-0024 exists to protect. paginated = not (limit is None or limit <= 0 or self._has_trailing_limit(sql)) - if not paginated: - paginated_sql = sql - paginated_params = params or [] - else: - paginated_sql = f"{sql}\nLIMIT %s OFFSET %s" - paginated_params = (params or []) + [limit, offset] + paginated_sql, paginated_params = self._paginate_sql(sql, params, limit, offset, paginated) cursor.execute(paginated_sql, paginated_params) # FAILS LOUD rows = [dict(row) for row in cursor.fetchall()] diff --git a/tina4_python/database/postgres.py b/tina4_python/database/postgres.py index 0b882f56..04cb039d 100644 --- a/tina4_python/database/postgres.py +++ b/tina4_python/database/postgres.py @@ -440,12 +440,7 @@ def fetch(self, sql: str, params: list = None, # a syntax error MEASURED on a live PostgreSQL. It worked on sqlite and # crashed on the server, which is the swap ADR-0024 exists to protect. paginated = not (is_write or limit is None or limit <= 0 or self._has_trailing_limit(sql)) - if not paginated: - paginated_sql = sql - paginated_params = params or [] - else: - paginated_sql = f"{sql}\nLIMIT %s OFFSET %s" - paginated_params = (params or []) + [limit, offset] + paginated_sql, paginated_params = self._paginate_sql(sql, params, limit, offset, paginated) self._exec_with_handling(cursor, paginated_sql, paginated_params) columns = [d[0] for d in cursor.description] if cursor.description else [] indexes = range(len(columns)) diff --git a/tina4_python/database/sqlite.py b/tina4_python/database/sqlite.py index aa2156d3..562072c6 100644 --- a/tina4_python/database/sqlite.py +++ b/tina4_python/database/sqlite.py @@ -235,15 +235,7 @@ def fetch(self, sql: str, params: list = None, # never paginated and never COUNT-probed (a probe repeats the write). is_write = self._is_write_statement(sql) paginated = not (is_write or limit is None or limit <= 0 or self._has_trailing_limit(sql)) - if not paginated: - paginated_sql = sql - paginated_params = params or [] - else: - # The clause goes on a NEW LINE. Appended inline it lands INSIDE a - # trailing `-- comment` and is swallowed, which is the same bug at - # the append site rather than the detector. - paginated_sql = f"{sql}\nLIMIT ? OFFSET ?" - paginated_params = (params or []) + [limit, offset] + paginated_sql, paginated_params = self._paginate_sql(sql, params, limit, offset, paginated) # Hydrate via a tuple cursor + a column list computed ONCE, rather than # dict(sqlite3.Row) per row. The connection's row_factory is sqlite3.Row, # which builds a Row object per row that dict() then copies -- two From 89bcb24282b04278cf0463086878141a1b004738 Mon Sep 17 00:00:00 2001 From: Andre van Zuydam Date: Wed, 30 Sep 2026 04:54:46 +0200 Subject: [PATCH 7/8] test(rate-limit): lock the middleware enforcement entry points MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit RateLimiter.apply() and the two before_rate_limit staticmethods (the class-based middleware entry points) had no coverage — the production dispatch path uses check()/apply_headers() directly, so nothing exercised apply() or the wrappers. The dedup that routed both wrappers through apply() could have regressed with the suite still green. Adds real-object tests (real Response, plain request holder, no doubles) that drive apply(), RateLimiter.before_rate_limit and RateLimiterMiddleware.before_rate_limit over the limit and assert the 429, the retry-after header and the rate-limit headers. Mutation-proven: gutting apply()'s 429 fails 3, and cutting either wrapper's delegation fails its test. Co-Authored-By: Claude Opus 4.8 Co-Authored-By: Tina4 <82961293+tina4stack@users.noreply.github.com> --- tests/test_rate_limiter.py | 66 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 66 insertions(+) diff --git a/tests/test_rate_limiter.py b/tests/test_rate_limiter.py index 6398edab..9154a29c 100644 --- a/tests/test_rate_limiter.py +++ b/tests/test_rate_limiter.py @@ -147,3 +147,69 @@ def test_cleanup_keeps_active_ips(self, monkeypatch): now = time.monotonic() rl._cleanup(now) assert "10.0.0.1" in rl._requests + + +class TestRateLimiterEnforcement: + """Lock the enforcement entry points, not just check()/apply_headers(). + + RateLimiter.apply() holds the enforcement path, and the two before_rate_limit + staticmethods (the class-based middleware entry points) delegate to it. These + entry points had no coverage before, so a refactor could quietly stop + returning the 429 with nothing going red. Real Response object, plain request + holder — no doubles. + """ + + def _request(self, ip="10.0.0.1"): + import types + return types.SimpleNamespace(ip=ip) + + def test_apply_allows_under_limit(self, monkeypatch): + monkeypatch.setenv("TINA4_RATE_LIMIT", "3") + monkeypatch.setenv("TINA4_RATE_WINDOW", "60") + rl = RateLimiter() + resp = Response() + _, out = rl.apply(self._request(), resp) + assert out.status_code == 200 + assert dict(resp._headers).get("x-ratelimit-limit") == "3" + + def test_apply_refuses_over_limit_with_429_and_retry_after(self, monkeypatch): + monkeypatch.setenv("TINA4_RATE_LIMIT", "3") + monkeypatch.setenv("TINA4_RATE_WINDOW", "60") + rl = RateLimiter() + req = self._request() + statuses = [] + for _ in range(4): + resp = Response() + _, out = rl.apply(req, resp) + statuses.append(out.status_code) + assert statuses == [200, 200, 200, 429] + # the refused response carries retry-after and the rate-limit headers + assert dict(resp._headers).get("retry-after") is not None + assert dict(resp._headers).get("x-ratelimit-remaining") == "0" + + def test_shared_before_rate_limit_refuses_over_limit(self, monkeypatch): + monkeypatch.setenv("TINA4_RATE_LIMIT", "2") + monkeypatch.setenv("TINA4_RATE_WINDOW", "60") + RateLimiter._shared_instance = None + req = self._request("10.0.0.2") + statuses = [] + for _ in range(3): + resp = Response() + _, out = RateLimiter.before_rate_limit(req, resp) + statuses.append(out.status_code) + RateLimiter._shared_instance = None + assert statuses == [200, 200, 429] + + def test_middleware_before_rate_limit_refuses_over_limit(self, monkeypatch): + from tina4_python.core.middleware import RateLimiterMiddleware + monkeypatch.setenv("TINA4_RATE_LIMIT", "2") + monkeypatch.setenv("TINA4_RATE_WINDOW", "60") + RateLimiterMiddleware._limiter = None + req = self._request("10.0.0.3") + statuses = [] + for _ in range(3): + resp = Response() + _, out = RateLimiterMiddleware.before_rate_limit(req, resp) + statuses.append(out.status_code) + RateLimiterMiddleware._limiter = None + assert statuses == [200, 200, 429] From db9cd76c37aa5063bdcdaa4cddc2d1f21a42c40d Mon Sep 17 00:00:00 2001 From: Andre van Zuydam Date: Wed, 30 Sep 2026 05:03:35 +0200 Subject: [PATCH 8/8] test(queue): cover retry_job() on the shared topic dead-letter mixin MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit dead_letters() on the Kafka/RabbitMQ backends was already gated by the deterministic dead-letter test, but retry_job() — the other method now shared through TopicDeadLetterMixin — had only a surface "it resolves" check. Adds two live-Kafka regressions: retry(id) reports True and clears the dead-letter topic; retry(unknown-id) reports False and leaves the real dead letter in place. Re-consumption of the re-produced record is a Kafka offset concern outside this method's contract, so it is not asserted. Mutation-proven on live Kafka: forcing retry_job to return False fails the revive test; skipping its found-is-None guard fails the unknown-id test. Co-Authored-By: Claude Opus 4.8 Co-Authored-By: Tina4 <82961293+tina4stack@users.noreply.github.com> --- ...t_queue_kafka_dead_letter_deterministic.py | 40 +++++++++++++++++++ 1 file changed, 40 insertions(+) diff --git a/tests/test_queue_kafka_dead_letter_deterministic.py b/tests/test_queue_kafka_dead_letter_deterministic.py index 865e4de1..a0155d10 100644 --- a/tests/test_queue_kafka_dead_letter_deterministic.py +++ b/tests/test_queue_kafka_dead_letter_deterministic.py @@ -86,3 +86,43 @@ def test_reject_dead_letter_is_observable(self): assert job_id in ids, f"rejected job {job_id} not observable (got {ids})" finally: q.close() + + def test_retry_job_clears_the_dead_letter_and_reports_true(self): + """retry(job_id) revives a dead-lettered job: it returns True and the + job leaves the dead-letter topic. + + This exercises the shared TopicDeadLetterMixin.retry_job() on live + Kafka — the same method the RabbitMQ backend inherits. (Whether the + re-produced record is then re-consumed is a Kafka offset concern, not + part of this method's contract, so it is not asserted here.) + """ + q = _make_queue() + try: + job_id = q.push({"task": "revive-me"}) + job = q.pop() + assert job is not None + job.reject("dead for now") + assert job_id in [d.id for d in q.dead_letters()], "prime: not dead-lettered" + + assert q.retry(job_id) is True, "retry() must report it revived the job" + assert job_id not in [d.id for d in q.dead_letters()], ( + "revived job must no longer be in the dead-letter topic" + ) + finally: + q.close() + + def test_retry_job_unknown_id_returns_false(self): + """retry() on an id that is not dead-lettered reports False, without + disturbing the real dead letter.""" + q = _make_queue() + try: + job_id = q.push({"task": "stays-dead"}) + job = q.pop() + assert job is not None + job.reject("poison") + assert q.retry("no-such-id") is False + assert job_id in [d.id for d in q.dead_letters()], ( + "a failed retry lookup must leave the real dead letter in place" + ) + finally: + q.close()