Only kinds below their target share are generated and run; one-controller flock guard
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -77,3 +77,11 @@ def accepted_task_counts(logs_dir=None):
|
|||||||
if k:
|
if k:
|
||||||
out[k] = out.get(k, 0) + 1
|
out[k] = out.get(k, 0) + 1
|
||||||
return out
|
return out
|
||||||
|
|
||||||
|
|
||||||
|
def below_target(counts, share=TYPE_SHARE):
|
||||||
|
"""Kinds whose share of `counts` is below the target for the next item (deficit > 0).
|
||||||
|
Since 2026-10-05 (Kral + Opus): only these kinds are generated and run; CLAS and FUNC wait until they are at or
|
||||||
|
below their target share."""
|
||||||
|
total, ss = sum(counts.values()) + 1, sum(share.values())
|
||||||
|
return {k for k, v in share.items() if total * v / ss - counts.get(k, 0) > 0}
|
||||||
|
|||||||
@@ -129,16 +129,18 @@ class Pipeline:
|
|||||||
done = {(r["task"], r["attempt"]) for r in rows}
|
done = {(r["task"], r["attempt"]) for r in rows}
|
||||||
ev0 = {r["task"]: self.evaluate(r) for r in rows if r["attempt"] == 0}
|
ev0 = {r["task"]: self.evaluate(r) for r in rows if r["attempt"] == 0}
|
||||||
tasks = trajectories.accepted_tasks()
|
tasks = trajectories.accepted_tasks()
|
||||||
cands = [t for t in tasks if (t, 0) not in done and (t, 0) not in self.inflight]
|
counts = {} # accepted trajectories (and runs in flight) per kind
|
||||||
if cands: # first attempts: the kind with the biggest deficit against the target mix goes first
|
for r in rows:
|
||||||
counts = {}
|
if self.evaluate(r)["accepted"]:
|
||||||
for r in rows:
|
k = mix.kind_of_task_dir(r["task"])
|
||||||
if self.evaluate(r)["accepted"]:
|
|
||||||
k = mix.kind_of_task_dir(r["task"])
|
|
||||||
counts[k] = counts.get(k, 0) + 1
|
|
||||||
for (t, _a) in self.inflight:
|
|
||||||
k = mix.kind_of_task_dir(t)
|
|
||||||
counts[k] = counts.get(k, 0) + 1
|
counts[k] = counts.get(k, 0) + 1
|
||||||
|
for (t, _a) in self.inflight:
|
||||||
|
k = mix.kind_of_task_dir(t)
|
||||||
|
counts[k] = counts.get(k, 0) + 1
|
||||||
|
below = mix.below_target(counts) # only kinds below their target share run (Kral + Opus 2026-10-05)
|
||||||
|
cands = [t for t in tasks if (t, 0) not in done and (t, 0) not in self.inflight
|
||||||
|
and mix.kind_of_task_dir(t) in below]
|
||||||
|
if cands: # first attempts: the kind with the biggest deficit against the target mix goes first
|
||||||
total, tot_share = sum(counts.values()) + 1, sum(mix.TYPE_SHARE.values())
|
total, tot_share = sum(counts.values()) + 1, sum(mix.TYPE_SHARE.values())
|
||||||
best = max(cands, key=lambda t: (total * mix.TYPE_SHARE.get(mix.kind_of_task_dir(t), 0) / tot_share
|
best = max(cands, key=lambda t: (total * mix.TYPE_SHARE.get(mix.kind_of_task_dir(t), 0) / tot_share
|
||||||
- counts.get(mix.kind_of_task_dir(t), 0), -int(t[1:])))
|
- counts.get(mix.kind_of_task_dir(t), 0), -int(t[1:])))
|
||||||
@@ -146,7 +148,8 @@ class Pipeline:
|
|||||||
return best, 0
|
return best, 0
|
||||||
for t in tasks: # second attempt: first failed, or accepted without a repair
|
for t in tasks: # second attempt: first failed, or accepted without a repair
|
||||||
e = ev0.get(t)
|
e = ev0.get(t)
|
||||||
if e and (t, 1) not in done and (t, 1) not in self.inflight and (not e["accepted"] or not e["repair"]):
|
if e and (t, 1) not in done and (t, 1) not in self.inflight and (not e["accepted"] or not e["repair"]) \
|
||||||
|
and mix.kind_of_task_dir(t) in below:
|
||||||
self.inflight.add((t, 1))
|
self.inflight.add((t, 1))
|
||||||
return t, 1
|
return t, 1
|
||||||
return None
|
return None
|
||||||
@@ -358,6 +361,26 @@ trajectory runs {len(rows)}, accepted trajectories {len(accepted)} (acceptance {
|
|||||||
log("pipeline ended:", text.replace("\n", " | "))
|
log("pipeline ended:", text.replace("\n", " | "))
|
||||||
|
|
||||||
|
|
||||||
|
def acquire_controller_lock():
|
||||||
|
"""One controller only. An exclusive flock on runs/pipeline/controller.lock lives as long as the process (the
|
||||||
|
kernel drops it when the process dies, so a crash never leaves a stale lock); the pid is written for people."""
|
||||||
|
import fcntl
|
||||||
|
os.makedirs(PIPE, exist_ok=True)
|
||||||
|
path = os.path.join(PIPE, "controller.lock")
|
||||||
|
f = open(path, "a+")
|
||||||
|
try:
|
||||||
|
fcntl.flock(f, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||||
|
except OSError:
|
||||||
|
f.seek(0)
|
||||||
|
sys.exit("a controller is already running (pid %s, %s); not starting a second one" % (
|
||||||
|
f.read().strip() or "?", path))
|
||||||
|
f.seek(0)
|
||||||
|
f.truncate()
|
||||||
|
f.write("%d\n" % os.getpid())
|
||||||
|
f.flush()
|
||||||
|
return f
|
||||||
|
|
||||||
|
|
||||||
def main():
|
def main():
|
||||||
load_env(os.path.join(ROOT, ".env"))
|
load_env(os.path.join(ROOT, ".env"))
|
||||||
ap = argparse.ArgumentParser()
|
ap = argparse.ArgumentParser()
|
||||||
@@ -379,7 +402,9 @@ def main():
|
|||||||
break
|
break
|
||||||
time.sleep(30)
|
time.sleep(30)
|
||||||
os.makedirs(OUT, exist_ok=True)
|
os.makedirs(OUT, exist_ok=True)
|
||||||
|
guard = acquire_controller_lock()
|
||||||
Pipeline(a).run()
|
Pipeline(a).run()
|
||||||
|
guard.close()
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
|
|||||||
@@ -279,7 +279,7 @@ def run_balanced(part, parts, deadline):
|
|||||||
if blocked:
|
if blocked:
|
||||||
print("kinds skipped (low acceptance):", sorted(blocked), flush=True)
|
print("kinds skipped (low acceptance):", sorted(blocked), flush=True)
|
||||||
waiting = {k for k, n in backlog_by_kind().items() if n > BAL_KIND_BACKLOG}
|
waiting = {k for k, n in backlog_by_kind().items() if n > BAL_KIND_BACKLOG}
|
||||||
allowed = set(mix.TYPE_SHARE) - blocked - waiting
|
allowed = (set(mix.TYPE_SHARE) - blocked - waiting) & mix.below_target(counts) # only kinds below their share
|
||||||
if not allowed:
|
if not allowed:
|
||||||
time.sleep(120) # every kind has a backlog: the trajectories are the slower side
|
time.sleep(120) # every kind has a backlog: the trajectories are the slower side
|
||||||
continue
|
continue
|
||||||
|
|||||||
Reference in New Issue
Block a user