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
17 changes: 17 additions & 0 deletions cron/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -2225,6 +2225,7 @@ def create_job(
monitor_script: Optional[str] = None,
monitor_url: Optional[str] = None,
reasoning_effort: Optional[str] = None,
failure_deliver: Optional[str] = None,
) -> Dict[str, Any]:
"""
Create a new cron job.
Expand Down Expand Up @@ -2327,6 +2328,18 @@ def create_job(
normalized_no_agent = bool(no_agent)
normalized_attach = attach_to_session if isinstance(attach_to_session, bool) else None
normalized_reasoning_effort = _normalize_reasoning_effort(reasoning_effort)
# failure_deliver shares deliver's value grammar; the str/list
# flatten below mirrors the tool layer's _normalize_deliver_param for
# direct create_job callers (the tool pre-normalizes). Semantic
# validation happens at resolution time via the shared deliver path.
normalized_failure_deliver = (
str(failure_deliver).strip() if isinstance(failure_deliver, str) else None
)
if isinstance(failure_deliver, (list, tuple)):
normalized_failure_deliver = ",".join(
str(p).strip() for p in failure_deliver if str(p).strip()
)
normalized_failure_deliver = normalized_failure_deliver or None
normalized_monitor_script = str(monitor_script).strip() if isinstance(monitor_script, str) else None
normalized_monitor_script = normalized_monitor_script or None
normalized_monitor_url = str(monitor_url).strip() if isinstance(monitor_url, str) else None
Expand Down Expand Up @@ -2443,6 +2456,10 @@ def create_job(
# absent key = job follows config resolution (pre-feature behavior).
if normalized_reasoning_effort is not None:
job["reasoning_effort"] = normalized_reasoning_effort
# Conditional-persist for failure_deliver too: absent key = failures
# follow deliver, byte-identical to pre-feature jobs (NS-788).
if normalized_failure_deliver is not None:
job["failure_deliver"] = normalized_failure_deliver

with _jobs_lock():
jobs = load_jobs()
Expand Down
59 changes: 49 additions & 10 deletions cron/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -2837,7 +2837,20 @@ def _expand_routing_tokens(part: str) -> List[str]:
return expanded


def _resolve_delivery_targets(job: dict) -> List[dict]:
def _delivery_lane_value(job: dict, *, for_failure: bool = False):
"""Raw deliver-lane value for a run outcome: the failure lane when
``for_failure`` and the job overrides it, else ``deliver``. Keeps
delivery bookkeeping (outcome classification, unresolved-origin,
incident 'alerted' marking) reading the SAME lane the notice was
actually routed through (NS-788 review finding B1)."""
if for_failure:
failure_deliver = job.get("failure_deliver")
if failure_deliver is not None and str(failure_deliver).strip():
return failure_deliver
return job.get("deliver", "local")


def _resolve_delivery_targets(job: dict, *, for_failure: bool = False) -> List[dict]:
"""Resolve all concrete auto-delivery targets for a cron job.

Accepts the legacy comma-separated ``deliver`` string plus the
Expand All @@ -2846,8 +2859,17 @@ def _resolve_delivery_targets(job: dict) -> List[dict]:
targets: ``origin,all`` and ``all,telegram:-100:17`` both work.
Duplicate (platform, chat_id, thread_id) tuples are collapsed by the
existing dedup pass.

``for_failure=True`` resolves failure-category engine notices
(failure summaries, interrupted-run notices, drift/preflight
alerts): when the job carries a ``failure_deliver`` value, targets
resolve from it INSTEAD of ``deliver`` — ``failure_deliver: local``
is the structural opt-out for shared channels (NS-788, Coatue).
Absent ``failure_deliver``, failure delivery follows ``deliver``
exactly as before.
"""
deliver = _normalize_deliver_value(job.get("deliver", "local"))
deliver_raw = _delivery_lane_value(job, for_failure=for_failure)
deliver = _normalize_deliver_value(deliver_raw)
if deliver == "local":
return []

Expand Down Expand Up @@ -3060,7 +3082,9 @@ def _is_channel_dm_topic(
return is_channel


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

Expand All @@ -3069,11 +3093,16 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option
the standalone HTTP path cannot encrypt. Falls back to standalone send if
the adapter path fails or is unavailable.

``for_failure=True`` routes failure-category engine notices through the
job's ``failure_deliver`` override when present (NS-788).

Returns None on success, or an error string on failure.
"""
targets = _resolve_delivery_targets(job)
targets = _resolve_delivery_targets(job, for_failure=for_failure)
if not targets:
deliver_value = _normalize_deliver_value(job.get("deliver", "local"))
deliver_value = _normalize_deliver_value(
_delivery_lane_value(job, for_failure=for_failure)
)
if deliver_value == "local":
return None # local-only jobs don't deliver — not a failure
# deliver=origin with no resolvable origin and no configured home
Expand Down Expand Up @@ -7422,8 +7451,9 @@ def _fire_claim_ownership_lost() -> bool:

if should_deliver:
unresolved_origin = (
_normalize_deliver_value(job.get("deliver", "local")) == "origin"
and not _resolve_delivery_targets(job)
_normalize_deliver_value(_delivery_lane_value(job, for_failure=not success))
== "origin"
and not _resolve_delivery_targets(job, for_failure=not success)
)
try:
with _side_effect_fence() as owns_delivery:
Expand All @@ -7435,6 +7465,10 @@ def _fire_claim_ownership_lost() -> bool:
deliver_content,
adapters=adapters,
loop=loop,
# Failure summaries (and drift/blocked-config alerts
# composed into deliver_content on the failure path)
# honor the job's failure_deliver override (NS-788).
for_failure=not success,
)
except Exception as de:
if isinstance(de, _FireClaimLostDuringSideEffect):
Expand Down Expand Up @@ -7522,7 +7556,9 @@ def _fire_claim_ownership_lost() -> bool:
error="Fire claim ownership lost before terminal completion.",
)
return True
normalized_deliver = _normalize_deliver_value(job.get("deliver", "local"))
normalized_deliver = _normalize_deliver_value(
_delivery_lane_value(job, for_failure=not success)
)
if delivery_error:
delivery_outcome = "failed"
elif should_deliver and unresolved_origin:
Expand Down Expand Up @@ -7579,7 +7615,7 @@ def _fire_claim_ownership_lost() -> bool:
and not _fire_claim_ownership_lost()
):
normalized_deliver = _normalize_deliver_value(
job.get("deliver", "local")
_delivery_lane_value(job, for_failure=True)
)
unresolved_origin = False
# Durable failure incident: same ack gate as the normal failure
Expand All @@ -7605,14 +7641,17 @@ def _fire_claim_ownership_lost() -> bool:
+ _failure_streak_nudge(job),
adapters=adapters,
loop=loop,
for_failure=True,
)
except Exception as delivery_exc:
delivery_error = str(delivery_exc)
logger.error(
"Delivery failed for job %s: %s", job["id"], delivery_exc
)
if not delivery_error and normalized_deliver == "origin":
unresolved_origin = not _resolve_delivery_targets(job)
unresolved_origin = not _resolve_delivery_targets(
job, for_failure=True
)
if delivery_error:
delivery_outcome = "failed"
elif unresolved_origin:
Expand Down
4 changes: 3 additions & 1 deletion gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -11470,7 +11470,9 @@ async def _notify_interrupted_cron_jobs(self, job_ids) -> int:
# deliver=local jobs — and deliver=origin jobs with no
# resolvable origin (#43014) — resolve to zero targets and
# must stay silent rather than fall back to a home channel.
targets = _resolve_delivery_targets(job)
# Interrupted notices are failure-category engine status, so
# they honor the job's failure_deliver override (NS-788).
targets = _resolve_delivery_targets(job, for_failure=True)
except Exception as e:
logger.debug("Cron interrupt targets unresolved for %s: %s", job_id, e)
continue
Expand Down
2 changes: 2 additions & 0 deletions hermes_cli/cron.py
Original file line number Diff line number Diff line change
Expand Up @@ -690,6 +690,7 @@ def cron_create(args):
prompt=args.prompt,
name=getattr(args, "name", None),
deliver=getattr(args, "deliver", None),
failure_deliver=getattr(args, "failure_deliver", None),
repeat=getattr(args, "repeat", None),
skill=getattr(args, "skill", None),
skills=_normalize_skills(getattr(args, "skill", None), getattr(args, "skills", None)),
Expand Down Expand Up @@ -766,6 +767,7 @@ def cron_edit(args):
prompt=getattr(args, "prompt", None),
name=getattr(args, "name", None),
deliver=getattr(args, "deliver", None),
failure_deliver=getattr(args, "failure_deliver", None),
repeat=getattr(args, "repeat", None),
skills=final_skills,
script=getattr(args, "script", None),
Expand Down
18 changes: 18 additions & 0 deletions hermes_cli/subcommands/cron.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,16 @@ def build_cron_parser(subparsers, *, cmd_cron: Callable) -> None:
"local profile's canonical Bot Chat as a message the bot responds to)"
),
)
cron_create.add_argument(
"--failure-deliver",
dest="failure_deliver",
help=(
"Override target for FAILURE notices only (same grammar as "
"--deliver). 'local' suppresses failure notices entirely; run "
"state stays visible in `hermes cron list`. Omit = failures "
"follow --deliver."
),
)
cron_create.add_argument("--repeat", type=int, help="Optional repeat count")
cron_create.add_argument(
"--skill",
Expand Down Expand Up @@ -142,6 +152,14 @@ def build_cron_parser(subparsers, *, cmd_cron: Callable) -> None:
cron_edit.add_argument("--prompt", help="New prompt/task instruction")
cron_edit.add_argument("--name", help="New job name")
cron_edit.add_argument("--deliver", help="New delivery target")
cron_edit.add_argument(
"--failure-deliver",
dest="failure_deliver",
help=(
"Override target for failure notices (same grammar as --deliver; "
"'local' suppresses; '' clears the override)"
),
)
cron_edit.add_argument("--repeat", type=int, help="New repeat count")
cron_edit.add_argument(
"--skill",
Expand Down
4 changes: 2 additions & 2 deletions tests/cron/test_cron_drift_alert_once.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ def _tick(job, tmp_path, current_provider, deliveries):
"""Run one run_one_job tick with the provider resolution pinned."""
fake_db = MagicMock()

def fake_deliver(job, content, adapters=None, loop=None):
def fake_deliver(job, content, adapters=None, loop=None, **kwargs):
deliveries.append(content)
return None

Expand Down Expand Up @@ -129,7 +129,7 @@ def test_non_drift_failures_untouched_by_the_bit(self, tmp_path):
job = _job(provider_snapshot=None, drift_alerted=True)
deliveries = []

def fake_deliver(jb, content, adapters=None, loop=None):
def fake_deliver(jb, content, adapters=None, loop=None, **kwargs):
deliveries.append(content)
return None

Expand Down
Loading
Loading