M3: channels/threads + cron delivery + search + outbox sync

Gateway plugin:
- ChannelDirectory (SQLite): channels + threads, friendly-name resolution
- FTS5 search bridge over state.db (scope all/chat, LIKE fallback)
- Outbox (monotonic cursor, 72h retention) for reconnect catch-up
- adapter: channel.* handlers, search, sync, cron target parsing
- ws_server: M3 frame routing + hello.ack sync cursor

App (KMP):
- ChannelStore (directory cache) + lane-aware ChatStore (per chat/thread)
- GatewayClient.sendFrame; IrisController channel/search/nav actions
- ChatScreen: channel drawer, thread toggle, topic switcher, search overlay
This commit is contained in:
ARIA committed 2026-08-19 14:53:48 +02:00
1 parent 218c50d688
commit 9b511fdd28
13 files changed
+2035 -259

No files matched your search

+294 -37
View File
@@ -92,6 +92,9 @@ from gateway.config import Platform # noqa: E402
from hermes_constants import get_hermes_home # noqa: E402
from . import protocol # noqa: E402
from . import search as search_bridge # noqa: E402
from .channels import get_directory # noqa: E402
from .outbox import Outbox # noqa: E402
from .pairing import ( # noqa: E402
DeviceRegistry,
generate_token,
@@ -471,24 +474,45 @@ def _parse_port(raw: str) -> int:
def _parse_target_ref(target_ref: str) -> Optional[tuple]:
"""Parse a raw target string into ``(chat_id, thread_id)`` or ``None``.
Recognises the native syntax ``android:<chat>[:<thread>]``. Returns
``None`` for anything else so the target proceeds to channel-directory
Recognises the native syntax ``android:<chat>[:<thread>]`` where the
chat_id itself carries the ``android:`` prefix (e.g. ``android:chan_7``)
and an optional thread is a trailing ``:t_<n>``. A bare friendly name
(e.g. ``Cron Reports``) is resolved against the channel directory so cron
/ ``send_message`` can target a channel by name immediately, without
waiting for the core directory's refresh timer. Returns ``None`` for
anything unrecognised so the target proceeds to the core channel-directory
resolution.
"""
if not target_ref or not target_ref.startswith("android:"):
if not target_ref:
return None
body = target_ref[len("android:"):]
if not body:
t = target_ref.strip()
if not t:
return None
if ":" in body:
chat_id, thread_id = body.split(":", 1)
thread_id = thread_id or None
else:
chat_id, thread_id = body, None
chat_id = chat_id.strip()
if not chat_id:
return None
return (chat_id, thread_id)
if t.startswith("android:"):
body = t[len("android:"):].strip()
if not body:
return None
thread_id: Optional[str] = None
if ":" in body:
head, tail = body.rsplit(":", 1)
if tail and tail.startswith("t_"):
thread_id = tail
body = head
return (f"android:{body}", thread_id)
# Bare friendly name -> resolve via the channel directory. A thread resolves
# to its session lane (parent_chat_id + thread_id); a channel/default to
# its chat_id.
try:
entry = get_directory().resolve_entry(t)
except Exception:
entry = None
if entry is not None:
if entry["kind"] == "thread":
return (entry["parent_chat_id"], entry["chat_id"])
return (entry["chat_id"], None)
return None
# ---------------------------------------------------------------------------
@@ -648,6 +672,12 @@ class AndroidAdapter(BasePlatformAdapter):
self._connected = False
# M2: per-chat turn state for outbound frame classification.
self._turns: Dict[str, _TurnState] = {}
# M3: channel directory (shared singleton) + offline outbox.
self._channels = get_directory()
self._outbox = Outbox(
get_hermes_home() / "android" / "outbox.db",
retention_hours=self.outbox_retention_hours,
)
def _turn_state(self, chat_id: str) -> _TurnState:
st = self._turns.get(chat_id)
@@ -695,6 +725,13 @@ class AndroidAdapter(BasePlatformAdapter):
self._connected = False
return False
# M3: ensure the default (home) channel exists in the directory so the
# app's channel list and cron home delivery have a stable anchor.
try:
self._channels.ensure_default(self.home_channel, self.home_channel_name)
except Exception:
logger.warning("android: ensure_default failed", exc_info=True)
self._connected = True
self._mark_connected()
logger.info("android: connected; WS server on %s:%s", self.host, self.port)
@@ -716,6 +753,10 @@ class AndroidAdapter(BasePlatformAdapter):
self._devices.close()
except Exception:
pass
try:
self._outbox.close()
except Exception:
pass
self._connected = False
self._mark_disconnected()
logger.info("android: disconnected")
@@ -984,10 +1025,16 @@ class AndroidAdapter(BasePlatformAdapter):
async def _broadcast_or_log(self, chat_id: str, frame: "protocol.Frame") -> None:
delivered = await self._ws_server.broadcast(frame)
if delivered == 0:
logger.info(
"android: no live devices for %s; %s frame not delivered (outbox lands in M3)",
chat_id, frame.type,
)
# M3: no live device -- park the frame in the outbox so a
# reconnecting app can `sync` it (M5 adds push to wake the device).
try:
cursor = self._outbox.append(chat_id, frame.to_json())
logger.info(
"android: no live devices for %s; %s frame parked in outbox (cursor=%s)",
chat_id, frame.type, cursor,
)
except Exception:
logger.warning("android: outbox append failed", exc_info=True)
async def send_typing(self, chat_id: str, metadata: Optional[Dict[str, Any]] = None) -> None:
"""Send a typing indicator (``typing`` frame, on=true)."""
@@ -1080,50 +1127,260 @@ class AndroidAdapter(BasePlatformAdapter):
)
await self.handle_message(event)
# ── M3: channel directory management (app -> agent) ───────────────────
#
# Each request is answered by broadcasting the matching ``channel.*``
# event carrying the request ``id``: the requester's pending request
# completes on the id, and every other device reconciles its local copy
# from the same frame (single broadcast serves as event + response).
async def on_channel_create(self, frame: protocol.Frame, device_id: str) -> None:
payload = frame.payload
name = payload.get("name")
if not isinstance(name, str) or not name.strip():
await self._ws_server.send_to(
device_id,
protocol.error(protocol.ERR_UNSUPPORTED, "channel.create requires a name", id=frame.id),
)
return
kind = payload.get("kind")
kind = kind if kind in ("channel", "thread") else "channel"
parent_chat_id = payload.get("parent_chat_id")
if not isinstance(parent_chat_id, str) or not parent_chat_id.strip():
parent_chat_id = None
if kind == "thread" and not parent_chat_id:
parent_chat_id = frame.chat_id or self.home_channel
try:
entry = self._channels.create(name=name, kind=kind, parent_chat_id=parent_chat_id)
except ValueError as e:
await self._ws_server.send_to(
device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id)
)
return
resp = protocol.channel_created(entry)
resp.id = frame.id
await self._ws_server.broadcast(resp)
async def on_channel_rename(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._ws_server.send_to(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, "channel.rename requires chat_id", id=frame.id),
)
return
name = frame.payload.get("name")
if not isinstance(name, str) or not name.strip():
await self._ws_server.send_to(
device_id,
protocol.error(protocol.ERR_UNSUPPORTED, "channel.rename requires a name", id=frame.id),
)
return
try:
entry = self._channels.rename(chat_id, name)
except ValueError as e:
await self._ws_server.send_to(
device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id)
)
return
if entry is None:
await self._ws_server.send_to(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id),
)
return
resp = protocol.channel_renamed(entry)
resp.id = frame.id
await self._ws_server.broadcast(resp)
async def on_channel_set_default(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._ws_server.send_to(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, "channel.set_default requires chat_id", id=frame.id),
)
return
entry = self._channels.set_default(chat_id)
if entry is None:
await self._ws_server.send_to(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id),
)
return
# Reuse the renamed event shape: it carries the full entry (incl. the
# new is_default flag) so every device reconciles the default change.
resp = protocol.channel_renamed(entry)
resp.id = frame.id
await self._ws_server.broadcast(resp)
async def on_channel_delete(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._ws_server.send_to(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, "channel.delete requires chat_id", id=frame.id),
)
return
entry = self._channels.delete(chat_id)
if entry is None:
await self._ws_server.send_to(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, f"cannot delete {chat_id} (unknown or default)", id=frame.id),
)
return
resp = protocol.channel_deleted(chat_id)
resp.id = frame.id
await self._ws_server.broadcast(resp)
async def on_channel_list(self, frame: protocol.Frame, device_id: str) -> None:
channels = self._channels.list(include_archived=False)
resp = protocol.channel_list(channels)
resp.id = frame.id
await self._ws_server.send_to(device_id, resp)
# ── M3: search (app -> agent) ─────────────────────────────────────────
async def on_search(self, frame: protocol.Frame, device_id: str) -> None:
payload = frame.payload
query = payload.get("query")
if not isinstance(query, str) or not query.strip():
await self._ws_server.send_to(
device_id, protocol.error(protocol.ERR_UNSUPPORTED, "search requires a query", id=frame.id)
)
return
scope = payload.get("scope")
scope = scope if scope in ("all", "chat") else "all"
chat_id = payload.get("chat_id") or frame.chat_id
if not isinstance(chat_id, str) or not chat_id.strip():
chat_id = None
thread_id = payload.get("thread_id") or frame.thread_id
if not isinstance(thread_id, str) or not thread_id.strip():
thread_id = None
limit = payload.get("limit")
try:
limit = int(limit) if limit is not None else 20
except (TypeError, ValueError):
limit = 20
db_path = get_hermes_home() / "state.db"
hits = search_bridge.search(
db_path, query, scope=scope, chat_id=chat_id, thread_id=thread_id, limit=limit
)
resp = protocol.search_results(query, scope, hits, id=frame.id)
await self._ws_server.send_to(device_id, resp)
# ── M3: sync (reconnect catch-up) ─────────────────────────────────────
async def on_sync(self, frame: protocol.Frame, device_id: str) -> None:
payload = frame.payload
cursor = payload.get("cursor")
try:
cursor = int(cursor) if cursor is not None else 0
except (TypeError, ValueError):
cursor = 0
for e in self._outbox.replay(cursor):
raw = e["frame"]
replayed = protocol.Frame(
type=raw.get("type", ""),
payload=raw.get("payload", {}) if isinstance(raw.get("payload"), dict) else {},
id=raw.get("id") if isinstance(raw.get("id"), int) else None,
chat_id=raw.get("chat_id") if isinstance(raw.get("chat_id"), str) else e.get("chat_id"),
thread_id=raw.get("thread_id") if isinstance(raw.get("thread_id"), str) else None,
v=raw.get("v") if isinstance(raw.get("v"), int) else protocol.PROTOCOL_VERSION,
)
await self._ws_server.send_to(device_id, replayed)
done = protocol.sync_done(self._outbox.latest_cursor(), id=frame.id)
await self._ws_server.send_to(device_id, done)
# ── Chat info ─────────────────────────────────────────────────────────
def _channel_name(self, chat_id: str) -> str:
"""Channel display name. M1: home channel only (directory is M3)."""
"""Channel display name (M3: from the channel directory)."""
if not chat_id:
return "chat"
entry = self._channels.get(chat_id)
if entry is not None:
return entry["name"]
if chat_id in (self.home_channel, DEFAULT_HOME_CHANNEL):
return self.home_channel_name
return chat_id or "chat"
return chat_id
async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
"""Return ``{name, type, chat_id}`` for a chat.
M1: the channel directory is not persisted yet, so report the home
channel name for the default chat and the raw id otherwise.
"""
"""Return ``{name, type, chat_id}`` for a chat (M3: directory-backed)."""
entry = self._channels.get(chat_id)
kind = entry["kind"] if entry else "channel"
return {
"name": self._channel_name(chat_id),
"type": "channel",
"type": "dm" if kind == "default" else "channel",
"chat_id": chat_id,
}
# ── hello.ack helpers ─────────────────────────────────────────────────
def server_caps(self) -> Dict[str, Any]:
"""Capability flags advertised in ``hello.ack`` (M2 surface)."""
"""Capability flags advertised in ``hello.ack`` (M3 surface)."""
return {
"streaming": True, # M2: message.start/update/stop
"reasoning": True, # M2: reasoning field on message / message.stop
"tools": True, # M2: tool.start/progress/end
"media": False, # M4
"search": False, # M3
"search": True, # M3: search frame
"push": self.push_backend,
"pickers": False, # M2+
}
def channel_list(self) -> List[Dict[str, Any]]:
"""Channel directory for ``hello.ack``. M1: home channel only (M3)."""
return [
{
"chat_id": self.home_channel,
"name": self.home_channel_name,
"kind": "default",
"is_default": True,
}
]
"""Channel directory for ``hello.ack`` (M3: full non-archived list)."""
return self._channels.list(include_archived=False)
# ── M3: core channel-directory hook (cron / send_message name resolution)
async def list_channels(self) -> List[Dict[str, Any]]:
"""Expose the directory to the gateway's core channel directory.
``gateway/channel_directory.build_channel_directory`` calls this to
populate ``channel_directory.json``, which ``resolve_channel_name``
reads for friendly-name -> chat_id resolution (cron + send_message).
Threads are addressed via the explicit ``android:<chat>:<thread>``
syntax (see ``_parse_target_ref``), so only channels are listed here.
"""
out: List[Dict[str, Any]] = []
for entry in self._channels.list(include_archived=False):
if entry["kind"] == "thread":
continue
out.append(
{
"id": entry["chat_id"],
"name": entry["name"],
"type": "dm" if entry["kind"] == "default" else "channel",
}
)
return out
# ── M3: thread handoff (gateway create_handoff_thread) ────────────────
async def create_handoff_thread(
self, parent_chat_id: str, name: str
) -> Optional[str]:
"""Mint a named thread under *parent_chat_id* (gateway handoff path).
Returns the new ``thread_id`` (``t_<n>``) so the handed-off session is
isolated in its own lane, or ``None`` when the parent is unknown.
"""
parent = self._channels.get(parent_chat_id)
if parent is None:
# Unknown parent: still mint a thread under it so the handoff has a
# lane (the directory row is created lazily on first use).
parent_chat_id = parent_chat_id or self.home_channel
try:
entry = self._channels.create(
name=name or "Handoff", kind="thread", parent_chat_id=parent_chat_id
)
except Exception:
logger.warning("android: create_handoff_thread failed", exc_info=True)
return None
await self._ws_server.broadcast(protocol.channel_created(entry))
return entry["chat_id"]
# ---------------------------------------------------------------------------
+353
View File
@@ -0,0 +1,353 @@
"""Channel directory (SQLite) -- the source of truth for the app's channel list.
Maps app concepts onto hermes' existing ``chat_id`` / ``thread_id`` primitives
(docs/06-channels-cron-search.md §6.1):
* **default chat** -> the home channel (``ANDROID_HOME_CHANNEL``, default
``android:default``), ``kind="default"``, ``is_default=1``.
* **user channel** -> a minted ``chat_id = android:chan_<n>``, ``kind="channel"``.
* **thread** -> a minted ``thread_id = t_<n>`` under a ``chat_id``,
``kind="thread"`` (stored with its ``parent_chat_id``).
The directory is the single source of truth for the app's channel list and for
cron name resolution. It is also exposed to the gateway's core channel
directory (``gateway/channel_directory.py``) via the adapter's
``list_channels()`` hook, so ``send_message`` / cron can resolve a friendly
name (e.g. "Cron Reports") to a chat_id.
Storage: ``get_hermes_home()/"android"/channels.db``.
Milestone M3.
"""
import logging
import sqlite3
import threading
import time
from pathlib import Path
from typing import Any, Dict, List, Optional
logger = logging.getLogger(__name__)
# Entry ``kind`` values.
KIND_DEFAULT = "default"
KIND_CHANNEL = "channel"
KIND_THREAD = "thread"
# chat_id / thread_id minting prefixes.
CHANNEL_PREFIX = "android:chan_"
THREAD_PREFIX = "t_"
class ChannelDirectory:
"""Persistent channel directory under ``get_hermes_home()/"android"``.
Thread-safe (single connection + lock); all operations are small and fast
enough to run inline on the gateway's asyncio loop. Mirrors the
``DeviceRegistry`` pattern (``pairing.py``).
"""
def __init__(self, db_path: Path):
self._db_path = Path(db_path)
self._db_path.parent.mkdir(parents=True, exist_ok=True)
self._lock = threading.Lock()
self._conn = sqlite3.connect(str(self._db_path), check_same_thread=False)
self._conn.row_factory = sqlite3.Row
with self._lock:
self._conn.execute("PRAGMA journal_mode=WAL")
self._conn.execute(
"""
CREATE TABLE IF NOT EXISTS channels (
chat_id TEXT PRIMARY KEY,
name TEXT NOT NULL,
kind TEXT NOT NULL DEFAULT 'channel',
parent_chat_id TEXT,
is_default INTEGER NOT NULL DEFAULT 0,
archived INTEGER NOT NULL DEFAULT 0,
created REAL NOT NULL DEFAULT 0
)
"""
)
self._conn.execute(
"""
CREATE TABLE IF NOT EXISTS counters (
name TEXT PRIMARY KEY,
value INTEGER NOT NULL DEFAULT 0
)
"""
)
self._conn.commit()
# ── id minting ────────────────────────────────────────────────────────
def _next_counter(self, name: str) -> int:
with self._lock:
self._conn.execute(
"INSERT INTO counters (name, value) VALUES (?, 1) "
"ON CONFLICT(name) DO UPDATE SET value = value + 1",
(name,),
)
row = self._conn.execute(
"SELECT value FROM counters WHERE name = ?", (name,)
).fetchone()
self._conn.commit()
return int(row["value"]) if row else 1
def _mint_chat_id(self, kind: str) -> str:
if kind == KIND_THREAD:
return f"{THREAD_PREFIX}{self._next_counter('thread')}"
return f"{CHANNEL_PREFIX}{self._next_counter('chan')}"
# ── default channel ───────────────────────────────────────────────────
def ensure_default(self, chat_id: str, name: str) -> Dict[str, Any]:
"""Ensure the default (home) channel exists. Idempotent.
If a row already exists for *chat_id* it is kept (name refreshed only
when it was still the auto default); if another row is marked default
it is cleared so exactly one default exists.
"""
chat_id = (chat_id or "android:default").strip() or "android:default"
name = (name or "Default").strip() or "Default"
with self._lock:
existing = self._conn.execute(
"SELECT name, is_default FROM channels WHERE chat_id = ?",
(chat_id,),
).fetchone()
if existing is None:
self._conn.execute(
"INSERT INTO channels (chat_id, name, kind, is_default, created) "
"VALUES (?, ?, ?, 1, ?)",
(chat_id, name, KIND_DEFAULT, time.time()),
)
else:
# Refresh the name only if it was never renamed by the user
# (heuristic: still equals the previous default name is not
# trackable, so leave user-renamed names alone).
self._conn.execute(
"UPDATE channels SET kind = ?, is_default = 1 WHERE chat_id = ?",
(KIND_DEFAULT, chat_id),
)
# Exactly one default: clear any other default flag.
self._conn.execute(
"UPDATE channels SET is_default = 0 WHERE chat_id != ?", (chat_id,)
)
self._conn.commit()
entry = self.get(chat_id)
if entry is not None:
return entry
return {
"chat_id": chat_id, "name": name, "kind": KIND_DEFAULT,
"parent_chat_id": None, "is_default": True, "archived": False,
"created": time.time(),
}
# ── create / rename / set_default / delete ────────────────────────────
def create(
self,
name: str,
kind: str = KIND_CHANNEL,
parent_chat_id: Optional[str] = None,
) -> Dict[str, Any]:
"""Mint a new channel (or thread) and store it. Returns the entry."""
name = (name or "").strip()
if not name:
raise ValueError("channel name required")
if kind not in (KIND_CHANNEL, KIND_THREAD):
kind = KIND_CHANNEL
if kind == KIND_THREAD and not parent_chat_id:
raise ValueError("thread requires a parent_chat_id")
chat_id = self._mint_chat_id(kind)
now = time.time()
with self._lock:
self._conn.execute(
"INSERT INTO channels (chat_id, name, kind, parent_chat_id, created) "
"VALUES (?, ?, ?, ?, ?)",
(chat_id, name, kind, parent_chat_id, now),
)
self._conn.commit()
entry = self.get(chat_id)
if entry is not None:
return entry
return {
"chat_id": chat_id, "name": name, "kind": kind,
"parent_chat_id": parent_chat_id, "is_default": False,
"archived": False, "created": now,
}
def rename(self, chat_id: str, name: str) -> Optional[Dict[str, Any]]:
name = (name or "").strip()
if not name:
raise ValueError("channel name required")
with self._lock:
cur = self._conn.execute(
"UPDATE channels SET name = ? WHERE chat_id = ? AND archived = 0",
(name, chat_id),
)
self._conn.commit()
if cur.rowcount == 0:
return None
return self.get(chat_id)
def set_default(self, chat_id: str) -> Optional[Dict[str, Any]]:
"""Mark *chat_id* as the default channel (clears the previous one)."""
with self._lock:
row = self._conn.execute(
"SELECT 1 FROM channels WHERE chat_id = ? AND archived = 0",
(chat_id,),
).fetchone()
if row is None:
return None
self._conn.execute("UPDATE channels SET is_default = 0")
self._conn.execute(
"UPDATE channels SET is_default = 1 WHERE chat_id = ?", (chat_id,)
)
self._conn.commit()
return self.get(chat_id)
def delete(self, chat_id: str) -> Optional[Dict[str, Any]]:
"""Soft-delete (archive) a channel. History stays for search.
The default channel cannot be deleted. Returns the (archived) entry,
or ``None`` when the id is unknown / is the default.
"""
with self._lock:
row = self._conn.execute(
"SELECT is_default FROM channels WHERE chat_id = ?", (chat_id,)
).fetchone()
if row is None or row["is_default"]:
return None
self._conn.execute(
"UPDATE channels SET archived = 1 WHERE chat_id = ?", (chat_id,)
)
self._conn.commit()
return self.get(chat_id)
# ── reads ─────────────────────────────────────────────────────────────
def get(self, chat_id: str) -> Optional[Dict[str, Any]]:
with self._lock:
row = self._conn.execute(
"SELECT * FROM channels WHERE chat_id = ?", (chat_id,)
).fetchone()
return _row_to_entry(row) if row else None
def list(self, include_archived: bool = False) -> List[Dict[str, Any]]:
"""Directory listing. Default channel first, then creation order."""
sql = "SELECT * FROM channels"
if not include_archived:
sql += " WHERE archived = 0"
sql += " ORDER BY is_default DESC, created ASC"
with self._lock:
rows = self._conn.execute(sql).fetchall()
return [_row_to_entry(r) for r in rows]
def default(self) -> Optional[Dict[str, Any]]:
with self._lock:
row = self._conn.execute(
"SELECT * FROM channels WHERE is_default = 1 LIMIT 1"
).fetchone()
return _row_to_entry(row) if row else None
def threads_for(self, chat_id: str) -> List[Dict[str, Any]]:
"""All (non-archived) threads under *chat_id*, oldest first."""
with self._lock:
rows = self._conn.execute(
"SELECT * FROM channels WHERE kind = ? AND parent_chat_id = ? "
"AND archived = 0 ORDER BY created ASC",
(KIND_THREAD, chat_id),
).fetchall()
return [_row_to_entry(r) for r in rows]
def resolve_entry(self, name: str) -> Optional[Dict[str, Any]]:
"""Resolve a friendly name to a directory entry (case-insensitive).
Matches non-archived channels/threads by exact name first, then by
unique prefix. Returns the entry dict, or ``None`` when nothing (or
more than one) matches.
"""
query = (name or "").strip().lower()
if not query:
return None
with self._lock:
rows = self._conn.execute(
"SELECT * FROM channels WHERE archived = 0"
).fetchall()
entries = [_row_to_entry(r) for r in rows]
exact = [e for e in entries if (e["name"] or "").strip().lower() == query]
if len(exact) == 1:
return exact[0]
if len(exact) > 1:
return None
prefix = [
e for e in entries if (e["name"] or "").strip().lower().startswith(query)
]
if len(prefix) == 1:
return prefix[0]
return None
def resolve_name(self, name: str) -> Optional[str]:
"""Resolve a friendly name to a valid chat_id (case-insensitive).
For a thread, returns the *parent* chat_id (the thread's session lane
is ``parent_chat_id`` + ``thread_id``; the bare thread id is not a
standalone chat). Used by the core channel directory.
"""
entry = self.resolve_entry(name)
if entry is None:
return None
if entry["kind"] == KIND_THREAD:
return entry["parent_chat_id"]
return entry["chat_id"]
def close(self) -> None:
with self._lock:
try:
self._conn.close()
except Exception:
pass
def _row_to_entry(row: sqlite3.Row) -> Dict[str, Any]:
return {
"chat_id": row["chat_id"],
"name": row["name"],
"kind": row["kind"],
"parent_chat_id": row["parent_chat_id"],
"is_default": bool(row["is_default"]),
"archived": bool(row["archived"]),
"created": row["created"],
}
# ---------------------------------------------------------------------------
# Module-level lazy singleton (profile-aware)
#
# ``parse_target_ref_fn`` is registered at plugin load (before any adapter is
# constructed) and needs name resolution, so it reaches the directory through
# this singleton. The adapter uses the same singleton so both paths agree.
# Keyed on ``get_hermes_home()`` so a profile switch rebuilds it.
# ---------------------------------------------------------------------------
_directory: Optional[ChannelDirectory] = None
_directory_home: Optional[Path] = None
_directory_lock = threading.Lock()
def get_directory() -> ChannelDirectory:
"""Return the process-wide channel directory for the active profile."""
global _directory, _directory_home
from hermes_constants import get_hermes_home
home = Path(get_hermes_home())
with _directory_lock:
if _directory is None or _directory_home != home:
if _directory is not None:
try:
_directory.close()
except Exception:
pass
_directory = ChannelDirectory(home / "android" / "channels.db")
_directory_home = home
return _directory
+154 -5
View File
@@ -1,10 +1,159 @@
"""SQLite offline outbox + monotonic sync cursor.
"""SQLite offline outbox + monotonic sync cursor (M3).
Undelivered frames are appended per ``chat_id`` with a monotonic cursor so a
reconnecting app can ``sync {cursor}`` the delta without re-reading full
history. Retention prunes old entries (``outbox_retention_hours``).
Undelivered frames (sent while no device is live) are appended with a
**monotonic** cursor so a reconnecting app can ``sync {cursor}`` the delta
without re-reading full history. The cursor is a separate high-water counter
that only ever increases -- pruning old rows never resets it, so a late
reconnect can't be handed a cursor lower than one it already saw.
Retention prunes rows older than ``outbox_retention_hours`` (default 72h). A
device offline longer than the window misses those frames; it recovers full
context via ``history`` (M5 wires push so the device is woken to sync).
Storage: ``get_hermes_home()/"android"/outbox.db``.
Milestone M3 (built), extended in M5 (push integration).
"""
"""
import json
import logging
import sqlite3
import threading
import time
from pathlib import Path
from typing import Any, Dict, List, Optional
logger = logging.getLogger(__name__)
DEFAULT_RETENTION_HOURS = 72
_REPLAY_LIMIT = 1000
_PRUNE_INTERVAL_S = 3600.0
class Outbox:
"""Persistent outbox under ``get_hermes_home()/"android"``.
Thread-safe (single connection + lock); operations are small and fast
enough to run inline on the gateway's asyncio loop (mirrors
``DeviceRegistry`` / ``ChannelDirectory``).
"""
def __init__(self, db_path: Path, retention_hours: int = DEFAULT_RETENTION_HOURS):
self._db_path = Path(db_path)
self._db_path.parent.mkdir(parents=True, exist_ok=True)
self._retention_hours = max(1, int(retention_hours))
self._lock = threading.Lock()
self._last_prune = 0.0
self._conn = sqlite3.connect(str(self._db_path), check_same_thread=False)
self._conn.row_factory = sqlite3.Row
with self._lock:
self._conn.execute("PRAGMA journal_mode=WAL")
self._conn.execute(
"""
CREATE TABLE IF NOT EXISTS outbox (
cursor INTEGER PRIMARY KEY,
chat_id TEXT,
frame TEXT NOT NULL,
created REAL NOT NULL DEFAULT 0
)
"""
)
self._conn.execute(
"CREATE INDEX IF NOT EXISTS idx_outbox_created ON outbox (created)"
)
self._conn.execute(
"""
CREATE TABLE IF NOT EXISTS counters (
name TEXT PRIMARY KEY,
value INTEGER NOT NULL DEFAULT 0
)
"""
)
self._conn.commit()
# ── append / cursor ───────────────────────────────────────────────────
def append(self, chat_id: Optional[str], frame_json: str) -> int:
"""Append a frame; returns the (monotonic) cursor assigned to it."""
now = time.time()
with self._lock:
self._conn.execute(
"INSERT INTO counters (name, value) VALUES ('cursor', 1) "
"ON CONFLICT(name) DO UPDATE SET value = value + 1"
)
row = self._conn.execute(
"SELECT value FROM counters WHERE name = 'cursor'"
).fetchone()
cursor = int(row["value"]) if row else 1
self._conn.execute(
"INSERT INTO outbox (cursor, chat_id, frame, created) "
"VALUES (?, ?, ?, ?)",
(cursor, chat_id, frame_json, now),
)
self._conn.commit()
self._maybe_prune()
return cursor
def latest_cursor(self) -> int:
"""The high-water cursor (0 when nothing has been appended)."""
with self._lock:
row = self._conn.execute(
"SELECT value FROM counters WHERE name = 'cursor'"
).fetchone()
return int(row["value"]) if row else 0
# ── replay ────────────────────────────────────────────────────────────
def replay(self, cursor: int, limit: int = _REPLAY_LIMIT) -> List[Dict[str, Any]]:
"""Frames with ``cursor > `cursor```, oldest first.
Each entry: ``{cursor, chat_id, frame}`` where ``frame`` is the parsed
frame dict (the caller re-serializes / forwards it to the device).
"""
cursor = max(0, int(cursor or 0))
limit = max(1, min(int(limit or _REPLAY_LIMIT), _REPLAY_LIMIT))
with self._lock:
rows = self._conn.execute(
"SELECT cursor, chat_id, frame FROM outbox "
"WHERE cursor > ? ORDER BY cursor ASC LIMIT ?",
(cursor, limit),
).fetchall()
out: List[Dict[str, Any]] = []
for r in rows:
try:
frame = json.loads(r["frame"])
except (json.JSONDecodeError, TypeError):
continue
if not isinstance(frame, dict):
continue
out.append(
{"cursor": int(r["cursor"]), "chat_id": r["chat_id"], "frame": frame}
)
return out
# ── retention ─────────────────────────────────────────────────────────
def _maybe_prune(self) -> None:
now = time.time()
if now - self._last_prune < _PRUNE_INTERVAL_S:
return
self._last_prune = now
cutoff = now - self._retention_hours * 3600
with self._lock:
try:
self._conn.execute("DELETE FROM outbox WHERE created < ?", (cutoff,))
self._conn.commit()
except sqlite3.Error as e:
logger.debug("android outbox: prune failed: %s", e)
def prune(self) -> None:
"""Force a retention prune (ignores the interval throttle)."""
self._last_prune = 0.0
self._maybe_prune()
def close(self) -> None:
with self._lock:
try:
self._conn.close()
except Exception:
pass
+93 -1
View File
@@ -17,7 +17,7 @@ Milestone M5: notification, fcm.register, read.receipt.
import json
from dataclasses import dataclass, field
from typing import Any, Dict, Optional
from typing import Any, Dict, List, Optional
PROTOCOL_VERSION = 1
@@ -50,6 +50,24 @@ TYPE_TOOL_END = "tool.end"
# Intermediate assistant beat (M2)
TYPE_COMMENTARY = "commentary"
# Channels / threads (M3)
TYPE_CHANNEL_CREATE = "channel.create"
TYPE_CHANNEL_RENAME = "channel.rename"
TYPE_CHANNEL_SET_DEFAULT = "channel.set_default"
TYPE_CHANNEL_DELETE = "channel.delete"
TYPE_CHANNEL_CREATED = "channel.created"
TYPE_CHANNEL_RENAMED = "channel.renamed"
TYPE_CHANNEL_DELETED = "channel.deleted"
TYPE_CHANNEL_LIST = "channel.list"
# Search (M3)
TYPE_SEARCH = "search"
TYPE_SEARCH_RESULTS = "search.results"
# Reconnect catch-up (M3 outbox; extended by M5 push)
TYPE_SYNC = "sync"
TYPE_SYNC_DONE = "sync.done"
# ---------------------------------------------------------------------------
# Error codes (``error`` frame payload.code)
# ---------------------------------------------------------------------------
@@ -352,6 +370,80 @@ def commentary(
)
# ---------------------------------------------------------------------------
# Channel directory frames (M3)
# ---------------------------------------------------------------------------
def _channel_payload(entry: Dict[str, Any]) -> Dict[str, Any]:
"""Project a directory entry onto the wire shape."""
payload: Dict[str, Any] = {
"chat_id": entry.get("chat_id"),
"name": entry.get("name"),
"kind": entry.get("kind", "channel"),
}
if entry.get("parent_chat_id") is not None:
payload["parent_chat_id"] = entry["parent_chat_id"]
if entry.get("is_default"):
payload["is_default"] = True
if entry.get("archived"):
payload["archived"] = True
return payload
def channel_created(entry: Dict[str, Any]) -> Frame:
"""Broadcast: a channel/thread was created."""
return Frame(type=TYPE_CHANNEL_CREATED, payload=_channel_payload(entry))
def channel_renamed(entry: Dict[str, Any]) -> Frame:
"""Broadcast: a channel/thread was renamed."""
return Frame(type=TYPE_CHANNEL_RENAMED, payload=_channel_payload(entry))
def channel_deleted(chat_id: str) -> Frame:
"""Broadcast: a channel was archived (soft-deleted)."""
return Frame(type=TYPE_CHANNEL_DELETED, payload={"chat_id": chat_id})
def channel_list(channels: List[Dict[str, Any]]) -> Frame:
"""Full directory (response to a ``channel.list`` request)."""
return Frame(
type=TYPE_CHANNEL_LIST,
payload={"channels": [_channel_payload(c) for c in channels]},
)
# ---------------------------------------------------------------------------
# Search frames (M3)
# ---------------------------------------------------------------------------
def search_results(
query: str,
scope: str,
hits: List[Dict[str, Any]],
*,
id: Optional[int] = None,
) -> Frame:
"""Response to a ``search`` request.
Each hit: ``{message_id, chat_id, thread_id, role, snippet, ts}``.
"""
return Frame(
type=TYPE_SEARCH_RESULTS,
id=id,
payload={"query": query, "scope": scope, "hits": hits},
)
# ---------------------------------------------------------------------------
# Sync frames (M3 outbox)
# ---------------------------------------------------------------------------
def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame:
"""Terminal frame of a ``sync`` replay: the new cursor to persist."""
return Frame(type=TYPE_SYNC_DONE, id=id, payload={"cursor": cursor})
def error(code: str, message: str, *, id: Optional[int] = None) -> Frame:
return Frame(type=TYPE_ERROR, id=id, payload={"code": code, "message": message})
+233 -5
View File
@@ -1,8 +1,236 @@
"""FTS5 session search bridge.
"""FTS5 session search bridge (M3).
Bridges the ``search {query, scope, chat_id?, thread_id?}`` frame to the
hermes session store (SQLite + FTS5, ``hermes_state_search.py``) and returns
``search.results``. Scope: "all" (everywhere) or "chat" (this chat/channel).
Bridges the ``search {query, scope, chat_id?, thread_id?}`` frame to the hermes
session store (SQLite + FTS5, ``hermes_state.py`` / ``hermes_state_search.py``)
and returns ``search.results`` hits.
The session DB (``get_hermes_home()/"state.db"``) is opened **read-only** --
search never writes to the store. Two query paths:
* **FTS5** (primary): ``messages_fts MATCH <sanitized>`` with BM25 ranking.
* **LIKE** (fallback): when the FTS5 table is absent (FTS disabled / fresh DB)
or the MATCH raises, a substring scan over ``messages.content``.
Scope:
* ``"all"`` -- every channel/thread/session.
* ``"chat"`` -- restrict to the given ``chat_id`` (and optional ``thread_id``).
Each hit: ``{message_id, chat_id, thread_id, role, snippet, ts}`` where
``message_id`` is the session-store row id (string) and ``ts`` is epoch
milliseconds. The app navigates to the hit's channel/thread and matches the
message by timestamp to scroll + highlight.
Privacy: search is local to the user's own hermes home; no data leaves the
machine.
Milestone M3.
"""
"""
import logging
import re
import sqlite3
from pathlib import Path
from typing import Any, Dict, List, Optional, Tuple
logger = logging.getLogger(__name__)
MAX_QUERY_CHARS = 200
DEFAULT_LIMIT = 20
MAX_LIMIT = 100
# FTS5 special chars (mirror of hermes_state_search._FTS5_SPECIAL_CHARS) for the
# fallback sanitizer when the real one can't be imported.
_FTS5_SPECIAL_CHARS = '+{}():"^@/#&|~[]<>,;!?$=\\\''
_FTS5_SPECIAL_RE = re.compile(f"[{re.escape(_FTS5_SPECIAL_CHARS)}]")
def _sanitize(query: str) -> str:
"""Sanitize user input for a safe FTS5 MATCH.
Prefers the gateway's own sanitizer (exact parity with session_search);
falls back to a simplified strip-and-quote pass when it can't be imported.
"""
q = (query or "").strip()
if not q:
return ""
try:
from hermes_state_search import SessionSearchMixin
return SessionSearchMixin._sanitize_fts5_query(q)
except Exception:
return _sanitize_fallback(q)
def _sanitize_fallback(query: str) -> str:
q = query[:MAX_QUERY_CHARS]
q = _FTS5_SPECIAL_RE.sub(" ", q)
if "%" in q:
q = q.replace("%", " ")
q = re.sub(r"\*+", "*", q)
q = re.sub(r"(^|\s)\*", r"\1", q)
q = re.sub(r"(?i)^(AND|OR|NOT)\b\s*", "", q.strip())
q = re.sub(r"(?i)\s+(AND|OR|NOT)\s*$", "", q.strip())
q = re.sub(r"\b(\w+(?:[._-]\w+)+)\b", r'"\1"', q)
return q.strip()
def _fts_available(conn: sqlite3.Connection) -> bool:
try:
row = conn.execute(
"SELECT 1 FROM sqlite_master WHERE type = 'table' "
"AND name = 'messages_fts' LIMIT 1"
).fetchone()
return row is not None
except sqlite3.Error:
return False
def _scope_clauses(
scope: str, chat_id: Optional[str], thread_id: Optional[str]
) -> Tuple[List[str], List[Any]]:
"""Build the scope WHERE clauses + params (empty for scope='all')."""
clauses: List[str] = []
params: List[Any] = []
if scope == "chat" and chat_id:
clauses.append("s.chat_id = ?")
params.append(chat_id)
if thread_id:
clauses.append("s.thread_id = ?")
params.append(thread_id)
return clauses, params
def _row_to_hit(row: sqlite3.Row) -> Dict[str, Any]:
ts = row["timestamp"]
try:
ts_ms = int(float(ts) * 1000)
except (TypeError, ValueError):
ts_ms = 0
return {
"message_id": str(row["id"]),
"chat_id": row["chat_id"],
"thread_id": row["thread_id"],
"role": row["role"],
"snippet": row["snippet"] or "",
"ts": ts_ms,
}
def _fts_query(
conn: sqlite3.Connection,
query: str,
scope: str,
chat_id: Optional[str],
thread_id: Optional[str],
limit: int,
) -> List[Dict[str, Any]]:
where = ["messages_fts MATCH ?", "(m.active = 1 OR m.compacted = 1)"]
params: List[Any] = [query]
scope_clauses, scope_params = _scope_clauses(scope, chat_id, thread_id)
where.extend(scope_clauses)
params.extend(scope_params)
params.extend([limit])
sql = f"""
SELECT
m.id,
m.role,
snippet(messages_fts, -1, '>>>', '<<<', '...', 40) AS snippet,
m.timestamp,
s.chat_id,
s.thread_id
FROM messages_fts
JOIN messages m ON m.id = messages_fts.rowid
JOIN sessions s ON s.id = m.session_id
WHERE {' AND '.join(where)}
ORDER BY rank
LIMIT ?
"""
rows = conn.execute(sql, params).fetchall()
return [_row_to_hit(r) for r in rows]
def _like_query(
conn: sqlite3.Connection,
query: str,
scope: str,
chat_id: Optional[str],
thread_id: Optional[str],
limit: int,
) -> List[Dict[str, Any]]:
"""Substring fallback when FTS5 is unavailable."""
# First plain word of the query is the LIKE needle (best-effort).
needle = re.split(r"\s+", query.strip(), maxsplit=1)[0].strip('"')
if not needle:
return []
like = f"%{needle}%"
where = ["(m.active = 1 OR m.compacted = 1)", "m.content LIKE ?"]
params: List[Any] = [like]
scope_clauses, scope_params = _scope_clauses(scope, chat_id, thread_id)
where.extend(scope_clauses)
params.extend(scope_params)
params.extend([limit])
sql = f"""
SELECT
m.id,
m.role,
substr(m.content, max(1, instr(m.content, ?) - 40), 120) AS snippet,
m.timestamp,
s.chat_id,
s.thread_id
FROM messages m
JOIN sessions s ON s.id = m.session_id
WHERE {' AND '.join(where)}
ORDER BY m.timestamp DESC
LIMIT ?
"""
# The needle appears twice (LIKE + instr); params order: like, scope..., needle, limit
full_params = [like, *scope_params, needle, limit]
rows = conn.execute(sql, full_params).fetchall()
return [_row_to_hit(r) for r in rows]
def search(
db_path: Path,
query: str,
scope: str = "all",
chat_id: Optional[str] = None,
thread_id: Optional[str] = None,
limit: int = DEFAULT_LIMIT,
) -> List[Dict[str, Any]]:
"""Run a scoped search over the session store. Returns a list of hits.
Never raises: any DB/FTS error yields an empty result (the caller sends an
empty ``search.results``).
"""
sanitized = _sanitize(query)
if not sanitized:
return []
db_path = Path(db_path)
if not db_path.exists():
return []
limit = max(1, min(int(limit or DEFAULT_LIMIT), MAX_LIMIT))
scope = (scope or "all").strip().lower()
if scope not in ("all", "chat"):
scope = "all"
try:
conn = sqlite3.connect(f"file:{db_path}?mode=ro", uri=True)
except sqlite3.Error as e:
logger.warning("android search: open failed: %s", e)
return []
conn.row_factory = sqlite3.Row
try:
if _fts_available(conn):
try:
return _fts_query(conn, sanitized, scope, chat_id, thread_id, limit)
except sqlite3.Error as e:
logger.debug("android search: FTS5 failed, using LIKE: %s", e)
return _like_query(conn, sanitized, scope, chat_id, thread_id, limit)
except sqlite3.Error as e:
logger.warning("android search: query failed: %s", e)
return []
finally:
try:
conn.close()
except Exception:
pass
+15 -1
View File
@@ -240,7 +240,7 @@ class WsServer:
ack = protocol.hello_ack(
server_caps=self._adapter.server_caps(),
sync_cursor=0, # outbox lands in M3; cursor starts at 0
sync_cursor=self._adapter._outbox.latest_cursor(),
channels=self._adapter.channel_list(),
)
try:
@@ -276,6 +276,20 @@ class WsServer:
await self._send_quiet(ws, protocol.pong(ts if isinstance(ts, int) else None))
elif frame.type == protocol.TYPE_MESSAGE_SEND:
await self._adapter.on_message_send(frame, device_id)
elif frame.type == protocol.TYPE_CHANNEL_CREATE:
await self._adapter.on_channel_create(frame, device_id)
elif frame.type == protocol.TYPE_CHANNEL_RENAME:
await self._adapter.on_channel_rename(frame, device_id)
elif frame.type == protocol.TYPE_CHANNEL_SET_DEFAULT:
await self._adapter.on_channel_set_default(frame, device_id)
elif frame.type == protocol.TYPE_CHANNEL_DELETE:
await self._adapter.on_channel_delete(frame, device_id)
elif frame.type == protocol.TYPE_CHANNEL_LIST:
await self._adapter.on_channel_list(frame, device_id)
elif frame.type == protocol.TYPE_SEARCH:
await self._adapter.on_search(frame, device_id)
elif frame.type == protocol.TYPE_SYNC:
await self._adapter.on_sync(frame, device_id)
elif frame.type == "fcm.register":
fcm_token = frame.payload.get("fcm_token")
ntfy_topic = frame.payload.get("ntfy_topic")