diff --git a/cron/scheduler_delivery.py b/cron/scheduler_delivery.py index 8ecbde4ff0..54b9a728d9 100644 --- a/cron/scheduler_delivery.py +++ b/cron/scheduler_delivery.py @@ -1444,8 +1444,12 @@ def _live_send_text( platform=t.platform, chat_id=str(t.chat_id), thread_id=route_thread_id, is_explicit=True) # Thread routing goes via the target, not a bare metadata "thread_id": the router only applies # its Telegram DM-topic detection when thread_id/message_thread_id are absent from metadata. + # Send through the already-authorized transport: re-resolving from the plain target_adapters + # dict cannot re-derive the SharedRouteAdapters satellite grant (the satellite owned + # platforms.

block is disabled), yields None, and drops the delivery (#115656). future = safe_schedule_threadsafe( - router._deliver_to_platform(route_target, text_to_send, route_metadata), t.loop) + router._deliver_to_platform( + route_target, text_to_send, route_metadata, transport=t.transport), t.loop) if future is None: target_errors.append("live adapter event loop scheduling failed") return False, False, None diff --git a/gateway/delivery.py b/gateway/delivery.py index 6585cfbe1e..0aac015a1c 100644 --- a/gateway/delivery.py +++ b/gateway/delivery.py @@ -252,9 +252,19 @@ class DeliveryRouter: return content[:max(0, MAX_PLATFORM_OUTPUT - len(footer))] + footer async def _deliver_to_platform(self, target: DeliveryTarget, content: str, - metadata: Optional[Dict[str, Any]]) -> Dict[str, Any]: - """Deliver content to a messaging platform.""" - transport = resolve_delivery_transport(target.platform, self.config, self.adapters) + metadata: Optional[Dict[str, Any]], + transport: Optional[DeliveryTransport] = None, + ) -> Dict[str, Any]: + """Deliver content to a messaging platform. + + ``transport`` carries an already-authorized transport past resolution: + the cron live lane resolved and authorized it per target (including the + SharedRouteAdapters satellite grant), and re-resolving from the plain + adapters dict cannot re-derive that grant under satellite config + (#115656). Omitted (None) preserves resolution for every other caller. + """ + if transport is None: + transport = resolve_delivery_transport(target.platform, self.config, self.adapters) if transport is None: raise ValueError(f"No adapter configured for {target.platform.value}") if not target.chat_id: diff --git a/tests/cron/test_cron_live_delivery_confirmation.py b/tests/cron/test_cron_live_delivery_confirmation.py index cdce4a8e4e..6702bc7879 100644 --- a/tests/cron/test_cron_live_delivery_confirmation.py +++ b/tests/cron/test_cron_live_delivery_confirmation.py @@ -157,7 +157,7 @@ def _run(job, content, send_result, relay=False, standalone_result=None, cron_cf router = MagicMock() - async def _deliver_to_platform(target, text, metadata): + async def _deliver_to_platform(target, text, metadata, transport=None): router_calls.append({"target": target, "text": text, "metadata": metadata}) return send_result diff --git a/tests/cron/test_cron_send_authorized_transport.py b/tests/cron/test_cron_send_authorized_transport.py new file mode 100644 index 0000000000..df745ebb92 --- /dev/null +++ b/tests/cron/test_cron_send_authorized_transport.py @@ -0,0 +1,125 @@ +"""Live-lane text send uses the already-authorized transport (#115656). + +``_resolve_target_transport`` authorizes a credentialless satellite's exact +profile_routes target through the PRIMARY adapter (the SharedRouteAdapters +grant, #101113), building a transport with the satellite's own +``platforms.

`` block force-enabled. ``_live_send_text`` then rebuilt a +``DeliveryRouter`` over the plain ``target_adapters`` dict and re-resolved: +the satellite's real (disabled) block vetoed the grant, the second +resolution yielded None, and delivery fell through to the credentialless +standalone lane and failed. + +The live send must go through the already-authorized ``t.transport`` — no +second resolution. These tests pin it with a REAL DeliveryRouter and a +satellite config whose own ``platforms.discord`` block is disabled: the only +way the send succeeds is via the authorized transport. +""" +import asyncio +from concurrent.futures import Future +from unittest.mock import MagicMock, patch + +import yaml + +from cron.scheduler import _deliver_result +from cron.scheduler_preflight import SharedRouteAdapters, _primary_profile_routes_for_current_home +from gateway.config import Platform, PlatformConfig +from hermes_constants import reset_hermes_home_override, set_hermes_home_override + +PRIMARY_YAML = { + "gateway": { + "multiplex_profiles": True, + "profile_routes": [ + {"name": "fit", "platform": "discord", "chat_id": "1543065293755256852", "profile": "fitness"}, + ], + } +} + +CHAT_ID = "1543065293755256852" + + +def _job(chat_id: str) -> dict: + return {"id": "a7ae1520356c", "name": "brief", "deliver": f"discord:{chat_id}"} + + +def _run(job, adapters): + """Drive ``_deliver_result`` with a live loop and a real DeliveryRouter. + + The satellite's own ``platforms.discord`` block is DISABLED (a connector + the credentialless satellite never runs) — the shape under which a second + router-side resolution yields None. + """ + loop = MagicMock() + loop.is_running.return_value = True + + def fake_run_coro(coro, _loop): + future = Future() + future.set_result(asyncio.run(coro)) + return future + + standalone = [] + + async def _fake_send_to_platform(platform, pconfig, chat_id, text, **kwargs): + standalone.append(chat_id) + return {"success": False, "error": "DISCORD_BOT_TOKEN is not set"} + + config = MagicMock() + config.platforms = {Platform.DISCORD: PlatformConfig(enabled=False)} + config.get_home_channel = lambda p: None + with patch("gateway.config.load_gateway_config", return_value=config), \ + patch("cron.scheduler.load_config", return_value={"cron": {"wrap_response": False}}), \ + patch("tools.send_message_tool._send_to_platform", _fake_send_to_platform), \ + patch("asyncio.run_coroutine_threadsafe", side_effect=fake_run_coro): + error = _deliver_result(job, "hello", adapters=adapters, loop=loop) + return error, standalone + + +def _primary_adapter(): + adapter = MagicMock() + adapter.sent = [] + + async def send(chat_id, content, metadata=None): + adapter.sent.append(chat_id) + return {"success": True, "message_id": "m1"} + + adapter.send = send + return adapter + + +def _shared_view(tmp_path, monkeypatch, primary): + root = tmp_path / "root" + fitness_home = root / "profiles" / "fitness" + fitness_home.mkdir(parents=True) + (root / "config.yaml").write_text(yaml.safe_dump(PRIMARY_YAML), encoding="utf-8") + monkeypatch.setattr("hermes_constants.get_default_hermes_root", lambda: root) + token = set_hermes_home_override(str(fitness_home)) + try: + yield SharedRouteAdapters( + {Platform.DISCORD: primary}, _primary_profile_routes_for_current_home() + ) + finally: + reset_hermes_home_override(token) + + +def test_satellite_disabled_block_sends_through_authorized_transport(tmp_path, monkeypatch): + """Exact route + disabled satellite block: live lane delivers, no standalone.""" + primary = _primary_adapter() + for shared in _shared_view(tmp_path, monkeypatch, primary): + error, standalone = _run(_job(CHAT_ID), shared) + assert error is None, error + assert primary.sent == [CHAT_ID] + assert standalone == [] + + +def test_satellite_disabled_block_still_fails_closed_off_route(tmp_path, monkeypatch): + """Unmatched chat: the primary bot is NEVER used (#101113 fail-closed). + + With the satellite's own block disabled, an off-route chat fails closed + at transport resolution (never reaching any sender) — the pin is that the + primary adapter stays silent and the run reports an error. + """ + primary = _primary_adapter() + for shared in _shared_view(tmp_path, monkeypatch, primary): + error, standalone = _run(_job("424242"), shared) + assert error is not None + assert primary.sent == [] + assert standalone == [] diff --git a/tests/cron/test_relay_fronted_delivery.py b/tests/cron/test_relay_fronted_delivery.py index 18048a10de..63be423d08 100644 --- a/tests/cron/test_relay_fronted_delivery.py +++ b/tests/cron/test_relay_fronted_delivery.py @@ -138,7 +138,7 @@ class TestRelayDeliveryGate: router = MagicMock() - async def _deliver_to_platform(target, content, metadata): + async def _deliver_to_platform(target, content, metadata, transport=None): return {"success": True, "raw_response": None} router._deliver_to_platform = _deliver_to_platform diff --git a/tests/cron/test_scheduler.py b/tests/cron/test_scheduler.py index 3a9ceade70..e0745909dd 100644 --- a/tests/cron/test_scheduler.py +++ b/tests/cron/test_scheduler.py @@ -2641,7 +2641,7 @@ class TestCronContinuableSurfaceInChannel: def __init__(self, *a, **k): pass - async def _deliver_to_platform(self, target, text, metadata): + async def _deliver_to_platform(self, target, text, metadata, transport=None): captured["target"] = target return {"success": True, "message_id": "msg_1"} @@ -2680,7 +2680,7 @@ class TestCronContinuableSurfaceInChannel: def __init__(self, *a, **k): pass - async def _deliver_to_platform(self, target, text, metadata): + async def _deliver_to_platform(self, target, text, metadata, transport=None): captured["metadata"] = metadata return {"success": True, "message_id": "msg_1"} @@ -2712,7 +2712,7 @@ class TestCronContinuableSurfaceInChannel: def __init__(self, *a, **k): pass - async def _deliver_to_platform(self, target, text, metadata): + async def _deliver_to_platform(self, target, text, metadata, transport=None): captured["metadata"] = metadata return {"success": True, "message_id": "msg_1"}