fix(kanban): preserve telegram dm topic metadata

This commit is contained in:
embwl0x 2026-07-08 03:58:44 -04:00 committed by Teknium
parent 48bdde1deb
commit 1bdec6f065
13 changed files with 214 additions and 16 deletions

View file

@ -413,8 +413,13 @@ class GatewayKanbanWatchersMixin:
# internal transition. They are also excluded from
# _WAKE_KINDS below, so they never wake the creator.
continue
metadata: dict[str, Any] = {}
if sub.get("thread_id"):
delivery_metadata = sub.get("delivery_metadata")
metadata: dict[str, Any] = (
dict(delivery_metadata)
if isinstance(delivery_metadata, dict)
else {}
)
if sub.get("thread_id") and not metadata.get("thread_id"):
metadata["thread_id"] = sub["thread_id"]
# Adapters with no push channel (the API server —
# ``supports_async_delivery = False``) can NEVER
@ -638,13 +643,21 @@ class GatewayKanbanWatchersMixin:
# from group/thread, so the old hardcoded
# "group" mis-routed DM/thread creators into a
# fresh session. Legacy rows written before the
# column existed store NULL/'' — fall back to
# "group" for them (the historical default that
# suits the dashboard/group flows).
# column existed may still carry chat_type in
# delivery_metadata (#60600 rows) — fall back
# to that, then to "group" (the historical
# default that suits the dashboard/group flows).
# handle_message() get_or_create_session's the
# target, so a mismatch only ever degrades to a
# fresh session, never an exception.
_chat_type = str(sub.get("chat_type") or "").strip() or "group"
_chat_type = str(sub.get("chat_type") or "").strip()
if not _chat_type:
_delivery_meta = sub.get("delivery_metadata")
if isinstance(_delivery_meta, dict):
_chat_type = str(
_delivery_meta.get("chat_type") or ""
).strip()
_chat_type = _chat_type or "group"
_source = SessionSource(
platform=plat,
chat_id=sub["chat_id"],

View file

@ -17450,6 +17450,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
str(context.source.chat_type) if context.source.chat_type else ""
),
chat_name=context.source.chat_name or "",
chat_type=context.source.chat_type or "",
thread_id=str(context.source.thread_id) if context.source.thread_id else "",
user_id=str(context.source.user_id) if context.source.user_id else "",
user_name=str(context.source.user_name) if context.source.user_name else "",

View file

