diff --git a/cron/job_definition.py b/cron/job_definition.py new file mode 100644 index 0000000000..1f8b84df16 --- /dev/null +++ b/cron/job_definition.py @@ -0,0 +1,45 @@ +"""Job-definition schema shared by importers of a foreign cron store. + +Cron owns which persisted fields are *authored* (create_job) versus *advanced by the +scheduler* (next_run_at, state, run counters...). Callers merging an authored store +into a live one (profile distributions) import this rather than duplicating the list. +""" +from typing import Any, Dict + +from cron.jobs import _apply_schedule_update +from cron.quota_hold import clear_state as _clear_quota_hold + +# 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", + "schedule", "schedule_display", "deliver", "origin", "enabled_toolsets", + "workdir", "attach_to_session", "reasoning_effort", "failure_deliver", +}) + + +def merge_job_definition(local: Dict[str, Any], authored: Dict[str, Any]) -> Dict[str, Any]: + """Refresh authored fields while preserving this store's scheduler-owned state.""" + merged = { + key: value for key, value in local.items() + if key not in JOB_DEFINITION_FIELDS and key != "repeat" + } + 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), + "times": (authored.get("repeat") or {}).get("times"), + } + + if local.get("schedule") != merged.get("schedule"): + merged.pop("pending_slot", None) + _clear_quota_hold(merged) + if merged.get("enabled", True) and merged.get("state") != "paused": + _apply_schedule_update( + merged, + {"schedule": merged["schedule"], "schedule_display": merged.get("schedule_display")}, + str(merged.get("id") or "imported job"), + ) + else: + merged["next_run_at"] = None + return merged diff --git a/cron/jobs.py b/cron/jobs.py index 14f213d483..20e3ae5a8f 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -407,15 +407,6 @@ def fire_claim_fence(job_id: str, *, expected_owner: str): # update could leak ``../escape``/absolute/nested values into output writes/deletes. _IMMUTABLE_JOB_FIELDS = frozenset({"id"}) -# 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", - "schedule", "schedule_display", "deliver", "origin", "enabled_toolsets", - "workdir", "attach_to_session", "reasoning_effort", "failure_deliver", -}) - def _job_output_dir(job_id: str) -> Path: """Resolve a job's output directory, rejecting any path-escape attempt (``..``, absolute @@ -2031,33 +2022,6 @@ def _fill_missing_next_run(updated: Dict[str, Any]) -> None: updated["next_run_at"] = next_run -def merge_job_definition(local: Dict[str, Any], authored: Dict[str, Any]) -> Dict[str, Any]: - """Refresh authored fields while preserving this store's scheduler-owned state.""" - merged = { - key: value for key, value in local.items() - if key not in JOB_DEFINITION_FIELDS and key != "repeat" - } - 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), - "times": (authored.get("repeat") or {}).get("times"), - } - - if local.get("schedule") != merged.get("schedule"): - merged.pop("pending_slot", None) - from cron.quota_hold import clear_state as _clear_quota_hold - _clear_quota_hold(merged) - if merged.get("enabled", True) and merged.get("state") != "paused": - _apply_schedule_update( - merged, - {"schedule": merged["schedule"], "schedule_display": merged.get("schedule_display")}, - str(merged.get("id") or "imported job"), - ) - else: - merged["next_run_at"] = None - return merged - - def update_job(job_id: str, updates: Dict[str, Any]) -> Optional[Dict[str, Any]]: """Update a job by ID, refreshing derived schedule fields when needed.""" # ``id`` is a path component under OUTPUT_DIR — changing it would leak path-escape values. diff --git a/hermes_cli/profile_distribution.py b/hermes_cli/profile_distribution.py index 7d700f18f2..c0d93966a6 100644 --- a/hermes_cli/profile_distribution.py +++ b/hermes_cli/profile_distribution.py @@ -396,6 +396,7 @@ def _replace_entry(src: Path, dest: Path) -> None: def _merge_cron_store(src: Path, dest: Path) -> None: """Merge a distribution's cron store by job id; new jobs arrive paused.""" from cron import jobs as cron_jobs + from cron.job_definition import merge_job_definition try: with tempfile.TemporaryDirectory(prefix="hermes_dist_cron_") as tmp: @@ -421,9 +422,9 @@ def _merge_cron_store(src: Path, dest: Path) -> None: merged = [] for local in cron_jobs.load_jobs(): incoming = shipped.pop(local.get("id"), None) - merged.append(local if incoming is None else cron_jobs.merge_job_definition(local, incoming)) + merged.append(local if incoming is None else merge_job_definition(local, incoming)) merged.extend( - cron_jobs.merge_job_definition({"id": job_id, **new_state}, incoming) + merge_job_definition({"id": job_id, **new_state}, incoming) for job_id, incoming in shipped.items() ) cron_jobs.save_jobs(merged)