refactor(cron): own the store-import loop in cron/job_definition.py

hermes_cli still reached into cron's private _jobs_lock and hand-built the
"created paused" record one function away from the module created so callers
never duplicate cron's schema. import_job_definitions() now holds the lock,
loads, merges and saves, and labels a merge ValueError with the job name;
_merge_cron_store keeps only the temp-store parse and the DistributionError
wrap, and takes the profile home instead of deriving it from dest.parent.parent.

Also drops the dead `and key != "repeat"` (repeat is reassigned right after)
and the comment that restated the module docstring.
This commit is contained in:
kshitijk4poor
2026-09-24 14:34:23 +05:30
committed by kshitij
parent bb0650e541
commit ca1ef789a8
2 changed files with 47 additions and 42 deletions

View File

@@ -6,11 +6,10 @@ into a live one (profile distributions) import this rather than duplicating the
"""
from typing import Any, Dict
from cron.jobs import _apply_schedule_update, parse_schedule
from cron.jobs import _apply_schedule_update, _jobs_lock, load_jobs, parse_schedule, save_jobs
from cron.quota_hold import clear_state as _clear_quota_hold
from hermes_time import now as _hermes_now
# Persisted fields authored by create_job rather than advanced by the scheduler.
# Cron owns this schema so import/update callers never need to duplicate it.
JOB_DEFINITION_FIELDS = frozenset({
"name", "prompt", "skills", "skill", "model", "provider", "base_url",
"script", "no_agent", "monitor_script", "monitor_url", "context_from",
@@ -24,10 +23,7 @@ def merge_job_definition(local: Dict[str, Any], authored: Dict[str, Any]) -> Dic
Raises ValueError when the authored schedule cannot be scheduled (unparseable string,
past one-shot for a live job)."""
merged = {
key: value for key, value in local.items()
if key not in JOB_DEFINITION_FIELDS and key != "repeat"
}
merged = {key: value for key, value in local.items() if key not in JOB_DEFINITION_FIELDS}
merged.update((key, authored[key]) for key in JOB_DEFINITION_FIELDS if key in authored)
merged["repeat"] = {
"completed": (local.get("repeat") or {}).get("completed", 0),
@@ -47,3 +43,35 @@ def merge_job_definition(local: Dict[str, Any], authored: Dict[str, Any]) -> Dic
else:
merged["next_run_at"] = None
return merged
def import_job_definitions(shipped: Dict[str, Dict[str, Any]], *, paused_reason: str) -> None:
"""Merge *shipped* (job id -> authored record) into the active store under its lock.
Unknown ids arrive with the marker set ``create_job(paused=True)`` writes; known ids keep
their scheduler state. Nothing is written when a record cannot be merged: the ValueError
is re-raised naming the job."""
now = _hermes_now().isoformat()
seed = {
"enabled": False,
"state": "paused",
"paused_at": now,
"paused_reason": paused_reason,
"created_at": now,
"next_run_at": None,
}
pending = dict(shipped)
with _jobs_lock():
merged = []
for local in load_jobs():
incoming = pending.pop(local.get("id"), None)
merged.append(local if incoming is None else _merge_or_name(local, incoming))
merged.extend(_merge_or_name({"id": job_id, **seed}, incoming) for job_id, incoming in pending.items())
save_jobs(merged)
def _merge_or_name(local: Dict[str, Any], authored: Dict[str, Any]) -> Dict[str, Any]:
try:
return merge_job_definition(local, authored)
except ValueError as exc:
raise ValueError(f"cron job {authored.get('name') or local.get('id')!r}: {exc}") from exc

View File

@@ -403,24 +403,15 @@ def _shipped_cron_store(entries: List[Tuple[Path, Tuple[str, ...]]]) -> Optional
return None
def _merge_cron_store(src: Path, dest: Path) -> None:
"""Merge a distribution's cron store by job id; new jobs arrive paused.
def _merge_cron_store(src: Path, home: Path) -> None:
"""Merge a distribution's cron store into profile *home* by job id; new jobs arrive paused.
Nothing is written when a shipped job cannot be scheduled (unparseable schedule, past
one-shot for a job the installer resumed); the error names the job."""
from hermes_time import now as _hermes_now
from cron import jobs as cron_jobs
from cron.job_definition import merge_job_definition
def merge(local: Dict[str, Any], incoming: Dict[str, Any]) -> Dict[str, Any]:
try:
return merge_job_definition(local, incoming)
except ValueError as exc:
raise DistributionError(
f"Could not merge cron job {incoming.get('name') or local.get('id')!r} into {dest}: {exc}"
) from exc
from cron.job_definition import import_job_definitions
dest = home.joinpath(*_CRON_STORE_REL)
try:
with tempfile.TemporaryDirectory(prefix="hermes_dist_cron_") as tmp:
staged_store = Path(tmp) / "cron"
@@ -431,27 +422,12 @@ def _merge_cron_store(src: Path, dest: Path) -> None:
job["id"]: job for job in cron_jobs.load_jobs()
if isinstance(job, dict) and job.get("id")
}
now = _hermes_now().isoformat()
new_state = {
"enabled": False,
"state": "paused",
"paused_at": now,
"paused_reason": "Installed from a profile distribution; review it, then resume.",
"created_at": now,
"next_run_at": None,
}
with cron_jobs.use_cron_store(dest.parent.parent), cron_jobs._jobs_lock():
merged = []
for local in cron_jobs.load_jobs():
incoming = shipped.pop(local.get("id"), None)
merged.append(local if incoming is None else merge(local, incoming))
merged.extend(
merge({"id": job_id, **new_state}, incoming)
for job_id, incoming in shipped.items()
)
cron_jobs.save_jobs(merged)
except RuntimeError as exc: # load_jobs: corrupt/unreadable store; OSError propagates as-is
with cron_jobs.use_cron_store(home):
import_job_definitions(
shipped, paused_reason="Installed from a profile distribution; review it, then resume.")
except (RuntimeError, ValueError) as exc:
# RuntimeError: load_jobs on a corrupt/unreadable store; ValueError: a job-labelled
# unschedulable definition. OSError propagates as-is like every other copy step.
raise DistributionError(f"Could not merge cron jobs into {dest}: {exc}") from exc
@@ -544,7 +520,8 @@ def _copy_dist_payload(staged: Path, target: Path, manifest: DistributionManifes
# (an unschedulable job), and rejecting before any file is replaced keeps the profile whole.
cron_store = _shipped_cron_store(entries)
if cron_store is not None:
_merge_cron_store(cron_store, _real_dir(target, _CRON_STORE_REL[:-1]) / _CRON_STORE_REL[-1])
_real_dir(target, _CRON_STORE_REL[:-1])
_merge_cron_store(cron_store, target)
for src, rel_parts in entries:
if rel_parts == _CRON_STORE_REL: