feat(box): recursive job workflow -- job-result, job-status, job-next, job-chain
This commit is contained in:
+309
@@ -25,6 +25,7 @@ import os
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
@@ -619,6 +620,263 @@ def act_job_trigger(name):
|
||||
out(True, name=name, triggered=True)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Recursive workflow: box -> agent -> box -> agent ...
|
||||
# The box stays in the loop: every RESULT recorded here advances the chain
|
||||
# through the box (shared dedupe state with response-harvester.py).
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
CHAINED_JOBS_FILE = BIN / "chained-jobs.json"
|
||||
|
||||
|
||||
def _chain_has(job_id):
|
||||
try:
|
||||
if CHAINED_JOBS_FILE.exists():
|
||||
return job_id in json.loads(CHAINED_JOBS_FILE.read_text())
|
||||
except Exception:
|
||||
pass
|
||||
return False
|
||||
|
||||
|
||||
def _chain_mark(job_id):
|
||||
try:
|
||||
c = []
|
||||
if CHAINED_JOBS_FILE.exists():
|
||||
c = json.loads(CHAINED_JOBS_FILE.read_text())
|
||||
if job_id not in c:
|
||||
c.append(job_id)
|
||||
c = c[-1000:]
|
||||
CHAINED_JOBS_FILE.write_text(json.dumps(c))
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _job_name_from_id(job_id):
|
||||
m = re.match(r"^(.*)-(\d{8}-\d{6}-[a-f0-9]{8})$", str(job_id))
|
||||
if m:
|
||||
return m.group(1)
|
||||
parts = str(job_id).split("-")
|
||||
if len(parts) < 3:
|
||||
return None
|
||||
return "-".join(parts[:-2])
|
||||
|
||||
|
||||
def _chain_target(job_name, success):
|
||||
"""Return (next_job, via_field) for a completed job, or (None, None)."""
|
||||
p = JOBS_DIR / f"{job_name}.json"
|
||||
if not p.exists():
|
||||
return None, None
|
||||
try:
|
||||
cfg = json.loads(p.read_text())
|
||||
except Exception:
|
||||
return None, None
|
||||
if success:
|
||||
nxt = cfg.get("on_success") or cfg.get("chain_next")
|
||||
via = "on_success" if cfg.get("on_success") else "chain_next"
|
||||
else:
|
||||
nxt = cfg.get("on_failure")
|
||||
via = "on_failure"
|
||||
if nxt and isinstance(nxt, str) and (JOBS_DIR / f"{nxt}.json").exists():
|
||||
return nxt, via
|
||||
return None, None
|
||||
|
||||
|
||||
def _dispatch_chained(next_job, prev_job_id, prev_result):
|
||||
"""Dispatch the next job in a chain. Returns (ok, err)."""
|
||||
env = os.environ.copy()
|
||||
env["CHAIN_PREV_JOB_ID"] = prev_job_id
|
||||
env["CHAIN_PREV_RESULT"] = (prev_result or "")[:1000]
|
||||
try:
|
||||
r = subprocess.run([sys.executable, str(DISPATCHER), next_job],
|
||||
capture_output=True, text=True, timeout=300, env=env)
|
||||
except Exception as e:
|
||||
return False, f"dispatch exception: {e}"
|
||||
if r.returncode != 0:
|
||||
return False, (r.stderr or r.stdout).strip()[-500:]
|
||||
return True, None
|
||||
|
||||
|
||||
def act_job_result(job_id):
|
||||
"""Agent -> box: record a RESULT and advance the chain (the recursion step)."""
|
||||
raw = sys.stdin.read()
|
||||
try:
|
||||
data = json.loads(raw)
|
||||
except Exception as e:
|
||||
fail("INVALID_RESULT", f"stdin is not valid JSON: {e}")
|
||||
agent = data.get("agent")
|
||||
if agent not in VALID_AGENTS:
|
||||
fail("BAD_NAME", f"agent must be one of {sorted(VALID_AGENTS)}")
|
||||
success = data.get("success", True)
|
||||
if not isinstance(success, bool):
|
||||
fail("INVALID_RESULT", "success must be a boolean")
|
||||
summary = data.get("summary", "")
|
||||
if not isinstance(summary, str) or len(summary) > 2000:
|
||||
fail("INVALID_RESULT", "summary must be a string of 0-2000 chars")
|
||||
rec = {
|
||||
"ts": utcnow(),
|
||||
"type": "job_result",
|
||||
"job_id": job_id,
|
||||
"agent": agent,
|
||||
"success": success,
|
||||
"result_snippet": summary[:300],
|
||||
"thread_id": None,
|
||||
"msg_id": None,
|
||||
"source": "box-ctl",
|
||||
}
|
||||
try:
|
||||
with open(JOB_LOG, "a") as f:
|
||||
f.write(json.dumps(rec) + "\n")
|
||||
except Exception as e:
|
||||
fail("LOG_FAILED", f"could not append job-log.jsonl: {e}")
|
||||
# recursion: the box dispatches the next step
|
||||
job_name = _job_name_from_id(job_id)
|
||||
chained_to = None
|
||||
chain_via = None
|
||||
chain_dispatched = False
|
||||
chain_error = None
|
||||
if job_name and not _chain_has(job_id):
|
||||
_chain_mark(job_id)
|
||||
nxt, via = _chain_target(job_name, success)
|
||||
if nxt:
|
||||
try:
|
||||
ccfg = json.loads((JOBS_DIR / f"{job_name}.json").read_text())
|
||||
except Exception:
|
||||
ccfg = {}
|
||||
try:
|
||||
step_delay = int(ccfg.get("step_delay", 0) or 0)
|
||||
except Exception:
|
||||
step_delay = 0
|
||||
if step_delay > 0:
|
||||
time.sleep(min(step_delay, 30))
|
||||
ok, err = _dispatch_chained(nxt, job_id, summary)
|
||||
chain_dispatched = ok
|
||||
chain_error = err
|
||||
if ok:
|
||||
chained_to = nxt
|
||||
chain_via = via
|
||||
audit("job-result", job_id)
|
||||
out(True, job_id=job_id, agent=agent, success=success, recorded=True,
|
||||
chained_to=chained_to, chain_via=chain_via,
|
||||
chain_dispatched=chain_dispatched, chain_error=chain_error)
|
||||
|
||||
|
||||
def act_job_status(name):
|
||||
"""Box-side view of a job's run state: dispatches, results, chain wiring."""
|
||||
events = []
|
||||
try:
|
||||
if JOB_LOG.exists():
|
||||
for line in JOB_LOG.read_text().splitlines()[-2000:]:
|
||||
try:
|
||||
e = json.loads(line)
|
||||
except Exception:
|
||||
continue
|
||||
jn = e.get("job_name")
|
||||
jid = str(e.get("job_id", ""))
|
||||
if jn == name or jid == name or jid.startswith(name + "-"):
|
||||
events.append(e)
|
||||
except Exception as e:
|
||||
fail("LOG_FAILED", f"could not read job-log.jsonl: {e}")
|
||||
dispatches = [e for e in events if e.get("type") == "job_dispatched"]
|
||||
results = [e for e in events if e.get("type") == "job_result"]
|
||||
failures = [e for e in events if e.get("type") in ("job_failed", "job_timeout")]
|
||||
last_result = results[-1] if results else None
|
||||
state = "pending"
|
||||
if last_result:
|
||||
state = "resulted"
|
||||
elif dispatches:
|
||||
state = "dispatched"
|
||||
chain = {}
|
||||
p = JOBS_DIR / f"{name}.json"
|
||||
if p.exists():
|
||||
try:
|
||||
cfg = json.loads(p.read_text())
|
||||
for k in ("chain_next", "on_success", "on_failure"):
|
||||
if cfg.get(k):
|
||||
chain[k] = cfg[k]
|
||||
except Exception:
|
||||
pass
|
||||
audit("job-status", name)
|
||||
out(True, name=name, state=state,
|
||||
dispatches=len(dispatches), results=len(results), failures=len(failures),
|
||||
last_dispatch=dispatches[-1].get("ts") if dispatches else None,
|
||||
last_result=({
|
||||
"ts": last_result.get("ts"),
|
||||
"job_id": last_result.get("job_id"),
|
||||
"agent": last_result.get("agent"),
|
||||
"success": last_result.get("success"),
|
||||
"snippet": (last_result.get("result_snippet") or "")[:200],
|
||||
} if last_result else None),
|
||||
chain=chain)
|
||||
|
||||
|
||||
def act_job_next(job_id, success=None):
|
||||
"""Dry-run: what WOULD the box dispatch next for this job_id? No dispatch."""
|
||||
job_name = _job_name_from_id(job_id)
|
||||
if not job_name:
|
||||
fail("BAD_NAME", f"cannot parse job name from job_id: {job_id}")
|
||||
if success is None:
|
||||
success = True
|
||||
try:
|
||||
if JOB_LOG.exists():
|
||||
for line in reversed(JOB_LOG.read_text().splitlines()[-2000:]):
|
||||
try:
|
||||
e = json.loads(line)
|
||||
except Exception:
|
||||
continue
|
||||
if e.get("type") == "job_result" and str(e.get("job_id")) == job_id:
|
||||
success = bool(e.get("success", True))
|
||||
break
|
||||
except Exception:
|
||||
pass
|
||||
nxt, via = _chain_target(job_name, success)
|
||||
already = _chain_has(job_id)
|
||||
audit("job-next", job_id)
|
||||
out(True, job_id=job_id, job_name=job_name, success=success,
|
||||
would_dispatch=bool(nxt) and not already,
|
||||
next_job=nxt, via=via, already_chained=already)
|
||||
|
||||
|
||||
def act_job_chain(frm, to, on_failure=False):
|
||||
"""Wire recursion: set chain_next (or on_failure) from one job to another."""
|
||||
check_name(frm)
|
||||
check_name(to)
|
||||
if frm == to:
|
||||
fail("INVALID_JOB", "a job cannot chain to itself")
|
||||
fp = JOBS_DIR / f"{frm}.json"
|
||||
if not fp.exists():
|
||||
fail("NOT_FOUND", f"no such job: {frm}")
|
||||
if not (JOBS_DIR / f"{to}.json").exists():
|
||||
fail("NOT_FOUND", f"no such job: {to}")
|
||||
try:
|
||||
cfg = json.loads(fp.read_text())
|
||||
except Exception as e:
|
||||
fail("INVALID_JOB", f"job JSON unreadable: {e}")
|
||||
field = "on_failure" if on_failure else "chain_next"
|
||||
cfg[field] = to
|
||||
# cycle guard
|
||||
seen = {frm}
|
||||
cur = to
|
||||
while cur:
|
||||
if cur in seen:
|
||||
fail("INVALID_JOB", f"chain would create a cycle at {cur!r}")
|
||||
seen.add(cur)
|
||||
try:
|
||||
c = json.loads((JOBS_DIR / f"{cur}.json").read_text())
|
||||
except Exception:
|
||||
break
|
||||
cur = c.get("chain_next") or c.get("on_success")
|
||||
fp.write_text(json.dumps(cfg, indent=2) + "\n")
|
||||
r = _git("add", f"jobs/{frm}.json")
|
||||
if r.returncode != 0:
|
||||
fail("DISPATCH_FAILED", f"git add failed: {r.stderr.strip()}")
|
||||
r = _git("-c", "user.name=box-ctl", "-c", "user.email=box-ctl@bl.local",
|
||||
"commit", "-m", f"Chain job {frm} -> {to} via box-ctl")
|
||||
if r.returncode != 0 and "nothing to commit" not in (r.stdout + r.stderr):
|
||||
fail("DISPATCH_FAILED", f"git commit failed: {r.stderr.strip()}")
|
||||
audit("job-chain", f"{frm}->{to}")
|
||||
out(True, **{"from": frm, "to": to, "field": field, "chained": True})
|
||||
|
||||
|
||||
def _valid_sidechat_name(name):
|
||||
# Sidechat names may contain spaces ("646 tasks"); reject only
|
||||
# control characters and enforce a sane length. dm.py resolves
|
||||
@@ -1197,6 +1455,10 @@ job actions:
|
||||
job-put <name> (job JSON on stdin)
|
||||
job-delete <name> [--force]
|
||||
job-trigger <name>
|
||||
job-result <job-id> (result JSON on stdin; records RESULT + advances chain)
|
||||
job-status <name-or-id> (dispatch/result/chain state)
|
||||
job-next <job-id> [--success|--fail] (dry-run: what would box dispatch next)
|
||||
job-chain <from> <to> [--on-failure] (wire chain_next between jobs)
|
||||
|
||||
variable actions:
|
||||
vars-list
|
||||
@@ -1350,6 +1612,27 @@ def main(argv):
|
||||
if len(rest) != 1:
|
||||
fail("BAD_NAME", "usage: job-trigger <name>")
|
||||
act_job_trigger(rest[0])
|
||||
elif action == "job-result":
|
||||
if len(rest) != 1:
|
||||
fail("BAD_NAME", "usage: job-result <job-id> (result JSON on stdin)")
|
||||
act_job_result(rest[0])
|
||||
elif action == "job-status":
|
||||
if len(rest) != 1:
|
||||
fail("BAD_NAME", "usage: job-status <name-or-job-id>")
|
||||
act_job_status(rest[0])
|
||||
elif action == "job-next":
|
||||
if not rest or len(rest) > 2:
|
||||
fail("BAD_NAME", "usage: job-next <job-id> [--success|--fail]")
|
||||
s = None
|
||||
if "--success" in rest[1:]:
|
||||
s = True
|
||||
if "--fail" in rest[1:]:
|
||||
s = False
|
||||
act_job_next(rest[0], s)
|
||||
elif action == "job-chain":
|
||||
if len(rest) < 2 or len(rest) > 3:
|
||||
fail("BAD_NAME", "usage: job-chain <from-job> <to-job> [--on-failure]")
|
||||
act_job_chain(rest[0], rest[1], on_failure=("--on-failure" in rest[2:]))
|
||||
elif action == "vars-list":
|
||||
act_vars_list()
|
||||
elif action == "vars-get":
|
||||
@@ -1611,6 +1894,32 @@ def main(argv):
|
||||
if op in ("pin", "unpin", "archive", "unarchive", "rename") and not thread_id:
|
||||
fail("BAD_NAME", f"usage: {action} --agent <agent> --thread <id> [...]")
|
||||
act_thread(op, agent, thread=thread_id, title=title, limit=limit, confirm=confirm)
|
||||
elif action == "job":
|
||||
if not rest or rest[0] not in ("result", "status", "next", "chain"):
|
||||
fail("BAD_NAME", "usage: job result|status|next|chain [...]")
|
||||
sub = rest[0]
|
||||
args = rest[1:]
|
||||
if sub == "result":
|
||||
if len(args) != 1:
|
||||
fail("BAD_ARGS", "usage: job result <job-id> (result JSON on stdin)")
|
||||
act_job_result(args[0])
|
||||
elif sub == "status":
|
||||
if len(args) != 1:
|
||||
fail("BAD_ARGS", "usage: job status <name-or-job-id>")
|
||||
act_job_status(args[0])
|
||||
elif sub == "next":
|
||||
if not args or len(args) > 2:
|
||||
fail("BAD_ARGS", "usage: job next <job-id> [--success|--fail]")
|
||||
s = None
|
||||
if "--success" in args[1:]:
|
||||
s = True
|
||||
if "--fail" in args[1:]:
|
||||
s = False
|
||||
act_job_next(args[0], s)
|
||||
elif sub == "chain":
|
||||
if len(args) < 2 or len(args) > 3:
|
||||
fail("BAD_ARGS", "usage: job chain <from-job> <to-job> [--on-failure]")
|
||||
act_job_chain(args[0], args[1], on_failure=("--on-failure" in args[2:]))
|
||||
elif action == "dm-log":
|
||||
limit = 50
|
||||
if rest:
|
||||
|
||||
Reference in New Issue
Block a user