Files
iris_x_hermes/gateway-plugin/channels.py
T

438 lines
17 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,
favorite INTEGER NOT NULL DEFAULT 0,
icon TEXT,
color TEXT,
automation INTEGER NOT NULL DEFAULT 0
)
"""
)
# Migrate existing DBs: add the cosmetic columns if missing.
existing = {
row[1] for row in self._conn.execute("PRAGMA table_info(channels)")
}
if "favorite" not in existing:
self._conn.execute(
"ALTER TABLE channels ADD COLUMN favorite INTEGER NOT NULL DEFAULT 0"
)
if "icon" not in existing:
self._conn.execute("ALTER TABLE channels ADD COLUMN icon TEXT")
if "color" not in existing:
self._conn.execute("ALTER TABLE channels ADD COLUMN color TEXT")
if "automation" not in existing:
self._conn.execute(
"ALTER TABLE channels ADD COLUMN automation INTEGER 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). The default
# channel is the user's chat surface, so it is never an
# automation channel.
self._conn.execute(
"UPDATE channels SET kind = ?, is_default = 1, automation = 0 "
"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).
The default channel is the user's chat surface, so the automation
flag is cleared on it (a read-only home channel would be unusable).
"""
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, automation = 0 "
"WHERE chat_id = ?",
(chat_id,),
)
self._conn.commit()
return self.get(chat_id)
def set_favorite(self, chat_id: str, on: bool) -> Optional[Dict[str, Any]]:
"""Toggle the cosmetic favorite flag (sorts to the top of the list)."""
with self._lock:
cur = self._conn.execute(
"UPDATE channels SET favorite = ? WHERE chat_id = ? AND archived = 0",
(1 if on else 0, chat_id),
)
self._conn.commit()
if cur.rowcount == 0:
return None
return self.get(chat_id)
def set_icon(self, chat_id: str, icon: Optional[str], color: Optional[str]) -> Optional[Dict[str, Any]]:
"""Set the channel's cosmetic icon (base64 image) and/or avatar color.
``icon`` is a base64-encoded image (or ``None`` to clear it); ``color``
is a ``#RRGGBB`` hex string (or ``None`` to clear it). Both are purely
cosmetic and independent of the name / default flag.
"""
with self._lock:
cur = self._conn.execute(
"UPDATE channels SET icon = ?, color = ? WHERE chat_id = ? AND archived = 0",
(icon, color, chat_id),
)
self._conn.commit()
if cur.rowcount == 0:
return None
return self.get(chat_id)
def set_automation(self, chat_id: str, on: bool) -> Optional[Dict[str, Any]]:
"""Mark *chat_id* as an automation channel (or clear the flag).
Automation channels are read-only for the user: they only receive
gateway-originated output (cron jobs, webhooks). The app hides the
composer and the gateway rejects ``message.send`` into them. The
default channel cannot be marked automation (it is the user's chat
surface), so this returns ``None`` for it, like ``delete``.
"""
with self._lock:
row = self._conn.execute(
"SELECT is_default FROM channels WHERE chat_id = ? AND archived = 0",
(chat_id,),
).fetchone()
if row is None or row["is_default"]:
return None
self._conn.execute(
"UPDATE channels SET automation = ? WHERE chat_id = ?",
(1 if on else 0, 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 first, then favorites, then creation order."""
sql = "SELECT * FROM channels"
if not include_archived:
sql += " WHERE archived = 0"
sql += " ORDER BY is_default DESC, favorite 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"],
"favorite": bool(row["favorite"]),
"icon": row["icon"],
"color": row["color"],
"automation": bool(row["automation"]),
}
# ---------------------------------------------------------------------------
# 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