adapter.py was a 3,493-line monolith. Split it into focused modules with clear separation of responsibilities, bringing it down to ~857 lines: - Module-level helpers: hooks, classify, pickers, commands, setup, defaults, secrets - Frame-handler mixins: inbound, tool_frames, push_frames, media_frames, picker_frames, channel_frames, query_frames - mixin_base: IrisAdapterBase (declaration-only base for shared attrs) - adapter.py now holds only IrisAdapter (the composition of the 7 mixins + BasePlatformAdapter), register(), and test-facing re-exports The mixins come before BasePlatformAdapter in the MRO so their methods override the base; super() calls (e.g. send_image) still resolve to BasePlatformAdapter. No circular imports; dispatch.py and http_server.py (instance-method callers) are unaffected. Ruff complexity ceilings (PLR0911/0912/0913/0915) restored to Ruff's built-in defaults (12/50/6/5) instead of "just above the current maxima", which ratchets the bar down as code grows. The existing genuinely-complex functions (frame builders mirroring the wire schema, the QR matrix builder, the dispatch table) carry an explicit `# noqa: PLR09xx` marking them as reviewed, frozen exceptions; new code is held to the default ceilings. All 125 tests green (94 test_android + 31 test_android_http); no new ruff errors introduced.
87 lines
3.4 KiB
Python
87 lines
3.4 KiB
Python
"""Shared inbound frame dispatch + inbound rate limit.
|
|
|
|
Extracted from the (now-removed) WS server so the HTTP transport has a
|
|
single home for the transport-agnostic dispatch chain and the per-device
|
|
token bucket. The HTTP leg (``http_server.py``) is the only transport;
|
|
this module is transport-neutral.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import time
|
|
from typing import Any
|
|
|
|
from . import protocol
|
|
|
|
# Inbound JSON control-frame rate limit (per device, token bucket).
|
|
# A legitimate app sends occasional user-initiated requests — far below
|
|
# 20/s sustained. Media uploads are exempt (they travel via
|
|
# ``POST /v1/media``, not the frame endpoint).
|
|
INBOUND_RATE_PER_S = 20.0
|
|
INBOUND_BURST = 40
|
|
|
|
# Max length of a client-supplied device_id.
|
|
MAX_DEVICE_ID_LEN = 128
|
|
|
|
|
|
async def dispatch_frame(adapter: Any, frame: protocol.Frame, device_id: str) -> None: # noqa: PLR0912
|
|
"""Shared inbound frame dispatch (docs/19 §19.4). Unknown types are
|
|
ignored (forward-compat)."""
|
|
if frame.type == protocol.TYPE_MESSAGE_SEND:
|
|
await adapter.on_message_send(frame, device_id)
|
|
elif frame.type == protocol.TYPE_CHANNEL_CREATE:
|
|
await adapter.on_channel_create(frame, device_id)
|
|
elif frame.type == protocol.TYPE_CHANNEL_RENAME:
|
|
await adapter.on_channel_rename(frame, device_id)
|
|
elif frame.type == protocol.TYPE_CHANNEL_SET_DEFAULT:
|
|
await adapter.on_channel_set_default(frame, device_id)
|
|
elif frame.type == protocol.TYPE_CHANNEL_FAVORITE:
|
|
await adapter.on_channel_favorite(frame, device_id)
|
|
elif frame.type == protocol.TYPE_CHANNEL_ICON:
|
|
await adapter.on_channel_icon(frame, device_id)
|
|
elif frame.type == protocol.TYPE_CHANNEL_SET_AUTOMATION:
|
|
await adapter.on_channel_set_automation(frame, device_id)
|
|
elif frame.type == protocol.TYPE_CHANNEL_DELETE:
|
|
await adapter.on_channel_delete(frame, device_id)
|
|
elif frame.type == protocol.TYPE_CHANNEL_LIST:
|
|
await adapter.on_channel_list(frame, device_id)
|
|
elif frame.type == protocol.TYPE_COMMANDS_CATALOG:
|
|
await adapter.on_commands_catalog(frame, device_id)
|
|
elif frame.type == protocol.TYPE_SEARCH:
|
|
await adapter.on_search(frame, device_id)
|
|
elif frame.type == protocol.TYPE_SYNC:
|
|
await adapter.on_sync(frame, device_id)
|
|
elif frame.type == protocol.TYPE_HISTORY:
|
|
await adapter.on_history(frame, device_id)
|
|
elif frame.type == protocol.TYPE_MESSAGE_DELETE:
|
|
await adapter.on_message_delete(frame, device_id)
|
|
elif frame.type == protocol.TYPE_FCM_REGISTER:
|
|
await adapter.on_fcm_register(frame, device_id)
|
|
elif frame.type == protocol.TYPE_PICKER_SELECT:
|
|
await adapter.on_picker_select(frame, device_id)
|
|
# Unknown types are ignored (forward-compat).
|
|
|
|
|
|
class _TokenBucket:
|
|
"""Minimal token bucket (stdlib only). One instance per device."""
|
|
|
|
__slots__ = ("rate", "burst", "tokens", "updated_at")
|
|
|
|
def __init__(self, rate: float, burst: int):
|
|
self.rate = rate
|
|
self.burst = burst
|
|
self.tokens = float(burst)
|
|
self.updated_at = time.monotonic()
|
|
|
|
def consume(self) -> bool:
|
|
"""Try to take one token. Refills at ``rate``/s up to ``burst``."""
|
|
now = time.monotonic()
|
|
elapsed = now - self.updated_at
|
|
if elapsed > 0:
|
|
self.tokens = min(self.burst, self.tokens + elapsed * self.rate)
|
|
self.updated_at = now
|
|
if self.tokens >= 1.0:
|
|
self.tokens -= 1.0
|
|
return True
|
|
return False
|