diff --git a/cron/scheduler.py b/cron/scheduler.py index 51dcd3a35a..c583d289ed 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -1297,6 +1297,30 @@ def _iter_home_target_platforms(): pass +def _relay_fronted_delivery_platforms(connected: set) -> set: + """Logical platforms deliverable through a connected relay connector. + + ``get_connected_platforms()`` only sees NATIVELY configured platforms. + On a relay-fronted deployment (relay in ``config.platforms``, the real + platform credential living in the connector) the fronted platforms are + absent from that set although fire-time routing delivers to them via + ``resolve_delivery_transport`` + ``RelayAdapter.fronts_platform``. This + keeps validation symmetric with routing by consulting the same + env-derived deploy stamp (``GATEWAY_RELAY_PLATFORMS``) the live + adapter's identity set is seeded from. No relay connected -> empty set, + so native topologies keep the strict credential check unchanged. + """ + if "relay" not in connected: + return set() + try: + from gateway.relay import relay_fronted_platforms + + return relay_fronted_platforms() + except Exception: + logger.debug("relay fronted-platform lookup failed", exc_info=True) + return set() + + def cron_delivery_targets() -> list[dict]: """Return the platforms a cron job can auto-deliver to. @@ -1317,6 +1341,7 @@ def cron_delivery_targets() -> list[dict]: gateway_config = load_gateway_config() connected = {p.value for p in gateway_config.get_connected_platforms()} + connected |= _relay_fronted_delivery_platforms(connected) except Exception: logger.debug("cron_delivery_targets: gateway config unavailable", exc_info=True) connected = set() @@ -1338,6 +1363,36 @@ def cron_delivery_targets() -> list[dict]: return targets +def _origin_thread_is_stale(origin: dict) -> bool: + """True when a Slack origin's thread is a stale creation-turn artifact. + + Relay-fronted Slack in thread-per-message mode stamps each top-level + message's own id as the session thread (a session KEY, not a durable + location). Jobs persisted before origin capture learned to drop that + stamp carry it as ``origin.thread_id`` forever. Heuristic that repairs + them at fire time without touching genuine threads: when the origin + chat IS the configured Slack home chat (the ``/sethome`` conversation), + a pinned origin thread is the creation-message artifact — the user's + delivery expectation for their home conversation is top-level (or the + home target's own configured thread). Non-home chats keep their + threads: a job deliberately created inside a working thread stays there. + """ + if str(origin.get("platform") or "").lower() != "slack": + return False + if not origin.get("thread_id"): + return False + home_chat = _get_home_target_chat_id("slack") + return bool(home_chat) and str(origin.get("chat_id")) == str(home_chat) + + +def _origin_delivery_thread(origin: dict): + """The thread a deliver=origin job should use, stale stamps dropped.""" + if _origin_thread_is_stale(origin): + home_thread = _get_home_target_thread_id("slack") + return home_thread if home_thread else None + return origin.get("thread_id") + + def _resolve_single_delivery_target(job: dict, deliver_value: str) -> Optional[dict]: """Resolve one concrete auto-delivery target for a cron job.""" @@ -1351,7 +1406,7 @@ def _resolve_single_delivery_target(job: dict, deliver_value: str) -> Optional[d return { "platform": origin["platform"], "chat_id": str(origin["chat_id"]), - "thread_id": origin.get("thread_id"), + "thread_id": _origin_delivery_thread(origin), } # Origin missing (e.g. job created via API/script) — try each # platform's home channel as a fallback instead of silently dropping. @@ -1398,6 +1453,7 @@ def _resolve_single_delivery_target(job: dict, deliver_value: str) -> Optional[d and str(origin.get("platform") or "").lower() == platform_key and str(origin.get("chat_id")) == str(chat_id) and origin.get("thread_id") + and not _origin_thread_is_stale(origin) ): thread_id = origin.get("thread_id") @@ -3128,6 +3184,7 @@ def _preflight_check_delivery(job: dict) -> Optional[str]: connected = { p.value for p in gateway_config.get_connected_platforms() } + connected |= _relay_fronted_delivery_platforms(connected) except Exception: logger.debug( "preflight: gateway config unavailable — skipping " diff --git a/gateway/relay/__init__.py b/gateway/relay/__init__.py index 1e6c25bd75..6c3050696a 100644 --- a/gateway/relay/__init__.py +++ b/gateway/relay/__init__.py @@ -91,6 +91,21 @@ def relay_platform_identities() -> list[tuple[str, str]]: return out +def relay_fronted_platforms() -> set[str]: + """The logical platform names the relay connector fronts for this gateway. + + Thin, env-derived wrapper over :func:`relay_platform_identities` (the + ``GATEWAY_RELAY_PLATFORMS`` deploy stamp) minus the generic ``relay`` + fallback. This is the SAME source ``ws_transport`` seeds the live + adapter's identity set from (``RelayAdapter.fronts_platform``), so + config-time validation (e.g. cron delivery preflight) and fire-time + routing can never disagree — and it needs no live adapter handle, so a + standalone scheduler process can consult it too. Empty set when the + relay fronts nothing. + """ + return {p for p, _ in relay_platform_identities() if p != "relay"} + + def _relay_bot_ids_map() -> dict: """Parse ``GATEWAY_RELAY_BOT_IDS`` (JSON keyed map). Never raises — a malformed map yields ``{}`` so a bad config degrades to empty bot ids (the connector diff --git a/tests/cron/test_cron_origin_synthetic_thread.py b/tests/cron/test_cron_origin_synthetic_thread.py new file mode 100644 index 0000000000..31733c44c1 --- /dev/null +++ b/tests/cron/test_cron_origin_synthetic_thread.py @@ -0,0 +1,82 @@ +"""Cron origin capture: Slack per-message session-key threads are not routing. + +Bug report (relay-fronted Slack, thread-per-message mode): creating a cron job +from a top-level Slack DM message persisted the creation message's own id as +``origin.thread_id`` — the relay adapter stamps ``source.thread_id = message_id`` +on every top-level Slack message purely for SESSION KEYING (native SlackAdapter +parity: ``thread_ts = event.thread_ts or ts``). Every subsequent cron delivery +then landed inside the ephemeral thread spawned around the creation message +instead of the top-level conversation / configured home. + +The stamp is recognizable at capture time: a Slack session whose thread id +equals the triggering message's own id is a synthetic per-message key, not a +durable thread. A genuine in-thread creation has thread_id == the parent +thread's id != the triggering message's own id, and must keep its thread. +""" + +from unittest.mock import patch + +from tools.cronjob_tools import _origin_from_env + + +def _session_env(env: dict): + """Patch gateway.session_context.get_session_env with a dict lookup.""" + return patch( + "gateway.session_context.get_session_env", + side_effect=lambda name, default="": env.get(name, default), + ) + + +class TestSlackSyntheticThreadCapture: + def test_synthetic_slack_thread_not_captured(self): + """thread_id == message_id on Slack = per-message session key: drop it.""" + env = { + "HERMES_SESSION_PLATFORM": "slack", + "HERMES_SESSION_CHAT_ID": "D0BJTDCSR7C", + "HERMES_SESSION_THREAD_ID": "1755043010.123456", + "HERMES_SESSION_MESSAGE_ID": "1755043010.123456", + } + with _session_env(env): + origin = _origin_from_env() + assert origin is not None + assert origin["platform"] == "slack" + assert origin["chat_id"] == "D0BJTDCSR7C" + assert origin["thread_id"] is None + + def test_genuine_slack_thread_preserved(self): + """A real in-thread creation (thread != own message id) keeps its thread.""" + env = { + "HERMES_SESSION_PLATFORM": "slack", + "HERMES_SESSION_CHAT_ID": "C0AGENERAL", + "HERMES_SESSION_THREAD_ID": "1755040000.000100", + "HERMES_SESSION_MESSAGE_ID": "1755043010.123456", + } + with _session_env(env): + origin = _origin_from_env() + assert origin is not None + assert origin["thread_id"] == "1755040000.000100" + + def test_non_slack_platform_thread_untouched(self): + """Telegram forum topics legitimately reuse ids; the rule is Slack-scoped.""" + env = { + "HERMES_SESSION_PLATFORM": "telegram", + "HERMES_SESSION_CHAT_ID": "-1003941067111", + "HERMES_SESSION_THREAD_ID": "2203", + "HERMES_SESSION_MESSAGE_ID": "2203", + } + with _session_env(env): + origin = _origin_from_env() + assert origin is not None + assert origin["thread_id"] == "2203" + + def test_slack_no_message_id_keeps_thread(self): + """Without a message id to compare, never guess: keep the thread.""" + env = { + "HERMES_SESSION_PLATFORM": "slack", + "HERMES_SESSION_CHAT_ID": "D0BJTDCSR7C", + "HERMES_SESSION_THREAD_ID": "1755040000.000100", + } + with _session_env(env): + origin = _origin_from_env() + assert origin is not None + assert origin["thread_id"] == "1755040000.000100" diff --git a/tests/cron/test_cron_relay_delivery_guards.py b/tests/cron/test_cron_relay_delivery_guards.py new file mode 100644 index 0000000000..71a062195c --- /dev/null +++ b/tests/cron/test_cron_relay_delivery_guards.py @@ -0,0 +1,148 @@ +"""Fire-time guards: stale Slack creation-thread routing + relay-fronted preflight. + +Two related defects on relay-fronted Slack deployments: + +1. Jobs persisted before the synthetic-thread capture fix carry the creation + message's own id as ``origin.thread_id``. At fire time ``deliver=origin`` + replayed it unconditionally, and the Slack origin-affinity re-attach put it + back even on explicit ``slack:`` targets. Guard: when the resolved + Slack chat IS the configured home chat, the origin thread is a stale + per-message artifact — deliver top-level (home thread config still wins). + +2. ``_preflight_check_delivery`` validated the ``slack:`` prefix against + natively-configured platforms only; in relay-only topology that set is + ``{relay}`` and the job was refused with "no gateway credentials configured" + although fire-time routing (resolve_delivery_transport + fronts_platform) + would have delivered it. Preflight must consult the relay's fronted set. +""" + +from unittest.mock import MagicMock, patch + +import pytest + +from cron import scheduler as sched +from cron.scheduler import ( + _preflight_check_delivery, + _resolve_single_delivery_target, + cron_delivery_targets, +) + + +def _slack_home(monkeypatch, chat_id="D0BJTDCSR7C", thread_id=None): + monkeypatch.setattr(sched, "_get_home_target_chat_id", + lambda p: chat_id if p == "slack" else None) + monkeypatch.setattr(sched, "_get_home_target_thread_id", + lambda p: thread_id if p == "slack" else None) + + +SYNTH = "1755043010.123456" + + +class TestOriginThreadStaleGuard: + def test_origin_thread_dropped_when_chat_is_home(self, monkeypatch): + """deliver=origin, slack origin chat == home chat: creation thread is stale.""" + _slack_home(monkeypatch) + job = {"origin": {"platform": "slack", "chat_id": "D0BJTDCSR7C", + "thread_id": SYNTH}} + target = _resolve_single_delivery_target(job, "origin") + assert target == {"platform": "slack", "chat_id": "D0BJTDCSR7C", + "thread_id": None} + + def test_origin_thread_kept_when_chat_not_home(self, monkeypatch): + """A non-home Slack origin thread may be a genuine working thread: keep it.""" + _slack_home(monkeypatch, chat_id="D_OTHER_HOME") + job = {"origin": {"platform": "slack", "chat_id": "C0AGENERAL", + "thread_id": "1755040000.000100"}} + target = _resolve_single_delivery_target(job, "origin") + assert target["thread_id"] == "1755040000.000100" + + def test_home_thread_config_still_wins(self, monkeypatch): + """When the home target itself pins a thread, deliver there, not top-level.""" + _slack_home(monkeypatch, thread_id="1755000000.000001") + job = {"origin": {"platform": "slack", "chat_id": "D0BJTDCSR7C", + "thread_id": SYNTH}} + target = _resolve_single_delivery_target(job, "origin") + assert target["thread_id"] == "1755000000.000001" + + def test_non_slack_origin_thread_untouched(self, monkeypatch): + """Telegram forum-topic origins replay their thread verbatim.""" + _slack_home(monkeypatch) + job = {"origin": {"platform": "telegram", "chat_id": "-1003941067111", + "thread_id": "2203"}} + target = _resolve_single_delivery_target(job, "origin") + assert target["thread_id"] == "2203" + + def test_explicit_target_no_reattach_when_chat_is_home(self, monkeypatch): + """slack: must not inherit the stale creation thread.""" + _slack_home(monkeypatch) + monkeypatch.setattr( + "tools.send_message_tool.prepare_send_message_platforms", lambda: None) + monkeypatch.setattr( + "tools.send_message_tool.resolve_send_target", + lambda platform, rest: (rest, None, None)) + job = {"origin": {"platform": "slack", "chat_id": "D0BJTDCSR7C", + "thread_id": SYNTH}} + target = _resolve_single_delivery_target(job, "slack:D0BJTDCSR7C") + assert target["thread_id"] is None + + def test_explicit_target_reattach_kept_for_non_home_chat(self, monkeypatch): + """Origin-affinity re-attach is preserved for genuine non-home threads.""" + _slack_home(monkeypatch, chat_id="D_OTHER_HOME") + monkeypatch.setattr( + "tools.send_message_tool.prepare_send_message_platforms", lambda: None) + monkeypatch.setattr( + "tools.send_message_tool.resolve_send_target", + lambda platform, rest: (rest, None, None)) + job = {"origin": {"platform": "slack", "chat_id": "C0AGENERAL", + "thread_id": "1755040000.000100"}} + target = _resolve_single_delivery_target(job, "slack:C0AGENERAL") + assert target["thread_id"] == "1755040000.000100" + + +def _gateway_config(connected_values): + config = MagicMock() + config.get_connected_platforms.return_value = [ + MagicMock(value=v) for v in connected_values + ] + return config + + +class TestPreflightRelayFronted: + def test_relay_fronted_slack_accepted(self, monkeypatch): + """Relay-only topology fronting slack: slack:CHAT passes preflight.""" + monkeypatch.setenv("GATEWAY_RELAY_PLATFORMS", "slack") + with patch("gateway.config.load_gateway_config", + return_value=_gateway_config({"relay"})): + assert _preflight_check_delivery( + {"deliver": "slack:D0BJTDCSR7C"}) is None + + def test_unfronted_platform_still_rejected(self, monkeypatch): + """The relay fronting slack does not whitelist other platforms.""" + monkeypatch.setenv("GATEWAY_RELAY_PLATFORMS", "slack") + with patch("gateway.config.load_gateway_config", + return_value=_gateway_config({"relay"})): + reason = _preflight_check_delivery({"deliver": "discord:12345"}) + assert reason is not None + assert "discord" in reason + + def test_native_strictness_without_relay(self, monkeypatch): + """No relay configured: the native credential check is unchanged.""" + monkeypatch.delenv("GATEWAY_RELAY_PLATFORMS", raising=False) + with patch("gateway.config.load_gateway_config", + return_value=_gateway_config({"telegram"})): + reason = _preflight_check_delivery( + {"deliver": "slack:D0BJTDCSR7C"}) + assert reason is not None + assert "slack" in reason + + def test_delivery_targets_include_relay_fronted(self, monkeypatch): + """The UI dropdown source offers relay-fronted platforms.""" + monkeypatch.setenv("GATEWAY_RELAY_PLATFORMS", "slack") + _slack_home(monkeypatch) + monkeypatch.setattr(sched, "_iter_home_target_platforms", + lambda: ["slack", "telegram"]) + with patch("gateway.config.load_gateway_config", + return_value=_gateway_config({"relay"})): + ids = {t["id"] for t in cron_delivery_targets()} + assert "slack" in ids + assert "telegram" not in ids diff --git a/tools/cronjob_tools.py b/tools/cronjob_tools.py index 19d6503165..6fb94b35b6 100644 --- a/tools/cronjob_tools.py +++ b/tools/cronjob_tools.py @@ -319,6 +319,23 @@ def _origin_from_env() -> Optional[Dict[str, str]]: origin_chat_id = get_session_env("HERMES_SESSION_CHAT_ID") if origin_platform and origin_chat_id: thread_id = get_session_env("HERMES_SESSION_THREAD_ID") or None + # Slack thread-per-message session keying (native parity: thread_ts = + # event.thread_ts or ts) stamps every TOP-LEVEL message's own id as + # the session thread. That stamp is a per-message session KEY, not a + # durable conversation location — persisting it as origin routing + # pins every future delivery inside the ephemeral thread spawned + # around the creation message. Recognize it at the source: a Slack + # thread id equal to the triggering message's own id is synthetic. + # A genuine in-thread creation (thread == the parent's id != this + # message's id) keeps its thread. + if thread_id and origin_platform == "slack": + message_id = get_session_env("HERMES_SESSION_MESSAGE_ID") or None + if message_id and str(thread_id) == str(message_id): + logger.debug( + "Cron origin: dropping synthetic per-message Slack " + "thread_id=%s (== creation message id)", thread_id, + ) + thread_id = None if thread_id: logger.debug( "Cron origin captured thread_id=%s for %s:%s",