fix(cron): live-lane text send uses the already-authorized transport (#115656)
_live_send_text rebuilt DeliveryRouter over the plain target_adapters dict and re-resolved, discarding the SharedRouteAdapters-authorized transport from _resolve_target_transport. Under satellite config (own platforms.<p> block disabled) the second resolution yields None and delivery falls through to the credentialless standalone lane and fails. DeliveryRouter._deliver_to_platform takes an optional transport; the live lane passes t.transport. All router pipeline behavior (oversize cap, relay home stamping, thread routing) is unchanged; other callers still resolve. (cherry picked from commit 9a2a67647cd823fe9120891d863a6fc6981fec71)
This commit is contained in:
@@ -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.<p> 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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
125
tests/cron/test_cron_send_authorized_transport.py
Normal file
125
tests/cron/test_cron_send_authorized_transport.py
Normal file
@@ -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.<p>`` 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 == []
|
||||
@@ -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
|
||||
|
||||
@@ -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"}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user