refactor(cron): move job-definition merge schema into cron/job_definition.py
cron/jobs.py is already ~3400 lines; JOB_DEFINITION_FIELDS and merge_job_definition are only used by importers of a foreign cron store (profile distributions), so they live in a small dedicated module instead of growing the scheduler file. profile_distribution imports the new module directly; jobs.py is left byte-identical to main (no re-export shim). Refs #120823 Co-authored-by: John Paul Soliva <soliva.johnpaul@icloud.com> Co-authored-by: JoaoMarcos44 <joaomarcosdias444@gmail.com>
This commit is contained in:
45
cron/job_definition.py
Normal file
45
cron/job_definition.py
Normal file
@@ -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
|
||||
36
cron/jobs.py
36
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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user