@ -75,6 +75,7 @@ _SESSION_SOURCE: ContextVar = ContextVar("HERMES_SESSION_SOURCE", default=_UNSET
_SESSION_CHAT_ID: ContextVar = ContextVar("HERMES_SESSION_CHAT_ID", default=_UNSET)
_SESSION_CHAT_TYPE: ContextVar = ContextVar("HERMES_SESSION_CHAT_TYPE", default=_UNSET)
_SESSION_CHAT_NAME: ContextVar = ContextVar("HERMES_SESSION_CHAT_NAME", default=_UNSET)
_SESSION_CHAT_TYPE: ContextVar = ContextVar("HERMES_SESSION_CHAT_TYPE", default=_UNSET)
_SESSION_THREAD_ID: ContextVar = ContextVar("HERMES_SESSION_THREAD_ID", default=_UNSET)
_SESSION_USER_ID: ContextVar = ContextVar("HERMES_SESSION_USER_ID", default=_UNSET)
_SESSION_USER_NAME: ContextVar = ContextVar("HERMES_SESSION_USER_NAME", default=_UNSET)
@ -126,6 +127,7 @@ _VAR_MAP = {
"HERMES_SESSION_CHAT_ID": _SESSION_CHAT_ID,
"HERMES_SESSION_CHAT_TYPE": _SESSION_CHAT_TYPE,
"HERMES_SESSION_CHAT_NAME": _SESSION_CHAT_NAME,
"HERMES_SESSION_CHAT_TYPE": _SESSION_CHAT_TYPE,
"HERMES_SESSION_THREAD_ID": _SESSION_THREAD_ID,
"HERMES_SESSION_USER_ID": _SESSION_USER_ID,
"HERMES_SESSION_USER_NAME": _SESSION_USER_NAME,
@ -161,6 +163,7 @@ def set_session_vars(
chat_id: str = "",
chat_type: str = "",
chat_name: str = "",
chat_type: str = "",
thread_id: str = "",
user_id: str = "",
user_name: str = "",
@ -198,6 +201,7 @@ def set_session_vars(
_SESSION_CHAT_ID.set(chat_id),
_SESSION_CHAT_TYPE.set(chat_type),
_SESSION_CHAT_NAME.set(chat_name),
_SESSION_CHAT_TYPE.set(chat_type),
_SESSION_THREAD_ID.set(thread_id),
_SESSION_USER_ID.set(user_id),
_SESSION_USER_NAME.set(user_name),
@ -234,6 +238,7 @@ def clear_session_vars(tokens: list) -> None:
_SESSION_CHAT_ID,
_SESSION_CHAT_TYPE,
_SESSION_CHAT_NAME,
_SESSION_CHAT_TYPE,
_SESSION_THREAD_ID,
_SESSION_USER_ID,
_SESSION_USER_NAME,

View file

@ -481,6 +481,13 @@ class GatewaySlashCommandsMixin:
chat_type = str(getattr(source, "chat_type", "") or "") or None
thread_id = str(getattr(source, "thread_id", "") or "")
user_id = str(getattr(source, "user_id", "") or "") or None
delivery_metadata = self._thread_metadata_for_source(
source, self._reply_anchor_for_event(event)
) or None
if isinstance(delivery_metadata, dict):
chat_type = str(getattr(source, "chat_type", "") or "")
if chat_type:
delivery_metadata.setdefault("chat_type", chat_type)
if platform_str and chat_id:
def _sub():
from hermes_cli import kanban_db as _kb
@ -493,6 +500,7 @@ class GatewaySlashCommandsMixin:
thread_id=thread_id or None,
user_id=user_id,
notifier_profile=getattr(self, "_kanban_notifier_profile", None) or self._active_profile_name(),
delivery_metadata=delivery_metadata,
)
finally:
conn.close()

View file

@ -87,7 +87,7 @@ import time
from contextvars import ContextVar, Token
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any, Iterable, Optional
from typing import Any, Iterable, Mapping, Optional
from hermes_cli.sqlite_util import add_column_if_missing as _add_column_if_missing
from toolsets import get_toolset_names
@ -1304,6 +1304,7 @@ CREATE TABLE IF NOT EXISTS kanban_notify_subs (
thread_id TEXT NOT NULL DEFAULT '',
user_id TEXT,
notifier_profile TEXT,
delivery_metadata TEXT,
created_at INTEGER NOT NULL,
last_event_id INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (task_id, platform, chat_id, thread_id)
@ -2450,6 +2451,10 @@ def _migrate_add_optional_columns(conn: sqlite3.Connection) -> None:
_add_column_if_missing(
conn, "kanban_notify_subs", "chat_type", "chat_type TEXT"
)
if "delivery_metadata" not in notify_cols:
_add_column_if_missing(
conn, "kanban_notify_subs", "delivery_metadata", "delivery_metadata TEXT"
)
# One-shot backfill: any task that is 'running' before runs existed
# had its claim_lock / claim_expires / worker_pid on the task row.
@ -2573,7 +2578,7 @@ _REBUILD_SPECS = {
"CREATE TABLE kanban_notify_subs ("
" task_id TEXT NOT NULL, platform TEXT NOT NULL, chat_id TEXT NOT NULL,"
" chat_type TEXT, thread_id TEXT NOT NULL DEFAULT '', user_id TEXT,"
" notifier_profile TEXT, created_at INTEGER NOT NULL,"
" notifier_profile TEXT, delivery_metadata TEXT, created_at INTEGER NOT NULL,"
" last_event_id INTEGER NOT NULL DEFAULT 0,"
" PRIMARY KEY (task_id, platform, chat_id, thread_id))",
("CREATE INDEX idx_notify_task ON kanban_notify_subs(task_id)",),
@ -9378,6 +9383,41 @@ def task_age(task: Task) -> dict:
# Notification subscriptions (used by the gateway kanban-notifier)
# ---------------------------------------------------------------------------
def _encode_notify_delivery_metadata(
metadata: Optional[Mapping[str, Any]],
) -> Optional[str]:
"""Serialize platform send metadata stored on notification subscriptions."""
if not isinstance(metadata, Mapping):
return None
clean: dict[str, Any] = {}
for key, value in metadata.items():
if value is None:
continue
if isinstance(value, (str, int, float, bool)):
clean[str(key)] = value
if not clean:
return None
return json.dumps(clean, sort_keys=True, separators=(",", ":"))
def _decode_notify_delivery_metadata(raw: Any) -> dict[str, Any]:
if isinstance(raw, Mapping):
return dict(raw)
if not raw:
return {}
try:
data = json.loads(str(raw))
except Exception:
return {}
if not isinstance(data, dict):
return {}
return {
str(key): value
for key, value in data.items()
if isinstance(value, (str, int, float, bool))
}
def add_notify_sub(
conn: sqlite3.Connection,
*,
@ -9388,16 +9428,19 @@ def add_notify_sub(
thread_id: Optional[str] = None,
user_id: Optional[str] = None,
notifier_profile: Optional[str] = None,
delivery_metadata: Optional[Mapping[str, Any]] = None,
) -> None:
"""Register a gateway source that wants terminal-state notifications
for ``task_id``. Idempotent on (task, platform, chat, thread)."""
now = int(time.time())
metadata_json = _encode_notify_delivery_metadata(delivery_metadata)
with write_txn(conn):
conn.execute(
"""
INSERT OR IGNORE INTO kanban_notify_subs
(task_id, platform, chat_id, chat_type, thread_id, user_id, notifier_profile, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
(task_id, platform, chat_id, chat_type, thread_id, user_id,
notifier_profile, delivery_metadata, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
task_id,
@ -9407,6 +9450,7 @@ def add_notify_sub(
thread_id or "",
user_id,
notifier_profile,
metadata_json,
now,
),
)
@ -9433,6 +9477,18 @@ def add_notify_sub(
""",
(notifier_profile, task_id, platform, chat_id, thread_id or ""),
)
if metadata_json:
# A duplicate subscribe from the same chat/thread should refresh
# the routing anchor. Telegram DM-topic notifications need the
# latest reply anchor to stay inside the visible topic lane.
conn.execute(
"""
UPDATE kanban_notify_subs
SET delivery_metadata = ?
WHERE task_id = ? AND platform = ? AND chat_id = ? AND thread_id = ?
""",
(metadata_json, task_id, platform, chat_id, thread_id or ""),
)
def list_notify_subs(
@ -9444,7 +9500,15 @@ def list_notify_subs(
).fetchall()
else:
rows = conn.execute("SELECT * FROM kanban_notify_subs").fetchall()
return [dict(r) for r in rows]
out: list[dict] = []
for row in rows:
item = dict(row)
if "delivery_metadata" in item:
item["delivery_metadata"] = _decode_notify_delivery_metadata(
item.get("delivery_metadata")
)
out.append(item)
return out
def remove_notify_sub(

View file

@ -179,6 +179,7 @@ _HERMES_BEHAVIORAL_VARS = frozenset({
"HERMES_SESSION_PLATFORM",
"HERMES_SESSION_CHAT_ID",
"HERMES_SESSION_CHAT_NAME",
"HERMES_SESSION_CHAT_TYPE",
"HERMES_SESSION_THREAD_ID",
"HERMES_SESSION_SOURCE",
"HERMES_SESSION_KEY",

View file

@ -110,6 +110,54 @@ def test_kanban_notifier_claim_prevents_second_watcher_send(tmp_path, monkeypatc
assert adapter2.sent == []
def test_kanban_notifier_replays_telegram_dm_topic_delivery_metadata(tmp_path, monkeypatch):
db_path = tmp_path / "dm-topic-metadata.db"
monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))
kb.init_db()
conn = kb.connect()
try:
tid = kb.create_task(
conn,
title="dm topic task",
assignee="worker",
session_id="agent:main:telegram:dm:chat-1",
)
kb.add_notify_sub(
conn,
task_id=tid,
platform="telegram",
chat_id="chat-1",
thread_id="20197",
delivery_metadata={
"chat_type": "dm",
"direct_messages_topic_id": "20197",
"telegram_dm_topic_reply_fallback": True,
"telegram_reply_to_message_id": "462",
"thread_id": "20197",
},
)
kb.complete_task(conn, tid, summary="done")
finally:
conn.close()
adapter = RecordingAdapter()
runner = _make_runner(adapter)
asyncio.run(_run_one_notifier_tick(monkeypatch, runner))
assert len(adapter.sent) == 1
assert adapter.sent[0]["metadata"] == {
"chat_type": "dm",
"direct_messages_topic_id": "20197",
"telegram_dm_topic_reply_fallback": True,
"telegram_reply_to_message_id": "462",
"thread_id": "20197",
}
assert len(adapter.handled) == 1
assert adapter.handled[0].source.chat_type == "dm"
assert adapter.handled[0].source.thread_id == "20197"
def test_kanban_notifier_rewinds_claim_if_adapter_disconnects(tmp_path, monkeypatch):
db_path = tmp_path / "adapter-disconnect.db"
monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))

View file

@ -48,6 +48,7 @@ def test_set_session_env_sets_contextvars(monkeypatch):
monkeypatch.delenv("HERMES_SESSION_SOURCE", raising=False)
monkeypatch.delenv("HERMES_SESSION_CHAT_ID", raising=False)
monkeypatch.delenv("HERMES_SESSION_CHAT_NAME", raising=False)
monkeypatch.delenv("HERMES_SESSION_CHAT_TYPE", raising=False)
monkeypatch.delenv("HERMES_SESSION_USER_ID", raising=False)
monkeypatch.delenv("HERMES_SESSION_USER_NAME", raising=False)
monkeypatch.delenv("HERMES_SESSION_THREAD_ID", raising=False)
@ -59,6 +60,7 @@ def test_set_session_env_sets_contextvars(monkeypatch):
assert get_session_env("HERMES_SESSION_SOURCE") == ""
assert get_session_env("HERMES_SESSION_CHAT_ID") == "-1001"
assert get_session_env("HERMES_SESSION_CHAT_NAME") == "Group"
assert get_session_env("HERMES_SESSION_CHAT_TYPE") == "group"
assert get_session_env("HERMES_SESSION_USER_ID") == "123456"
assert get_session_env("HERMES_SESSION_USER_NAME") == "alice"
assert get_session_env("HERMES_SESSION_THREAD_ID") == "17585"
@ -66,6 +68,7 @@ def test_set_session_env_sets_contextvars(monkeypatch):
# os.environ should NOT be touched
assert os.getenv("HERMES_SESSION_PLATFORM") is None
assert os.getenv("HERMES_SESSION_SOURCE") is None
assert os.getenv("HERMES_SESSION_CHAT_TYPE") is None
assert os.getenv("HERMES_SESSION_THREAD_ID") is None
# Clean up
@ -91,6 +94,7 @@ def test_clear_session_env_restores_previous_state(monkeypatch):
monkeypatch.delenv("HERMES_SESSION_PLATFORM", raising=False)
monkeypatch.delenv("HERMES_SESSION_CHAT_ID", raising=False)
monkeypatch.delenv("HERMES_SESSION_CHAT_NAME", raising=False)
monkeypatch.delenv("HERMES_SESSION_CHAT_TYPE", raising=False)
monkeypatch.delenv("HERMES_SESSION_USER_ID", raising=False)
monkeypatch.delenv("HERMES_SESSION_USER_NAME", raising=False)
monkeypatch.delenv("HERMES_SESSION_THREAD_ID", raising=False)
@ -109,6 +113,7 @@ def test_clear_session_env_restores_previous_state(monkeypatch):
tokens = runner._set_session_env(context)
assert get_session_env("HERMES_SESSION_PLATFORM") == "telegram"
assert get_session_env("HERMES_SESSION_USER_ID") == "123456"
assert get_session_env("HERMES_SESSION_CHAT_TYPE") == "group"
runner._clear_session_env(tokens)
@ -116,6 +121,7 @@ def test_clear_session_env_restores_previous_state(monkeypatch):
assert get_session_env("HERMES_SESSION_PLATFORM") == ""
assert get_session_env("HERMES_SESSION_CHAT_ID") == ""
assert get_session_env("HERMES_SESSION_CHAT_NAME") == ""
assert get_session_env("HERMES_SESSION_CHAT_TYPE") == ""
assert get_session_env("HERMES_SESSION_USER_ID") == ""
assert get_session_env("HERMES_SESSION_USER_NAME") == ""
assert get_session_env("HERMES_SESSION_THREAD_ID") == ""
@ -393,4 +399,3 @@ async def test_gateway_executor_refuses_resurrection_after_shutdown():
await runner._run_in_executor_with_context(lambda: "second")
finally:
runner._shutdown_executor()

View file

@ -522,16 +522,31 @@ def test_notify_sub_crud(kanban_home):
kb.add_notify_sub(
conn, task_id=tid, platform="telegram", chat_id="123", user_id="u1",
notifier_profile="default",
delivery_metadata={
"chat_type": "dm",
"telegram_reply_to_message_id": "42",
},
)
subs = kb.list_notify_subs(conn, tid)
assert len(subs) == 1
assert subs[0]["platform"] == "telegram"
assert subs[0]["notifier_profile"] == "default"
assert subs[0]["delivery_metadata"] == {
"chat_type": "dm",
"telegram_reply_to_message_id": "42",
}
# Duplicate add is a no-op.
kb.add_notify_sub(
conn, task_id=tid, platform="telegram", chat_id="123",
delivery_metadata={
"chat_type": "dm",
"telegram_reply_to_message_id": "43",
},
)
assert len(kb.list_notify_subs(conn, tid)) == 1
assert kb.list_notify_subs(conn, tid)[0]["delivery_metadata"][
"telegram_reply_to_message_id"
] == "43"
# Distinct thread is a new row.
kb.add_notify_sub(
conn, task_id=tid, platform="telegram", chat_id="123",

View file

@ -115,6 +115,7 @@ def test_legacy_text_pk_tables_rebuilt_to_integer_autoincrement(tmp_path, monkey
lei = {r["name"]: r for r in conn.execute("PRAGMA table_info(kanban_notify_subs)")}
assert lei["last_event_id"]["type"].upper() == "INTEGER"
assert "delivery_metadata" in lei
# Data preserved across the rebuild.
assert len(conn.execute("SELECT * FROM task_events").fetchall()) == 2

View file

@ -561,12 +561,15 @@ async def test_gateway_create_autosubscribes_on_explicit_board(kanban_home):
source = SimpleNamespace(
platform=Platform.TELEGRAM,
chat_id="chat1",
thread_id="th1",
chat_type="dm",
thread_id="20197",
user_id="u1",
)
event = SimpleNamespace(
text='/kanban --board projx create "hello" --assignee alice',
source=source,
message_id="462",
reply_to_message_id=None,
)
out = await GatewayRunner._handle_kanban_command(runner, event)
@ -583,7 +586,14 @@ async def test_gateway_create_autosubscribes_on_explicit_board(kanban_home):
assert [t.title for t in tasks] == ["hello"]
assert len(subs) == 1
assert subs[0]["chat_id"] == "chat1"
assert subs[0]["thread_id"] == "th1"
assert subs[0]["thread_id"] == "20197"
assert subs[0]["delivery_metadata"] == {
"chat_type": "dm",
"direct_messages_topic_id": "20197",
"telegram_dm_topic_reply_fallback": True,
"telegram_reply_to_message_id": "462",
"thread_id": "20197",
}
conn = kb.connect(board="default")
try:

View file

@ -2415,6 +2415,7 @@ def _sub_index(subs):
"chat_id": getattr(s, "chat_id", None),
"thread_id": getattr(s, "thread_id", None),
"user_id": getattr(s, "user_id", None),
"delivery_metadata": getattr(s, "delivery_metadata", None),
})
return out
@ -2426,8 +2427,10 @@ def test_create_subscribes_gateway_session(monkeypatch, worker_env):
from tools import kanban_tools as kt
monkeypatch.setenv("HERMES_SESSION_PLATFORM", "telegram")
monkeypatch.setenv("HERMES_SESSION_CHAT_ID", "chat-42")
monkeypatch.setenv("HERMES_SESSION_THREAD_ID", "thread-7")
monkeypatch.setenv("HERMES_SESSION_CHAT_TYPE", "dm")
monkeypatch.setenv("HERMES_SESSION_THREAD_ID", "20197")
monkeypatch.setenv("HERMES_SESSION_USER_ID", "user-9")
monkeypatch.setenv("HERMES_SESSION_MESSAGE_ID", "msg-11")
out = kt._handle_create({
"title": "auto-sub gateway",
@ -2443,8 +2446,15 @@ def test_create_subscribes_gateway_session(monkeypatch, worker_env):
s = subs[0]
assert s["platform"] == "telegram"
assert s["chat_id"] == "chat-42"
assert s["thread_id"] == "thread-7"
assert s["thread_id"] == "20197"
assert s["user_id"] == "user-9"
assert s["delivery_metadata"] == {
"chat_type": "dm",
"direct_messages_topic_id": "20197",
"telegram_dm_topic_reply_fallback": True,
"telegram_reply_to_message_id": "msg-11",
"thread_id": "20197",
}
def test_create_subscribes_tui_session_via_session_key(monkeypatch, worker_env):

View file

@ -1334,10 +1334,26 @@ def _maybe_auto_subscribe(conn: Any, task_id: str) -> bool:
thread_id = get_session_env("HERMES_SESSION_THREAD_ID", "") or None
user_id = get_session_env("HERMES_SESSION_USER_ID", "") or None
chat_type = get_session_env("HERMES_SESSION_CHAT_TYPE", "") or None
message_id = get_session_env("HERMES_SESSION_MESSAGE_ID", "") or ""
notifier_profile = (
get_session_env("HERMES_SESSION_PROFILE", "")
or os.environ.get("HERMES_PROFILE")
)
delivery_metadata: dict[str, Any] = {}
if thread_id:
delivery_metadata["thread_id"] = thread_id
if chat_type:
delivery_metadata["chat_type"] = chat_type
if (
platform.lower() == "telegram"
and thread_id
and (chat_type or "").lower() in {"dm", "direct", "private"}
):
delivery_metadata["telegram_dm_topic_reply_fallback"] = True
if str(thread_id) not in {"", "1"}:
delivery_metadata["direct_messages_topic_id"] = str(thread_id)
if message_id:
delivery_metadata["telegram_reply_to_message_id"] = str(message_id)
# Lazy-import to keep the module-level dependency light
from hermes_cli import kanban_db as _kb
@ -1347,6 +1363,7 @@ def _maybe_auto_subscribe(conn: Any, task_id: str) -> bool:
chat_type=chat_type,
thread_id=thread_id, user_id=user_id,
notifier_profile=notifier_profile,
delivery_metadata=delivery_metadata or None,
)
return True
except Exception as _exc: