feat(swarm): add box swarm prune --stale-hours command and archive stale swarms
This commit is contained in:
+49
-2
@@ -1606,7 +1606,7 @@ def act_thread(op, agent, thread=None, title=None, limit=None, confirm=False):
|
|||||||
SWARM_FILE = NETVM_ROOT / "swarms.json"
|
SWARM_FILE = NETVM_ROOT / "swarms.json"
|
||||||
SWARM_ID_RE = re.compile(r"^[a-z0-9-]{1,64}$")
|
SWARM_ID_RE = re.compile(r"^[a-z0-9-]{1,64}$")
|
||||||
SWARM_AGENT_RE = re.compile(r"^[A-Za-z0-9-]{1,64}$")
|
SWARM_AGENT_RE = re.compile(r"^[A-Za-z0-9-]{1,64}$")
|
||||||
SWARM_OPS = ("spawn", "list", "status", "attach", "report", "results", "kill")
|
SWARM_OPS = ("spawn", "list", "status", "attach", "report", "results", "kill", "prune")
|
||||||
SWARM_MAX_COUNT = 50
|
SWARM_MAX_COUNT = 50
|
||||||
SWARM_MAX_TASK_LEN = 2000
|
SWARM_MAX_TASK_LEN = 2000
|
||||||
|
|
||||||
@@ -1817,6 +1817,42 @@ def act_swarm_results(sid):
|
|||||||
counts=_swarm_counts(swarm), results=results)
|
counts=_swarm_counts(swarm), results=results)
|
||||||
|
|
||||||
|
|
||||||
|
def act_swarm_prune(stale_hours=6, confirm=False):
|
||||||
|
import datetime
|
||||||
|
swarms = _swarm_load()
|
||||||
|
cutoff = datetime.datetime.now(datetime.timezone.utc) - datetime.timedelta(hours=stale_hours)
|
||||||
|
pruned = []
|
||||||
|
|
||||||
|
for sid, swarm in swarms.items():
|
||||||
|
st = swarm.get("status")
|
||||||
|
up_ts = swarm.get("updated_ts") or swarm.get("created_ts")
|
||||||
|
if not up_ts:
|
||||||
|
continue
|
||||||
|
try:
|
||||||
|
dt = datetime.datetime.fromisoformat(up_ts.replace("Z", "+00:00"))
|
||||||
|
except Exception:
|
||||||
|
continue
|
||||||
|
|
||||||
|
if dt < cutoff and st in ("partial", "completed", "killed"):
|
||||||
|
pruned.append(sid)
|
||||||
|
|
||||||
|
if not confirm:
|
||||||
|
fail("CONFIRM_REQUIRED",
|
||||||
|
f"prune transitions {len(pruned)} stale swarms to 'archived' -- pass --confirm to proceed",
|
||||||
|
{"stale_hours": stale_hours, "matching_count": len(pruned), "sample": pruned[:10]})
|
||||||
|
|
||||||
|
now = utcnow()
|
||||||
|
archived_count = 0
|
||||||
|
for sid in pruned:
|
||||||
|
swarms[sid]["status"] = "archived"
|
||||||
|
swarms[sid]["updated_ts"] = now
|
||||||
|
archived_count += 1
|
||||||
|
|
||||||
|
_swarm_save(swarms)
|
||||||
|
audit("swarm-prune", f"stale_hours={stale_hours} count={archived_count}")
|
||||||
|
out(True, archived_count=archived_count, stale_hours=stale_hours)
|
||||||
|
|
||||||
|
|
||||||
def act_swarm_kill(sid, confirm=False):
|
def act_swarm_kill(sid, confirm=False):
|
||||||
swarms = _swarm_load()
|
swarms = _swarm_load()
|
||||||
swarm = _swarm_get(swarms, sid)
|
swarm = _swarm_get(swarms, sid)
|
||||||
@@ -2956,7 +2992,7 @@ def main(argv):
|
|||||||
act_thread(op, agent, thread=thread_id, title=title, limit=limit, confirm=confirm)
|
act_thread(op, agent, thread=thread_id, title=title, limit=limit, confirm=confirm)
|
||||||
elif action in ("swarm-spawn", "swarm-list", "swarm-status",
|
elif action in ("swarm-spawn", "swarm-list", "swarm-status",
|
||||||
"swarm-attach", "swarm-report", "swarm-results",
|
"swarm-attach", "swarm-report", "swarm-results",
|
||||||
"swarm-kill"):
|
"swarm-kill", "swarm-prune"):
|
||||||
op = action[len("swarm-"):]
|
op = action[len("swarm-"):]
|
||||||
if op == "spawn":
|
if op == "spawn":
|
||||||
if len(rest) < 2:
|
if len(rest) < 2:
|
||||||
@@ -3047,6 +3083,17 @@ def main(argv):
|
|||||||
if not args or len(args) > 2:
|
if not args or len(args) > 2:
|
||||||
fail("BAD_ARGS", "usage: swarm kill <swarm-id> --confirm")
|
fail("BAD_ARGS", "usage: swarm kill <swarm-id> --confirm")
|
||||||
act_swarm_kill(args[0], confirm=("--confirm" in args[1:]))
|
act_swarm_kill(args[0], confirm=("--confirm" in args[1:]))
|
||||||
|
elif sub == "prune":
|
||||||
|
stale_h = 6
|
||||||
|
confirm = "--confirm" in args
|
||||||
|
if "--stale-hours" in args:
|
||||||
|
i = args.index("--stale-hours")
|
||||||
|
if i + 1 < len(args):
|
||||||
|
try:
|
||||||
|
stale_h = float(args[i+1])
|
||||||
|
except ValueError:
|
||||||
|
pass
|
||||||
|
act_swarm_prune(stale_hours=stale_h, confirm=confirm)
|
||||||
elif action == "job":
|
elif action == "job":
|
||||||
if not rest or rest[0] not in ("result", "status", "next", "chain"):
|
if not rest or rest[0] not in ("result", "status", "next", "chain"):
|
||||||
fail("BAD_NAME", "usage: job result|status|next|chain [...]")
|
fail("BAD_NAME", "usage: job result|status|next|chain [...]")
|
||||||
|
|||||||
Reference in New Issue
Block a user