fix(profiles): merge distributed cron jobs safely
(cherry picked from commit 291dca0055855abcec3bfbedb92d8890baeb833a)
This commit is contained in:
36
cron/jobs.py
36
cron/jobs.py
@@ -407,6 +407,15 @@ 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
|
||||
@@ -2022,6 +2031,33 @@ 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.
|
||||
|
||||
@@ -55,6 +55,22 @@ USER_OWNED_EXCLUDE: frozenset = frozenset({
|
||||
"local",
|
||||
})
|
||||
|
||||
# Profile distributions own cron definitions, not scheduler state. The runtime has
|
||||
# one canonical multi-record store; every sibling under cron/ is runtime data.
|
||||
_CRON_STORE_REL = ("cron", "jobs.json")
|
||||
|
||||
|
||||
def _is_distribution_runtime_path(parts: Tuple[str, ...]) -> bool:
|
||||
"""Runtime-owned entries nested under otherwise distribution-owned roots."""
|
||||
if len(parts) < 2:
|
||||
return False
|
||||
if parts[0] == "cron":
|
||||
return parts[:2] != _CRON_STORE_REL
|
||||
# Root-level dot entries under skills are Hermes bookkeeping (.hub,
|
||||
# .usage.json, curator state, bundled manifest, locks, archives, ...).
|
||||
return parts[0] == "skills" and len(parts) == 2 and parts[1].startswith(".")
|
||||
|
||||
|
||||
|
||||
class DistributionError(Exception):
|
||||
"""Raised for distribution install/update failures."""
|
||||
@@ -302,8 +318,7 @@ class InstallPlan:
|
||||
|
||||
|
||||
def _has_cron_jobs(staged: Path) -> bool:
|
||||
cron_dir = staged / "cron"
|
||||
return cron_dir.is_dir() and (any(cron_dir.rglob("*.json")) or any(cron_dir.rglob("*.yaml")))
|
||||
return staged.joinpath(*_CRON_STORE_REL).is_file()
|
||||
|
||||
|
||||
def plan_install(source: str, workdir: Path, override_name: Optional[str] = None) -> InstallPlan:
|
||||
@@ -350,7 +365,7 @@ def _owned_entries(staged: Path, manifest: DistributionManifest):
|
||||
# Path-aware allowlist: copy exactly the declared paths.
|
||||
for rel in explicit_owned:
|
||||
rel_parts = PurePosixPath(rel).parts
|
||||
if not rel_parts or rel_parts[0] in USER_OWNED_EXCLUDE:
|
||||
if not rel_parts or rel_parts[0] in USER_OWNED_EXCLUDE or _is_distribution_runtime_path(rel_parts):
|
||||
continue
|
||||
if ".." in rel_parts or PurePosixPath(rel).is_absolute():
|
||||
continue
|
||||
@@ -378,6 +393,44 @@ def _replace_entry(src: Path, dest: Path) -> None:
|
||||
shutil.copy2(src, dest)
|
||||
|
||||
|
||||
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
|
||||
|
||||
try:
|
||||
with tempfile.TemporaryDirectory(prefix="hermes_dist_cron_") as tmp:
|
||||
staged_store = Path(tmp) / "cron"
|
||||
staged_store.mkdir()
|
||||
shutil.copy2(src, staged_store / "jobs.json")
|
||||
with cron_jobs.use_cron_store(tmp):
|
||||
shipped = {
|
||||
job["id"]: job for job in cron_jobs.load_jobs()
|
||||
if isinstance(job, dict) and job.get("id")
|
||||
}
|
||||
|
||||
now = datetime.now(timezone.utc).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 cron_jobs.merge_job_definition(local, incoming))
|
||||
merged.extend(
|
||||
cron_jobs.merge_job_definition({"id": job_id, **new_state}, incoming)
|
||||
for job_id, incoming in shipped.items()
|
||||
)
|
||||
cron_jobs.save_jobs(merged)
|
||||
except (OSError, RuntimeError, ValueError) as exc:
|
||||
raise DistributionError(f"Could not merge cron jobs into {dest}: {exc}") from exc
|
||||
|
||||
|
||||
def _real_dir(base: Path, parts: Tuple[str, ...]) -> Path:
|
||||
"""Return ``base/parts`` as a chain of real directories.
|
||||
|
||||
@@ -409,22 +462,28 @@ def _is_container(path: Path) -> bool:
|
||||
return path.is_dir() and not any(p.is_file() for p in path.iterdir())
|
||||
|
||||
|
||||
def _merge_dir(src: Path, dest: Path) -> None:
|
||||
"""Replace only the roots *src* ships inside *dest*; a nested container
|
||||
(``skills/<category>``) is merged, not replaced, so sibling roots the user
|
||||
added under the same category survive."""
|
||||
def _merge_dir(src: Path, dest: Path, rel: Tuple[str, ...]) -> None:
|
||||
"""Merge authored roots while leaving runtime-owned nested state untouched."""
|
||||
for child in src.iterdir():
|
||||
if _is_container(child):
|
||||
_merge_dir(child, _real_dir(dest, (child.name,)))
|
||||
parts = (*rel, child.name)
|
||||
if _is_distribution_runtime_path(parts):
|
||||
continue
|
||||
if parts == _CRON_STORE_REL:
|
||||
_merge_cron_store(child, dest / child.name)
|
||||
elif _is_container(child):
|
||||
_merge_dir(child, _real_dir(dest, (child.name,)), parts)
|
||||
else:
|
||||
_replace_entry(child, dest / child.name)
|
||||
|
||||
|
||||
def _refuse_symlinked_containers(src: Path, dest: Path) -> None:
|
||||
def _refuse_symlinked_containers(src: Path, dest: Path, rel: Tuple[str, ...]) -> None:
|
||||
for child in src.iterdir():
|
||||
parts = (*rel, child.name)
|
||||
if _is_distribution_runtime_path(parts):
|
||||
continue
|
||||
if _is_container(child):
|
||||
_refuse_symlink(dest / child.name)
|
||||
_refuse_symlinked_containers(child, dest / child.name)
|
||||
_refuse_symlinked_containers(child, dest / child.name, parts)
|
||||
|
||||
|
||||
def _refuse_symlinked_targets(target: Path, entries) -> None:
|
||||
@@ -440,7 +499,7 @@ def _refuse_symlinked_targets(target: Path, entries) -> None:
|
||||
path = path / part
|
||||
_refuse_symlink(path)
|
||||
if src.is_dir() and len(rel_parts) == 1:
|
||||
_refuse_symlinked_containers(src, path)
|
||||
_refuse_symlinked_containers(src, path, rel_parts)
|
||||
|
||||
|
||||
def _copy_dist_payload(staged: Path, target: Path, manifest: DistributionManifest, preserve_config: bool) -> None:
|
||||
@@ -450,9 +509,9 @@ def _copy_dist_payload(staged: Path, target: Path, manifest: DistributionManifes
|
||||
``preserve_config`` is False (fresh install / ``--force-config``). ``.env.template`` lands
|
||||
as ``.env.EXAMPLE`` so it never shadows a real ``.env``.
|
||||
|
||||
A top-level owned directory (``skills/``, ``cron/``, ...) is a container of roots: only
|
||||
the roots the payload ships are replaced, so roots the user added (or that an older
|
||||
version shipped) survive an update or forced reinstall."""
|
||||
A top-level owned directory is merged per authored root. ``cron/jobs.json`` is
|
||||
special: it is one multi-record runtime store, so shipped definitions merge by job id
|
||||
instead of replacing the file."""
|
||||
target.mkdir(parents=True, exist_ok=True)
|
||||
entries = list(_owned_entries(staged, manifest))
|
||||
_refuse_symlinked_targets(target, entries)
|
||||
@@ -467,10 +526,13 @@ def _copy_dist_payload(staged: Path, target: Path, manifest: DistributionManifes
|
||||
if name == "config.yaml" and preserve_config and (target / "config.yaml").exists():
|
||||
continue
|
||||
if src.is_dir():
|
||||
_merge_dir(src, _real_dir(target, rel_parts))
|
||||
_merge_dir(src, _real_dir(target, rel_parts), rel_parts)
|
||||
continue
|
||||
parent = _real_dir(target, rel_parts[:-1])
|
||||
_replace_entry(src, parent / rel_parts[-1])
|
||||
if rel_parts == _CRON_STORE_REL:
|
||||
_merge_cron_store(src, parent / rel_parts[-1])
|
||||
else:
|
||||
_replace_entry(src, parent / rel_parts[-1])
|
||||
|
||||
# Emit .env.EXAMPLE from manifest if the staged tree didn't ship one
|
||||
if manifest.env_requires and not (target / ENV_EXAMPLE_FILENAME).exists():
|
||||
|
||||
Reference in New Issue
Block a user