fix(cron): relay-fronted Slack delivery — synthetic creation-thread capture + preflight fronted-platform blindness
Bug 1: relay-fronted Slack in thread-per-message mode stamps each top-level
message's own id as source.thread_id (session KEYING, native thread_ts
parity). Cron origin capture persisted that stamp as durable routing, so
every delivery landed inside the ephemeral thread spawned around the
creation message instead of the top-level conversation. Fix at the source:
_origin_from_env drops a Slack thread id equal to the creation message's
own id (genuine in-thread creations keep theirs). Fire-time repair for
already-persisted jobs: deliver=origin and the explicit-target Slack
re-attach treat an origin thread as stale when the origin chat is the
configured Slack home chat — top-level (or the home target's configured
thread) wins; non-home working threads are preserved.
Bug 2: _preflight_check_delivery and cron_delivery_targets validated
deliver prefixes against get_connected_platforms(), which only sees
natively configured platforms — a relay-only deployment ({relay}) rejected
'slack:CHAT' with 'no gateway credentials configured' although fire-time
routing (resolve_delivery_transport + RelayAdapter.fronts_platform)
delivers it. New gateway.relay.relay_fronted_platforms() (env-derived from
GATEWAY_RELAY_PLATFORMS — the same source that seeds the live adapter's
identity set, so validation and routing cannot disagree) is unioned into
the connected set when the relay is connected. Native topologies keep the
strict credential check unchanged.
This commit is contained in:
committed by
Teknium
parent
8b243dff62
commit
58ff0fd302
@@ -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 "
|
||||
|
||||
@@ -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
|
||||
|
||||
82
tests/cron/test_cron_origin_synthetic_thread.py
Normal file
82
tests/cron/test_cron_origin_synthetic_thread.py
Normal file
@@ -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"
|
||||
148
tests/cron/test_cron_relay_delivery_guards.py
Normal file
148
tests/cron/test_cron_relay_delivery_guards.py
Normal file
@@ -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:<chat_id>`` 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:<home_chat> 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
|
||||
@@ -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",
|
||||
|
||||
Reference in New Issue
Block a user