diff --git a/bin/box-ctl.py b/bin/box-ctl.py index 4a796eb..4ecbbb3 100755 --- a/bin/box-ctl.py +++ b/bin/box-ctl.py @@ -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 (job JSON on stdin) job-delete [--force] job-trigger + job-result (result JSON on stdin; records RESULT + advances chain) + job-status (dispatch/result/chain state) + job-next [--success|--fail] (dry-run: what would box dispatch next) + job-chain [--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 ") act_job_trigger(rest[0]) + elif action == "job-result": + if len(rest) != 1: + fail("BAD_NAME", "usage: job-result (result JSON on stdin)") + act_job_result(rest[0]) + elif action == "job-status": + if len(rest) != 1: + fail("BAD_NAME", "usage: job-status ") + act_job_status(rest[0]) + elif action == "job-next": + if not rest or len(rest) > 2: + fail("BAD_NAME", "usage: job-next [--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 [--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 --thread [...]") 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 (result JSON on stdin)") + act_job_result(args[0]) + elif sub == "status": + if len(args) != 1: + fail("BAD_ARGS", "usage: job status ") + act_job_status(args[0]) + elif sub == "next": + if not args or len(args) > 2: + fail("BAD_ARGS", "usage: job next [--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 [--on-failure]") + act_job_chain(args[0], args[1], on_failure=("--on-failure" in args[2:])) elif action == "dm-log": limit = 50 if rest: