Gateway (docs/19): - Remove ws_server.py; frame dispatch factored into dispatch.py - http_server: media upload/pull, pairing over HTTP - protocol: media frames mirrored; tests + ws_probe updated for HTTP App: - HttpGateway: postFrame/uploadMedia/pullMedia no longer throw on network failure (PostResult ok=false / Result.failure) — uncaught SocketTimeoutException on Dispatchers.Default crashed the app - GatewayClient: dead-stream watchdog (health probe every 10s, 2 failures -> redial in ~20s instead of the 45s SSE read timeout); state flips to Reconnecting when the stream dies, restored from the last hello.ack on long-poll success; poke() + backoff reset on app resume (MainActivity.onResume) - Offline sends: composer enabled while disconnected; a send with no response (status 0) stays queued (Pending) and is auto-resent on the next (re)connect after a 2s outbox-replay grace; gateway 4xx rejections fail the bubble (tap to retry, no auto-loop) - ChatStore: echo-replace and thread-relocate also match Failed bubbles (POST response lost in a network drop); loadHistory dedupes local failed bubbles the server already has; failMessage() - MainActivity: poke() on resume so a backgrounded app reconnects promptly instead of waiting out the backoff
85 lines
3.3 KiB
Python
85 lines
3.3 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:
|
|
"""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)
|
|
# 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
|