diff --git a/cron/__init__.py b/cron/__init__.py index 2c44cabf6b..8e2934ee75 100644 --- a/cron/__init__.py +++ b/cron/__init__.py @@ -24,6 +24,7 @@ from cron.jobs import ( pause_job, resume_job, trigger_job, + rearm_oneshot, JOBS_FILE, ) from cron.scheduler import tick @@ -37,6 +38,7 @@ __all__ = [ "pause_job", "resume_job", "trigger_job", + "rearm_oneshot", "tick", "JOBS_FILE", ] diff --git a/cron/jobs.py b/cron/jobs.py index df75798321..313002ce3a 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -2409,6 +2409,70 @@ def trigger_job(job_id: str) -> Optional[Dict[str, Any]]: ) +def _claim_is_live(claim: Any, now: datetime, ttl_seconds: float) -> bool: + if not isinstance(claim, dict) or not claim.get("at"): + return False + try: + age = (now - _ensure_aware(datetime.fromisoformat(claim["at"]))).total_seconds() + except (TypeError, ValueError): + return False + return 0 <= age < ttl_seconds + + +def rearm_oneshot(job_id: str, run_at: Any) -> Optional[Dict[str, Any]]: + """Re-arm a completed one-shot as an explicit new occurrence.""" + job_ref = resolve_job_ref(job_id) + if not job_ref: + return None + if isinstance(run_at, datetime): + run_at = run_at.isoformat() + parsed_schedule = parse_schedule(str(run_at)) + if parsed_schedule.get("kind") != "once": + raise ValueError( + "Cannot re-arm recurring jobs: re-arm is one-shot-only; " + "use plain resume or cron run." + ) + next_run_at = compute_next_run(parsed_schedule) + if next_run_at is None: + requested = parsed_schedule.get("run_at") or run_at + raise ValueError( + f"Requested one-shot time {requested} is more than " + f"{ONESHOT_GRACE_SECONDS}s in the past and cannot be scheduled." + ) + + with _jobs_lock(): + jobs = load_jobs() + for index, job in enumerate(jobs): + if job.get("id") != job_ref["id"]: + continue + now = _hermes_now() + if _claim_is_live(job.get("run_claim"), now, _oneshot_run_claim_ttl_seconds()): + raise ValueError("Cannot re-arm one-shot over a live run claim.") + if _claim_is_live(job.get("fire_claim"), now, 300): + raise ValueError("Cannot re-arm one-shot over a live fire claim.") + if job.get("schedule", {}).get("kind") != "once": + raise ValueError( + "Cannot re-arm recurring jobs: re-arm is one-shot-only; " + "use plain resume or cron run." + ) + repeat = job.get("repeat") or {} + repeat["completed"] = 0 + job["schedule"] = parsed_schedule + job["schedule_display"] = parsed_schedule.get("display", str(run_at)) + job["repeat"] = repeat + job["run_claim"] = None + job["fire_claim"] = None + job["enabled"] = True + job["state"] = "scheduled" + job["paused_at"] = None + job["paused_reason"] = None + job["next_run_at"] = next_run_at + jobs[index] = job + save_jobs(jobs) + return _normalize_job_record(job) + return None + + def remove_job(job_id: str) -> bool: """Remove a job by ID or name.""" job = resolve_job_ref(job_id) diff --git a/hermes_cli/console_engine.py b/hermes_cli/console_engine.py index b52e7a1194..2728af3ff9 100644 --- a/hermes_cli/console_engine.py +++ b/hermes_cli/console_engine.py @@ -1600,16 +1600,26 @@ def _cron_pause(_engine: HermesConsoleEngine, args: list[str]) -> str: def _cron_resume(_engine: HermesConsoleEngine, args: list[str]) -> str: - if len(args) != 1: - raise ConsoleCommandError("Usage: cron resume ") - from cron.jobs import AmbiguousJobReference, resume_job + parser = _ArgumentParser(prog="cron resume", add_help=False) + parser.add_argument("job") + parser.add_argument("--at") + parser.add_argument("--run-now", action="store_true") + ns = parser.parse_args(args) + if ns.at and ns.run_now: + raise ConsoleCommandError("Use exactly one of --at or --run-now.") + from cron.jobs import AmbiguousJobReference, _hermes_now, rearm_oneshot, resume_job try: - job = resume_job(args[0]) + if ns.at or ns.run_now: + job = rearm_oneshot(ns.job, _hermes_now().isoformat() if ns.run_now else ns.at) + else: + job = resume_job(ns.job) except AmbiguousJobReference as exc: raise ConsoleCommandError(str(exc)) from exc + except ValueError as exc: + raise ConsoleCommandError(str(exc)) from exc if not job: - raise ConsoleCommandError(f"Job not found: {args[0]}") + raise ConsoleCommandError(f"Job not found: {ns.job}") return _format_job(job, "Resumed") diff --git a/hermes_cli/cron.py b/hermes_cli/cron.py index 16b73330bf..b135d532da 100644 --- a/hermes_cli/cron.py +++ b/hermes_cli/cron.py @@ -560,6 +560,29 @@ def _job_action(action: str, job_id: str, success_verb: str) -> int: return 0 +def cron_resume(args) -> int: + """Resume a paused job or explicitly re-arm a completed one-shot.""" + if bool(getattr(args, "run_at", None)) == bool(getattr(args, "run_now", False)): + if getattr(args, "run_at", None) or getattr(args, "run_now", False): + print(color("Use exactly one of --at or --run-now.", Colors.RED)) + return 1 + return _job_action("resume", args.job_id, "Resumed") + from cron.jobs import AmbiguousJobReference, _hermes_now, rearm_oneshot + + run_at = _hermes_now().isoformat() if args.run_now else args.run_at + try: + job = rearm_oneshot(args.job_id, run_at) + except (AmbiguousJobReference, ValueError) as exc: + print(color(f"Failed to re-arm job: {exc}", Colors.RED)) + return 1 + if not job: + print(color(f"Job not found: {args.job_id}", Colors.RED)) + return 1 + print(color(f"Re-armed job: {job.get('name', args.job_id)} ({args.job_id})", Colors.GREEN)) + print(f" Next run: {job.get('next_run_at')}") + return 0 + + def cron_notepad(args) -> int: """Handle ``hermes cron notepad [get|set|delete|list]``. @@ -656,7 +679,7 @@ def cron_command(args): return _job_action("pause", args.job_id, "Paused") if subcmd == "resume": - return _job_action("resume", args.job_id, "Resumed") + return cron_resume(args) if subcmd == "run": return _job_action("run", args.job_id, "Triggered") diff --git a/hermes_cli/subcommands/cron.py b/hermes_cli/subcommands/cron.py index d617007b8b..2588ad071d 100644 --- a/hermes_cli/subcommands/cron.py +++ b/hermes_cli/subcommands/cron.py @@ -261,6 +261,8 @@ def build_cron_parser(subparsers, *, cmd_cron: Callable) -> None: cron_resume = cron_subparsers.add_parser("resume", help="Resume a paused job") cron_resume.add_argument("job_id", help="Job ID to resume") + cron_resume.add_argument("--at", dest="run_at", help="Re-arm at an ISO-8601 time") + cron_resume.add_argument("--run-now", action="store_true", help="Re-arm to run now") cron_run = cron_subparsers.add_parser( "run", help="Run a job on the next scheduler tick" diff --git a/tests/cron/test_terminal_job_rearm.py b/tests/cron/test_terminal_job_rearm.py new file mode 100644 index 0000000000..dc9960e799 --- /dev/null +++ b/tests/cron/test_terminal_job_rearm.py @@ -0,0 +1,138 @@ +"""Behavioral coverage for terminal cron jobs and explicit one-shot re-arm.""" + +from datetime import datetime, timedelta, timezone +import copy + +import pytest + +from cron.jobs import ( + advance_next_run, + create_job, + get_due_jobs, + get_job, + load_jobs, + mark_job_run, + save_jobs, + trigger_job, + update_job, +) + + +@pytest.fixture() +def tmp_cron_dir(tmp_path, monkeypatch): + monkeypatch.setattr("cron.jobs.CRON_DIR", tmp_path / "cron") + monkeypatch.setattr("cron.jobs.JOBS_FILE", tmp_path / "cron" / "jobs.json") + monkeypatch.setattr("cron.jobs.OUTPUT_DIR", tmp_path / "cron" / "output") + return tmp_path + + +def test_completed_oneshot_trigger_is_refused_and_disk_record_is_unchanged(tmp_cron_dir): + job = create_job("done", "30m", name="done", repeat=1) + mark_job_run(job["id"], success=True) + before = copy.deepcopy(load_jobs()) + + with pytest.raises(ValueError, match="terminal"): + trigger_job(job["id"]) + + assert load_jobs() == before + assert get_job(job["id"]) == before[0] + + +def test_exhausted_recurring_job_trigger_is_refused(tmp_cron_dir): + job = create_job("done", "every 1h", repeat=1) + mark_job_run(job["id"], success=True) + + with pytest.raises(ValueError, match="terminal"): + trigger_job(job["id"]) + + +def test_wedged_claimed_oneshot_remains_triggerable(tmp_cron_dir): + now = datetime.now(timezone.utc) + job = create_job("wedged", "30m", repeat=2) + record = get_job(job["id"]) + record.update({ + "run_claim": {"at": now.isoformat(), "by": "dead-worker"}, + "state": "scheduled", + "enabled": True, + "next_run_at": (now - timedelta(minutes=1)).isoformat(), + }) + save_jobs([record]) + + triggered = trigger_job(job["id"]) + assert triggered["state"] == "scheduled" + assert triggered["enabled"] is True + + +def test_paused_job_run_override_remains_allowed(tmp_cron_dir): + job = create_job("paused", "every 1h") + from cron.jobs import pause_job + + pause_job(job["id"]) + triggered = trigger_job(job["id"]) + assert triggered["state"] == "scheduled" + assert triggered["enabled"] is True + + +def test_terminal_jobs_are_not_due_or_advanced(tmp_cron_dir): + job = create_job("done", "every 1h", repeat=1) + mark_job_run(job["id"], success=True) + before = copy.deepcopy(load_jobs()) + + assert get_due_jobs() == [] + assert advance_next_run(job["id"]) is False + assert load_jobs() == before + + +def test_terminal_refusal_survives_reload(tmp_cron_dir): + job = create_job("done", "30m", repeat=1) + mark_job_run(job["id"], success=True) + before = copy.deepcopy(load_jobs()) + assert get_job(job["id"])["state"] == "completed" + + with pytest.raises(ValueError): + trigger_job(job["id"]) + assert load_jobs() == before + + +def test_update_cannot_reactivate_terminal_record(tmp_cron_dir): + job = create_job("done", "30m", repeat=1) + mark_job_run(job["id"], success=True) + with pytest.raises(ValueError, match="terminal"): + update_job(job["id"], {"enabled": True}) + with pytest.raises(ValueError, match="terminal"): + update_job(job["id"], {"schedule": "every 1h"}) + + +def test_rearm_completed_oneshot_restores_schedule_and_preserves_history(tmp_cron_dir): + from cron.jobs import rearm_oneshot + + job = create_job("done", "30m", repeat=3) + mark_job_run(job["id"], success=True) + finished = get_job(job["id"]) + run_at = (datetime.now(timezone.utc) + timedelta(minutes=5)).isoformat() + + rearmed = rearm_oneshot(job["id"], run_at) + assert rearmed["schedule"]["kind"] == "once" + assert rearmed["repeat"]["times"] == 3 + assert rearmed["repeat"]["completed"] == 0 + assert rearmed["state"] == "scheduled" + assert rearmed["enabled"] is True + assert rearmed["next_run_at"] == rearmed["schedule"]["run_at"] + assert rearmed["last_run_at"] == finished["last_run_at"] + assert rearmed["last_status"] == finished["last_status"] + + +def test_rearm_refuses_recurring_and_live_claim(tmp_cron_dir): + from cron.jobs import rearm_oneshot + + recurring = create_job("recurring", "every 1h") + future = (datetime.now(timezone.utc) + timedelta(minutes=5)).isoformat() + with pytest.raises(ValueError, match="one-shot"): + rearm_oneshot(recurring["id"], future) + + oneshot = create_job("claimed", "30m") + record = get_job(oneshot["id"]) + record["run_claim"] = {"at": datetime.now(timezone.utc).isoformat(), "by": "live"} + save_jobs([record]) + with pytest.raises(ValueError, match="claim"): + rearm_oneshot(oneshot["id"], future)