Skip to content
Closed
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
36 changes: 36 additions & 0 deletions cron/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
96 changes: 79 additions & 17 deletions hermes_cli/profile_distribution.py
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,22 @@
"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."""
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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:
Expand All @@ -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:
Expand All @@ -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)
Expand All @@ -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():
Expand Down
85 changes: 72 additions & 13 deletions tests/hermes_cli/test_profile_distribution.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
import stat
import subprocess
import sys
from datetime import datetime, timedelta, timezone
from pathlib import Path

import pytest
Expand Down Expand Up @@ -57,7 +58,7 @@ def _make_staging_dir(root: Path, name: str = "src", *, manifest: DistributionMa
contain after .git is removed).

Lays down a minimal but representative tree: SOUL.md, config.yaml,
mcp.json, one skill, one cron file, plus the distribution.yaml manifest.
mcp.json, one skill, one cron job, plus the distribution.yaml manifest.
"""
staged = root / f"staging_{name}"
staged.mkdir(parents=True, exist_ok=True)
Expand All @@ -69,8 +70,9 @@ def _make_staging_dir(root: Path, name: str = "src", *, manifest: DistributionMa
(staged / "skills" / "demo" / "SKILL.md").write_text(
"---\nname: demo\ndescription: test\n---\n# Demo skill\n"
)
(staged / "cron").mkdir(exist_ok=True)
(staged / "cron" / "daily.json").write_text('{"schedule": "0 9 * * *"}')
from cron.jobs import create_job, use_cron_store
with use_cron_store(staged):
create_job("daily task", "0 9 * * *", name="daily")

mf = manifest or DistributionManifest(name=name, version="0.1.0")
write_manifest(staged, mf)
Expand Down Expand Up @@ -310,30 +312,23 @@ def test_install_omitted_allowlist_copies_everything(self, profile_env):
"omitted distribution_owned must keep copying undeclared dirs"

def test_install_allowlist_supports_nested_paths(self, profile_env):
"""Documented nested entries like skills/research/ and cron/digest.json
must select exactly that subtree/file, not be silently dropped."""
"""Nested skill paths and the canonical cron store can be allowlisted exactly."""
mf = DistributionManifest(
name="nested",
version="0.1.0",
distribution_owned=["SOUL.md", "skills/research/", "cron/digest.json"],
distribution_owned=["SOUL.md", "skills/research/", "cron/jobs.json"],
)
staged = _make_staging_dir(profile_env, "nested", manifest=mf)
(staged / "skills" / "research").mkdir()
(staged / "skills" / "research" / "SKILL.md").write_text(
"---\nname: research\ndescription: r\n---\n# R\n"
)
(staged / "cron" / "digest.json").write_text('{"schedule": "0 8 * * *"}')

plan = install_distribution(str(staged), name="nested")
# Nested allowlisted paths are installed
assert (plan.target_dir / "skills" / "research" / "SKILL.md").exists()
assert (plan.target_dir / "cron" / "digest.json").exists()
# Sibling paths under the same parents are NOT dragged along
assert (plan.target_dir / "cron" / "jobs.json").exists()
assert not (plan.target_dir / "skills" / "demo").exists(), \
"skills/demo is not allowlisted and must not be copied"
assert not (plan.target_dir / "cron" / "daily.json").exists(), \
"cron/daily.json is not allowlisted and must not be copied"
# Unrelated top-level entries stay out too
assert not (plan.target_dir / "mcp.json").exists()

def test_update_respects_distribution_owned_allowlist(self, profile_env):
Expand Down Expand Up @@ -370,6 +365,34 @@ def test_update_respects_distribution_owned_allowlist(self, profile_env):
assert not (plan.target_dir / "new_config.toml").exists(), \
"new_config.toml should not be copied (not in distribution_owned)"

def test_install_pauses_shipped_cron_jobs_and_skips_runtime_state(self, profile_env):
"""Distribution cron definitions are inert until the installer explicitly resumes them."""
from cron.jobs import get_due_jobs, is_job_runnable, list_jobs, update_job, use_cron_store

staged = _make_staging_dir(profile_env, "cron_src")
with use_cron_store(staged):
shipped = list_jobs(include_disabled=True)[0]
update_job(shipped["id"], {
"next_run_at": (datetime.now(timezone.utc) - timedelta(days=2)).isoformat()
})
(staged / "cron" / "future-runtime.bin").write_text("runtime", encoding="utf-8")
(staged / "skills" / ".future-runtime").write_text("runtime", encoding="utf-8")
(staged / "skills" / "demo" / ".authored-hidden").write_text("keep", encoding="utf-8")

plan = install_distribution(str(staged), name="cron_paused")

with use_cron_store(plan.target_dir):
installed = {job["id"]: job for job in list_jobs(include_disabled=True)}
due = get_due_jobs()
job = installed[shipped["id"]]
assert not is_job_runnable(job)
assert job["state"] == "paused"
assert due == []
assert not (plan.target_dir / "cron" / "future-runtime.bin").exists()
assert not (plan.target_dir / "skills" / ".future-runtime").exists()
assert (plan.target_dir / "skills" / "demo" / ".authored-hidden").read_text() == "keep"


def test_install_rejects_non_distribution_directory(self, profile_env, tmp_path):
bogus = tmp_path / "bogus_dir"
bogus.mkdir()
Expand Down Expand Up @@ -437,6 +460,42 @@ def test_update_and_force_install_merge_owned_dirs_per_root(self, profile_env):
assert (custom / "SKILL.md").read_text(encoding="utf-8") == "custom skill\n"
assert (plan.target_dir / "cron" / "mine.json").exists()

def test_update_merges_cron_jobs_without_losing_local_state(self, profile_env):
"""Updating one shipped definition cannot replace the profile's whole cron store."""
from cron.jobs import create_job, list_jobs, pause_job, resume_job, update_job, use_cron_store

staged = _make_staging_dir(profile_env, "cron_update")
with use_cron_store(staged):
shipped = list_jobs(include_disabled=True)[0]
plan = install_distribution(str(staged), name="cron_merge")

with use_cron_store(plan.target_dir):
resume_job(shipped["id"])
mine = create_job("water the plants", "0 9 * * *", name="mine")
pause_job(mine["id"], reason="my choice")
old_next_run = {
job["id"]: job for job in list_jobs(include_disabled=True)
}[shipped["id"]]["next_run_at"]

with use_cron_store(staged):
update_job(shipped["id"], {
"prompt": "updated upstream definition",
"schedule": "every 7d",
})

update_distribution("cron_merge")

with use_cron_store(plan.target_dir):
jobs = {job["id"]: job for job in list_jobs(include_disabled=True)}
assert mine["id"] in jobs
assert jobs[mine["id"]]["state"] == "paused"
assert jobs[mine["id"]]["paused_reason"] == "my choice"
assert jobs[shipped["id"]]["prompt"] == "updated upstream definition"
assert jobs[shipped["id"]]["schedule"]["minutes"] == 7 * 24 * 60
assert jobs[shipped["id"]]["next_run_at"] != old_next_run
assert jobs[shipped["id"]]["enabled"] is True


def test_update_refuses_symlinked_owned_container(self, profile_env):
staged = _make_staging_dir(profile_env, "src")
plan = install_distribution(str(staged), name="link_safe")
Expand Down
Loading