diff --git a/harness/mix.py b/harness/mix.py index 5d9a5d7..f9434de 100644 --- a/harness/mix.py +++ b/harness/mix.py @@ -77,3 +77,11 @@ def accepted_task_counts(logs_dir=None): if k: out[k] = out.get(k, 0) + 1 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} diff --git a/harness/pipeline.py b/harness/pipeline.py index 5781a48..9ae9148 100644 --- a/harness/pipeline.py +++ b/harness/pipeline.py @@ -129,16 +129,18 @@ class Pipeline: done = {(r["task"], r["attempt"]) for r in rows} ev0 = {r["task"]: self.evaluate(r) for r in rows if r["attempt"] == 0} tasks = trajectories.accepted_tasks() - cands = [t for t in tasks if (t, 0) not in done and (t, 0) not in self.inflight] - if cands: # first attempts: the kind with the biggest deficit against the target mix goes first - counts = {} - for r in rows: - 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 = {} # accepted trajectories (and runs in flight) per kind + for r in rows: + 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 + 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()) 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:]))) @@ -146,7 +148,8 @@ class Pipeline: return best, 0 for t in tasks: # second attempt: first failed, or accepted without a repair 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)) return t, 1 return None @@ -358,6 +361,26 @@ trajectory runs {len(rows)}, accepted trajectories {len(accepted)} (acceptance { 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(): load_env(os.path.join(ROOT, ".env")) ap = argparse.ArgumentParser() @@ -379,7 +402,9 @@ def main(): break time.sleep(30) os.makedirs(OUT, exist_ok=True) + guard = acquire_controller_lock() Pipeline(a).run() + guard.close() if __name__ == "__main__": diff --git a/harness/trainset.py b/harness/trainset.py index 005d693..1e14476 100644 --- a/harness/trainset.py +++ b/harness/trainset.py @@ -279,7 +279,7 @@ def run_balanced(part, parts, deadline): if blocked: print("kinds skipped (low acceptance):", sorted(blocked), flush=True) 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: time.sleep(120) # every kind has a backlog: the trajectories are the slower side continue