Slash commands with a finite set of options (/reasoning, /fast, ...) now
render a tappable card with buttons (2 per row, ✓ on the current value)
instead of a plain text status card. The mechanism is generic: any command
that calls the adapter's send_choice_picker() gets a picker automatically.
Wire protocol (docs/04, frames.schema.json):
- picker.choice (server→app): {picker_id, title, choices[]}
- picker.select (app→server): {picker_id, value}
- pickers capability flag now True in server_caps
gateway-plugin:
- protocol.py: picker.choice/picker.select frame types + picker_choice()
- dispatch.py: route picker.select → adapter.on_picker_select
- adapter.py: send_choice_picker() (fails cleanly with no live device so
hermes falls back to text), on_picker_select(), in-memory pending pickers
(gateway restart expires them; stale select is a no-op), pickers=True
app (KMP):
- Protocol.kt: PickerChoice/PickerChoicePayload + pickerSelectFrame()
- ChatStore.kt: PickerItem + onPickerChoice (idempotent) + resolvePicker
(optimistic, one-shot)
- ChatDb.kt: persist PickerItem in the messages table (polymorphic decode)
- IrisController.kt: picker.choice routing + selectPicker() action
- ChatScreen.kt: PickerCard composable (locks after selection)
Tests:
- python: 3 picker tests (roundtrip, no-device fallback, stale-select noop)
- kotlin: ChatStorePickerTest (add/idempotent/resolve/one-shot/noop/serialize)
- fixture fix: clear leaked IRIS_HTTP_PORT/IRIS_WS_HOST env so the adapter
binds the ephemeral port (a prior test's interactive_setup() polluted the
process env, colliding with a live gateway on 8791)
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
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:
|
|
"""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
|