Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 37 additions & 4 deletions cron/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -615,7 +615,13 @@ def _send_media_via_adapter(
logger.warning("Job '%s': failed to send media %s: %s", job.get("id", "?"), media_path, e)


def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Optional[str]:
def _deliver_result(
job: dict,
content: str,
adapters=None,
loop=None,
status_hint: str = "ok",
) -> Optional[str]:
"""
Deliver job output to the configured target(s) (origin chat, specific platform, etc.).

Expand Down Expand Up @@ -674,13 +680,29 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option

delivery_errors = []

# Structured side-channel so plugin adapters can recover cron context
# without regex-parsing the envelope text. Invariant across targets, so
# built once before the loop; only thread_id varies per target.
# Filter on `is not None` rather than truthiness so falsy-but-meaningful
# values (e.g. an explicitly empty schedule) reach adapters.
origin = _resolve_origin(job) or {}
cron_meta_full = {
"job_id": job.get("id", ""),
"job_name": job.get("name") or job.get("id", ""),
"schedule": job.get("schedule"),
"deliver": job.get("deliver"),
"origin": origin or None,
"status": status_hint,
"ran_at": _hermes_now().isoformat(),
}
cron_meta = {k: v for k, v in cron_meta_full.items() if v is not None}

for target in targets:
platform_name = target["platform"]
chat_id = target["chat_id"]
thread_id = target.get("thread_id")

# Diagnostic: log thread_id for topic-aware delivery debugging
origin = _resolve_origin(job) or {}
origin_thread = origin.get("thread_id")
if origin_thread and not thread_id:
logger.warning(
Expand Down Expand Up @@ -711,12 +733,17 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option
delivery_errors.append(msg)
continue

send_metadata: Optional[dict] = None
if thread_id:
send_metadata = {"thread_id": thread_id}
if cron_meta:
send_metadata = {**(send_metadata or {}), "cron": cron_meta}

# Prefer the live adapter when the gateway is running — this supports E2EE
# rooms (e.g. Matrix) where the standalone HTTP path cannot encrypt.
runtime_adapter = (adapters or {}).get(platform)
delivered = False
if runtime_adapter is not None and loop is not None and getattr(loop, "is_running", lambda: False)():
send_metadata = {"thread_id": thread_id} if thread_id else None
try:
# Send cleaned text (MEDIA tags stripped) — not the raw content
text_to_send = cleaned_delivery_content.strip()
Expand Down Expand Up @@ -1952,7 +1979,13 @@ def _process_job(job: dict) -> bool:
delivery_error = None
if should_deliver:
try:
delivery_error = _deliver_result(job, deliver_content, adapters=adapters, loop=loop)
delivery_error = _deliver_result(
job,
deliver_content,
adapters=adapters,
loop=loop,
status_hint="ok" if success else "error",
)
except Exception as de:
delivery_error = str(de)
logger.error("Delivery failed for job %s: %s", job["id"], de)
Expand Down
224 changes: 219 additions & 5 deletions tests/cron/test_scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -837,6 +837,216 @@ def test_origin_delivery_preserves_thread_id(self):
assert send_mock.call_args.kwargs["thread_id"] == "17585"


class TestDeliverResultCronMetadata:
"""Live adapter sends should receive structured cron context in metadata.

Without this, plugin adapters have to regex-parse the "Cronjob Response: ..."
envelope to recover job_id/job_name, which is brittle and breaks entirely
when cron.wrap_response=false. The structured side-channel under
metadata["cron"] makes job_id/job_name/schedule/deliver/origin available
via the existing metadata kwarg, alongside thread_id.
"""

def _build_adapter_send_mock(self):
from concurrent.futures import Future

adapter = AsyncMock()
adapter.send.return_value = MagicMock(success=True)

def fake_run_coro(coro, _loop):
future = Future()
future.set_result(MagicMock(success=True))
coro.close()
return future

loop = MagicMock()
loop.is_running.return_value = True
return adapter, loop, fake_run_coro

def test_live_adapter_send_receives_cron_metadata(self):
from gateway.config import Platform

adapter, loop, fake_run_coro = self._build_adapter_send_mock()
pconfig = MagicMock()
pconfig.enabled = True
mock_cfg = MagicMock()
mock_cfg.platforms = {Platform.TELEGRAM: pconfig}

job = {
"id": "abc123",
"name": "daily-summary",
"schedule": "0 9 * * *",
"deliver": "origin",
"origin": {"platform": "telegram", "chat_id": "999"},
}

with patch("gateway.config.load_gateway_config", return_value=mock_cfg), \
patch("cron.scheduler.load_config", return_value={"cron": {"wrap_response": False}}), \
patch("asyncio.run_coroutine_threadsafe", side_effect=fake_run_coro):
_deliver_result(
job,
"hello",
adapters={Platform.TELEGRAM: adapter},
loop=loop,
)

adapter.send.assert_called_once()
metadata = adapter.send.call_args.kwargs.get("metadata")
assert metadata is not None, "send_metadata should be populated with cron context"
assert "cron" in metadata, f"cron context missing: {metadata!r}"
cron_ctx = metadata["cron"]
assert cron_ctx["job_id"] == "abc123"
assert cron_ctx["job_name"] == "daily-summary"
assert cron_ctx["schedule"] == "0 9 * * *"
assert cron_ctx["deliver"] == "origin"
assert cron_ctx["origin"] == {"platform": "telegram", "chat_id": "999"}

def test_cron_metadata_falls_back_to_job_id_when_name_missing(self):
from gateway.config import Platform

adapter, loop, fake_run_coro = self._build_adapter_send_mock()
pconfig = MagicMock()
pconfig.enabled = True
mock_cfg = MagicMock()
mock_cfg.platforms = {Platform.TELEGRAM: pconfig}

job = {
"id": "noname-job",
"deliver": "origin",
"origin": {"platform": "telegram", "chat_id": "5"},
}

with patch("gateway.config.load_gateway_config", return_value=mock_cfg), \
patch("cron.scheduler.load_config", return_value={"cron": {"wrap_response": False}}), \
patch("asyncio.run_coroutine_threadsafe", side_effect=fake_run_coro):
_deliver_result(
job,
"hi",
adapters={Platform.TELEGRAM: adapter},
loop=loop,
)

cron_ctx = adapter.send.call_args.kwargs["metadata"]["cron"]
assert cron_ctx["job_name"] == "noname-job"

def test_cron_metadata_coexists_with_thread_id(self):
from gateway.config import Platform

adapter, loop, fake_run_coro = self._build_adapter_send_mock()
pconfig = MagicMock()
pconfig.enabled = True
mock_cfg = MagicMock()
mock_cfg.platforms = {Platform.TELEGRAM: pconfig}

job = {
"id": "topic-job",
"name": "topic-job",
"deliver": "origin",
"origin": {
"platform": "telegram",
"chat_id": "-1001",
"thread_id": "17585",
},
}

with patch("gateway.config.load_gateway_config", return_value=mock_cfg), \
patch("cron.scheduler.load_config", return_value={"cron": {"wrap_response": False}}), \
patch("asyncio.run_coroutine_threadsafe", side_effect=fake_run_coro):
_deliver_result(
job,
"hi",
adapters={Platform.TELEGRAM: adapter},
loop=loop,
)

metadata = adapter.send.call_args.kwargs["metadata"]
assert metadata["thread_id"] == "17585"
assert metadata["cron"]["job_id"] == "topic-job"

def test_cron_metadata_omits_empty_fields(self):
"""Optional fields (schedule, deliver) absent on the job should be omitted."""
from gateway.config import Platform

adapter, loop, fake_run_coro = self._build_adapter_send_mock()
pconfig = MagicMock()
pconfig.enabled = True
mock_cfg = MagicMock()
mock_cfg.platforms = {Platform.TELEGRAM: pconfig}

job = {
"id": "min-job",
"deliver": "origin",
"origin": {"platform": "telegram", "chat_id": "5"},
} # no name, no schedule

with patch("gateway.config.load_gateway_config", return_value=mock_cfg), \
patch("cron.scheduler.load_config", return_value={"cron": {"wrap_response": False}}), \
patch("asyncio.run_coroutine_threadsafe", side_effect=fake_run_coro):
_deliver_result(
job,
"hi",
adapters={Platform.TELEGRAM: adapter},
loop=loop,
)

cron_ctx = adapter.send.call_args.kwargs["metadata"]["cron"]
assert "schedule" not in cron_ctx
# job_id is present (required), job_name falls back to it,
# deliver and origin are present.
# status defaults to "ok"; ran_at is always set (ISO timestamp).
ran_at = cron_ctx.pop("ran_at", None)
assert isinstance(ran_at, str) and "T" in ran_at
assert cron_ctx == {
"job_id": "min-job",
"job_name": "min-job",
"deliver": "origin",
"origin": {"platform": "telegram", "chat_id": "5"},
"status": "ok",
}

def test_cron_metadata_attached_on_failed_run_with_wrap_disabled(self):
"""Failed run + wrap_response=false must still surface status=error and full context.

Regression guard for #26004 step 2: without this, adapters that rely on
`metadata["cron"]["status"]` to distinguish success from failure (e.g.
for routing failure summaries to a different chat) would always see
the default "ok" even when the agent raised.
"""
from gateway.config import Platform

adapter, loop, fake_run_coro = self._build_adapter_send_mock()
pconfig = MagicMock()
pconfig.enabled = True
mock_cfg = MagicMock()
mock_cfg.platforms = {Platform.TELEGRAM: pconfig}

job = {
"id": "fail-job",
"name": "nightly-report",
"schedule": "0 3 * * *",
"deliver": "origin",
"origin": {"platform": "telegram", "chat_id": "999"},
}

with patch("gateway.config.load_gateway_config", return_value=mock_cfg), \
patch("cron.scheduler.load_config", return_value={"cron": {"wrap_response": False}}), \
patch("asyncio.run_coroutine_threadsafe", side_effect=fake_run_coro):
_deliver_result(
job,
"Job hit an exception",
adapters={Platform.TELEGRAM: adapter},
loop=loop,
status_hint="error",
)

cron_ctx = adapter.send.call_args.kwargs["metadata"]["cron"]
assert cron_ctx["status"] == "error"
assert cron_ctx["job_id"] == "fail-job"
assert cron_ctx["job_name"] == "nightly-report"
assert cron_ctx["schedule"] == "0 3 * * *"
assert "ran_at" in cron_ctx and "T" in cron_ctx["ran_at"]


class TestDeliverResultErrorReturns:
"""Verify _deliver_result returns error strings on failure, None on success."""

Expand Down Expand Up @@ -2526,11 +2736,15 @@ def fake_run_coro(coro, _loop):
"configured thread_id 7072 for telegram:226252250 was not found; "
"delivered without thread_id"
)
adapter.send.assert_called_once_with(
"226252250",
"Hello world",
metadata={"thread_id": "7072"},
)
assert adapter.send.call_count == 1
args, kwargs = adapter.send.call_args
assert args == ("226252250", "Hello world")
metadata = kwargs["metadata"]
# thread_id is the load-bearing field for this test — it must survive
# alongside the structured cron context added in #26004.
assert metadata["thread_id"] == "7072"
assert metadata["cron"]["job_id"] == "thread-fallback-job"
assert metadata["cron"]["deliver"] == "telegram:226252250:7072"


class TestSendMediaTimeoutCancelsFuture:
Expand Down
Loading