diff --git a/gateway/kanban_watchers.py b/gateway/kanban_watchers.py index 4f53c0e1f8f..eef928650ef 100644 --- a/gateway/kanban_watchers.py +++ b/gateway/kanban_watchers.py @@ -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"], diff --git a/gateway/run.py b/gateway/run.py index 02677455d58..5030617d986 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -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 "", diff --git a/gateway/session_context.py b/gateway/session_context.py index cf227fac121..41665ec9574 100644 --- a/gateway/session_context.py +++ b/gateway/session_context.py @@ -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, diff --git a/gateway/slash_commands.py b/gateway/slash_commands.py index 65af38ab7be..3bd184f8c72 100644 --- a/gateway/slash_commands.py +++ b/gateway/slash_commands.py @@ -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() diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 4fcc4fd610e..3877a7f10b4 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -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( diff --git a/tests/conftest.py b/tests/conftest.py index 159afd4c2b4..2ae674cfa40 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -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", diff --git a/tests/gateway/test_kanban_notifier.py b/tests/gateway/test_kanban_notifier.py index 4f3d771722d..b5a51d113a6 100644 --- a/tests/gateway/test_kanban_notifier.py +++ b/tests/gateway/test_kanban_notifier.py @@ -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)) diff --git a/tests/gateway/test_session_env.py b/tests/gateway/test_session_env.py index f5392ab2c22..1183e755df3 100644 --- a/tests/gateway/test_session_env.py +++ b/tests/gateway/test_session_env.py @@ -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() - diff --git a/tests/hermes_cli/test_kanban_core_functionality.py b/tests/hermes_cli/test_kanban_core_functionality.py index 0e898999c1d..5bb66e0d4ec 100644 --- a/tests/hermes_cli/test_kanban_core_functionality.py +++ b/tests/hermes_cli/test_kanban_core_functionality.py @@ -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", diff --git a/tests/hermes_cli/test_kanban_db_init.py b/tests/hermes_cli/test_kanban_db_init.py index 7db5d2009e6..643c55ec3f0 100644 --- a/tests/hermes_cli/test_kanban_db_init.py +++ b/tests/hermes_cli/test_kanban_db_init.py @@ -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 diff --git a/tests/hermes_cli/test_kanban_notify.py b/tests/hermes_cli/test_kanban_notify.py index 15eb2671b24..37acc5ae8ff 100644 --- a/tests/hermes_cli/test_kanban_notify.py +++ b/tests/hermes_cli/test_kanban_notify.py @@ -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: diff --git a/tests/tools/test_kanban_tools.py b/tests/tools/test_kanban_tools.py index 38ba049eee9..20deaeae9e4 100644 --- a/tests/tools/test_kanban_tools.py +++ b/tests/tools/test_kanban_tools.py @@ -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): diff --git a/tools/kanban_tools.py b/tools/kanban_tools.py index acf3442150c..c781ff88457 100644 --- a/tools/kanban_tools.py +++ b/tools/kanban_tools.py @@ -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: