Skip to content
Open
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
53 changes: 38 additions & 15 deletions cron/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -1923,6 +1923,43 @@ def heartbeat_run_claim(job_id: str, *, expected_owner: str) -> bool:
return False


def advance_next_runs(job_ids) -> int:
"""Batch form of :func:`advance_next_run` for the due-dispatch loop.

One ``load_jobs()`` + at most one ``save_jobs()`` for the whole due
set, instead of one of each per job — the per-job form costs
O(N loads + N saves) for N due jobs (~110 ms at N=50, measured), the
batch form O(1 + 1) (~2 ms). ``job_ids`` may contain ids of one-shot
or unknown jobs; they are skipped exactly as the per-job form skips
them. Returns the number of jobs whose ``next_run_at`` was advanced.

Crash semantics: the batch persists once at the end, so a crash
mid-batch re-fires the whole set on restart (at-least-once burst)
rather than advancing a prefix — acceptable given the sub-10ms window,
and identical to the per-job form once the batch completes.
"""
ids = set(job_ids)
if not ids:
return 0
with _jobs_lock():
jobs = load_jobs()
now = _hermes_now().isoformat()
advanced = 0
for job in jobs:
if job["id"] not in ids:
continue
kind = job.get("schedule", {}).get("kind")
if kind not in {"cron", "interval"}:
continue
new_next = compute_next_run(job["schedule"], now)
if new_next and new_next != job.get("next_run_at"):
job["next_run_at"] = new_next
advanced += 1
if advanced:
save_jobs(jobs)
return advanced


def advance_next_run(job_id: str) -> bool:
"""Preemptively advance next_run_at for a recurring job before execution.

Expand All @@ -1935,21 +1972,7 @@ def advance_next_run(job_id: str) -> bool:

Returns True if next_run_at was advanced, False otherwise.
"""
with _jobs_lock():
jobs = load_jobs()
for job in jobs:
if job["id"] == job_id:
kind = job.get("schedule", {}).get("kind")
if kind not in {"cron", "interval"}:
return False
now = _hermes_now().isoformat()
new_next = compute_next_run(job["schedule"], now)
if new_next and new_next != job.get("next_run_at"):
job["next_run_at"] = new_next
save_jobs(jobs)
return True
return False
return False
return advance_next_runs([job_id]) == 1


def _machine_id() -> str:
Expand Down
8 changes: 4 additions & 4 deletions cron/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -281,7 +281,7 @@ def _resolve_cron_enabled_toolsets(job: dict, cfg: dict) -> list[str] | None:
"QQBOT_HOME_CHANNEL": "QQ_HOME_CHANNEL",
}

from cron.jobs import get_due_jobs, mark_job_run, save_job_output, advance_next_run, claim_dispatch, heartbeat_run_claim
from cron.jobs import get_due_jobs, mark_job_run, save_job_output, advance_next_run, advance_next_runs, claim_dispatch, heartbeat_run_claim
from cron.executions import create_execution, finish_execution, mark_execution_running

# Sentinel: when a cron agent has nothing new to report, it can start its
Expand Down Expand Up @@ -4153,11 +4153,11 @@ def tick(

# Advance next_run_at for all recurring jobs FIRST, under the file lock,
# before any execution begins. This preserves at-most-once semantics.
# For parallel jobs that are already running, advance_next_run keeps
# For parallel jobs that are already running, the advance keeps
# bumping next_run_at forward so the grace window never expires.
# mark_job_run() overwrites next_run_at on completion.
for job in due_jobs:
advance_next_run(job["id"])
# Batched: one load + one save for the whole due set, not one per job.
advance_next_runs([job["id"] for job in due_jobs])

# Resolve max parallel workers: env var > config.yaml > unbounded.
# Set HERMES_CRON_MAX_PARALLEL=1 to restore old serial behaviour.
Expand Down
2 changes: 1 addition & 1 deletion tests/cron/test_execution_ledger.py
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,7 @@ def submit(self, _callable):
lambda execution_id, **kwargs: finished.append((execution_id, kwargs)),
)
monkeypatch.setattr(scheduler, "get_due_jobs", lambda: [{"id": "submit-fail"}])
monkeypatch.setattr(scheduler, "advance_next_run", lambda _job_id: None)
monkeypatch.setattr(scheduler, "advance_next_runs", lambda _ids: 0)
monkeypatch.setattr(scheduler, "_get_parallel_pool", lambda _workers: BrokenPool())

assert scheduler.tick(verbose=False, sync=False) == 0
Expand Down
70 changes: 70 additions & 0 deletions tests/cron/test_jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -1069,3 +1069,73 @@ def test_load_jobs_bomless_regression(self, tmp_cron_dir):
assert [j["id"] for j in loaded] == ["plainjob01"]




class TestAdvanceNextRuns:
"""Tests for advance_next_runs() — the batched due-set advance.

The scheduler's pre-dispatch loop advanced each due job individually:
N due jobs = N full load_jobs() + N full save_jobs() of the jobs file
(~110 ms at N=50, measured). The batch form does one load + at most
one save (~2 ms). Imported inside test bodies so the pre-fix tree
fails with a real test failure (ImportError), not a collection error.
"""

def _make_due(self, tmp_cron_dir, n_recurring=3, n_oneshot=1):
rec = [create_job(prompt=f"rec {i}", schedule="every 1h")
for i in range(n_recurring)]
one = [create_job(prompt=f"one {i}", schedule="30m")
for i in range(n_oneshot)]
jobs = load_jobs()
old = (datetime.now() - timedelta(minutes=5)).isoformat()
for j in jobs:
j["next_run_at"] = old
save_jobs(jobs)
return [j["id"] for j in rec], [j["id"] for j in one]

def test_batch_advances_recurring_skips_oneshots(self, tmp_cron_dir):
from cron.jobs import advance_next_runs
rec_ids, one_ids = self._make_due(tmp_cron_dir)
advanced = advance_next_runs(rec_ids + one_ids)
assert advanced == len(rec_ids)
from cron.jobs import _ensure_aware, _hermes_now
for jid in rec_ids:
nxt = _ensure_aware(datetime.fromisoformat(get_job(jid)["next_run_at"]))
assert nxt > _hermes_now()
for jid in one_ids:
# one-shots keep their (past) next_run_at for restart retry
assert datetime.fromisoformat(get_job(jid)["next_run_at"]) < datetime.now()

def test_batch_single_load_and_save(self, tmp_cron_dir, monkeypatch):
"""I/O pin: the whole due set costs one load + one save, not N+N.
Fails pre-fix (function absent) and would fail on any regression
back to per-job I/O."""
from cron.jobs import advance_next_runs
rec_ids, _ = self._make_due(tmp_cron_dir, n_recurring=10, n_oneshot=0)
import cron.jobs as cj
counts = {"load": 0, "save": 0}
real_load, real_save = cj.load_jobs, cj.save_jobs
monkeypatch.setattr(cj, "load_jobs", lambda *a, **k: (
counts.__setitem__("load", counts["load"] + 1), real_load(*a, **k))[1])
monkeypatch.setattr(cj, "save_jobs", lambda *a, **k: (
counts.__setitem__("save", counts["save"] + 1), real_save(*a, **k))[1])
advance_next_runs(rec_ids)
assert counts == {"load": 1, "save": 1}

def test_batch_no_save_when_nothing_advances(self, tmp_cron_dir, monkeypatch):
from cron.jobs import advance_next_runs
rec_ids, one_ids = self._make_due(tmp_cron_dir, n_recurring=0, n_oneshot=2)
import cron.jobs as cj
saves = [0]
real_save = cj.save_jobs
monkeypatch.setattr(cj, "save_jobs", lambda *a, **k: (
saves.__setitem__(0, saves[0] + 1), real_save(*a, **k))[1])
assert advance_next_runs(one_ids + ["missing-id"]) == 0
assert saves[0] == 0

def test_wrapper_semantics_unchanged(self, tmp_cron_dir):
"""advance_next_run keeps its per-job contract over the batch."""
rec_ids, one_ids = self._make_due(tmp_cron_dir)
assert advance_next_run(rec_ids[0]) is True
assert advance_next_run(one_ids[0]) is False
assert advance_next_run("missing-id") is False
45 changes: 42 additions & 3 deletions tests/cron/test_parallel_pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,7 @@ def test_running_set_prevents_double_dispatch(self, tmp_path, monkeypatch):

dispatched = []
monkeypatch.setattr(sched, "get_due_jobs", lambda: [job])
monkeypatch.setattr(sched, "advance_next_run", lambda *_a, **_kw: None)
monkeypatch.setattr(sched, "advance_next_runs", lambda *_a, **_kw: 0)
monkeypatch.setattr(sched, "run_job", lambda j, **_kw: dispatched.append(j["id"]) or (True, "out", "resp", None))
monkeypatch.setattr(sched, "save_job_output", lambda *_a, **_kw: None)
monkeypatch.setattr(sched, "mark_job_run", lambda *_a, **_kw: None)
Expand Down Expand Up @@ -104,7 +104,7 @@ def test_sync_true_blocks_and_returns_correct_count(self, tmp_path, monkeypatch)
]

monkeypatch.setattr(sched, "get_due_jobs", lambda: jobs)
monkeypatch.setattr(sched, "advance_next_run", lambda *_a, **_kw: None)
monkeypatch.setattr(sched, "advance_next_runs", lambda *_a, **_kw: 0)
monkeypatch.setattr(sched, "run_job", lambda j, **_kw: (True, "out", "resp", None))
monkeypatch.setattr(sched, "save_job_output", lambda *_a, **_kw: "/tmp/out")
monkeypatch.setattr(sched, "mark_job_run", lambda *_a, **_kw: None)
Expand Down Expand Up @@ -151,7 +151,7 @@ def slow_run(j, *, defer_agent_teardown=None):
return True, "out", "resp", None

monkeypatch.setattr(sched, "get_due_jobs", lambda: [job])
monkeypatch.setattr(sched, "advance_next_run", lambda *_a, **_kw: None)
monkeypatch.setattr(sched, "advance_next_runs", lambda *_a, **_kw: 0)
monkeypatch.setattr(sched, "run_job", slow_run)
monkeypatch.setattr(sched, "save_job_output", lambda *_a, **_kw: "/tmp/out")
monkeypatch.setattr(sched, "mark_job_run", lambda *_a, **_kw: None)
Expand Down Expand Up @@ -180,3 +180,42 @@ def test_get_sequential_pool_is_persistent(self):

sched._shutdown_parallel_pool()
assert sched._sequential_pool is None


class TestTickBatchAdvance:
"""The tick's pre-dispatch advance must go through advance_next_runs
exactly once with the whole due set — a revert to the per-job loop
(or back to advance_next_run) must fail this test, not slip past the
helper-level I/O pin."""

def test_tick_calls_advance_next_runs_once_with_all_due_ids(self, tmp_path, monkeypatch):
import cron.scheduler as sched

sched._parallel_pool = None
sched._parallel_pool_max_workers = None
sched._running_job_ids.clear()

jobs = [
{"id": f"job-{i}", "name": f"Job {i}", "prompt": "test",
"schedule": "every 5m", "enabled": True,
"next_run_at": "2020-01-01T00:00:00", "deliver": "local"}
for i in range(4)
]

advance_calls = []
monkeypatch.setattr(sched, "get_due_jobs", lambda: jobs)
monkeypatch.setattr(
sched, "advance_next_runs",
lambda ids: advance_calls.append(list(ids)) or len(list(ids)))
monkeypatch.setattr(sched, "run_job", lambda j, **_kw: (True, "out", "resp", None))
monkeypatch.setattr(sched, "save_job_output", lambda *_a, **_kw: "/tmp/out")
monkeypatch.setattr(sched, "mark_job_run", lambda *_a, **_kw: None)
monkeypatch.setattr(sched, "_deliver_result", lambda *_a, **_kw: None)

n = sched.tick(verbose=False)

assert n == 4
assert advance_calls == [["job-0", "job-1", "job-2", "job-3"]], (
f"tick must batch-advance the due set in ONE call; got {advance_calls}")

sched._shutdown_parallel_pool()
2 changes: 1 addition & 1 deletion tests/cron/test_run_one_job.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ def test_tick_process_job_sequence(monkeypatch):
sequence run_job → save → deliver → mark, in that order."""
calls = _patch_pipeline(monkeypatch)
monkeypatch.setattr(s, "get_due_jobs", lambda: [{"id": "j1", "name": "t"}])
monkeypatch.setattr(s, "advance_next_run", lambda jid: True)
monkeypatch.setattr(s, "advance_next_runs", lambda ids: 1)

s.tick(verbose=False, sync=True)

Expand Down
6 changes: 3 additions & 3 deletions tests/cron/test_scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -545,7 +545,7 @@ def test_tick_skips_due_jobs_while_dispatch_is_paused(self, tmp_path):
"enabled": True,
}
with patch("cron.scheduler.get_due_jobs", return_value=[job]), patch(
"cron.scheduler.advance_next_run"
"cron.scheduler.advance_next_runs"
) as advance, patch("cron.scheduler.run_one_job") as run_one:
assert tick(verbose=False, sync=True, can_dispatch=lambda: False) == 0

Expand Down Expand Up @@ -1230,7 +1230,7 @@ def mock_run_job(job, *, defer_agent_teardown=None):
]

with patch("cron.scheduler.get_due_jobs", return_value=jobs), \
patch("cron.scheduler.advance_next_run"), \
patch("cron.scheduler.advance_next_runs"), \
patch("cron.scheduler.run_job", side_effect=mock_run_job), \
patch("cron.scheduler.save_job_output", return_value="/tmp/out.md"), \
patch("cron.scheduler._deliver_result", return_value=None), \
Expand Down Expand Up @@ -1275,7 +1275,7 @@ def mock_run_job(job, *, defer_agent_teardown=None):
]

with patch("cron.scheduler.get_due_jobs", return_value=jobs), \
patch("cron.scheduler.advance_next_run"), \
patch("cron.scheduler.advance_next_runs"), \
patch("cron.scheduler.run_job", side_effect=mock_run_job), \
patch("cron.scheduler.save_job_output", return_value="/tmp/out.md"), \
patch("cron.scheduler._deliver_result", return_value=None), \
Expand Down
2 changes: 1 addition & 1 deletion tests/cron/test_sessiondb_init_hang.py
Original file line number Diff line number Diff line change
Expand Up @@ -222,7 +222,7 @@ def test_guard_is_released_and_job_refires_after_sessiondb_hang(self, tmp_path,
side_effect=_session_db_executor(timeouts),
), \
patch.object(sched, "get_due_jobs", return_value=[job]), \
patch.object(sched, "advance_next_run"), \
patch.object(sched, "advance_next_runs"), \
patch.object(sched, "save_job_output", return_value="/tmp/out"), \
patch.object(sched, "mark_job_run"), \
patch.object(sched, "_deliver_result", return_value=None):
Expand Down
Loading