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
353 lines
14 KiB
Python
353 lines
14 KiB
Python
"""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 |