Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 4 additions & 8 deletions pipeline/governor.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,14 +43,10 @@ def project(state, corpus_rows, conflict_rate=None):
# a shard that hit the 24h wall is not zero progress -- it was salvaged
done[sh["stage"]] = done.get(sh["stage"], 0) + sh.get("salvaged", 0)
elif sh.get("status") == "submitted" and sh.get("job"):
# in-flight: count live completionStats as "paid for" so remaining-work
# estimates don't double-count rows the job has already processed
try:
import batch
cs = batch.get_job(sh["job"]).get("completionStats", {})
done[sh["stage"]] = done.get(sh["stage"], 0) + int(cs.get("successfulCount", 0))
except Exception:
pass
# Unharvested jobs are absent from the spend ledger. Keep their rows in
# the future-cost estimate until harvest records the actual cost; live
# completionStats must not make reserved work disappear from projection.
continue
rem_a = max(0, corpus_rows - done["pass_a"])
rem_b = max(0, corpus_rows - done["pass_b"])
rem_c = max(0, int(corpus_rows * conflict_rate) - done["pass_c"])
Expand Down
110 changes: 82 additions & 28 deletions pipeline/session.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
python3 pipeline/session.py status what is in flight / done / pending
python3 pipeline/session.py submit N build+submit up to N shards (governor-gated)
python3 pipeline/session.py harvest collect finished jobs -> BQ + ledger
python3 pipeline/session.py run harvest, then submit to keep MAX_INFLIGHT busy
python3 pipeline/session.py run harvest, then submit up to MAX_INFLIGHT new shards

State lives in state/shards.json, mirrored to GCS after every transition, so a
session that dies mid-run resumes rather than restarts.
Expand All @@ -20,6 +20,7 @@
# Measured 2026-08-31: Pro and Flash draw from SEPARATE batch throughput pools --
# running Pass B alongside Pass A left Pro at +11.9 rows/min (unchanged). So each
# stage gets its own in-flight budget rather than sharing one.
# Legacy name: this caps new submissions per `run`; stage caps below bound live jobs.
MAX_INFLIGHT = int(os.environ.get("PW_MAX_INFLIGHT", "3"))
INFLIGHT_BY_STAGE = {"pass_a": int(os.environ.get("PW_INFLIGHT_A", "4")),
"pass_b": int(os.environ.get("PW_INFLIGHT_B", "6")),
Expand Down Expand Up @@ -128,42 +129,95 @@ def cmd_harvest():
if changed: save(st); ledger.sync()
return st

def cmd_submit(n=None):
def cmd_submit(n=1):
if isinstance(n, bool) or not isinstance(n, int) or n < 0:
raise ValueError("submission limit must be a non-negative integer")
if n == 0:
print("submission limit is 0; nothing to submit")
return True

st = load()
gov = governor.project(st, st["corpus"])
action, notes = governor.decide(gov)
if action == "HALT":
print("GOVERNOR HALT:", notes[0]); return
# Fail closed: CUT describes changes an operator must apply before spending.
# Treating it as advisory would submit the uncut workload that exceeded target.
if action not in ("PROCEED", "STRETCH"):
print(f"GOVERNOR {action}: submission blocked")
for note in notes:
print(" " + note)
return False
print(f"governor: {action} (projected THB {gov['projected_total_thb']})")
todo = []
submitted = 0
for stage, cap in INFLIGHT_BY_STAGE.items():
infl = sum(1 for s in st["shards"] if s["status"] == "submitted" and s["stage"] == stage)
free = max(0, cap - infl)
pend = [s for s in st["shards"] if s["status"] == "pending" and s["stage"] == stage]
if pend and free == 0:
print(f" {stage}: no slots ({infl}/{cap} in flight)")
todo += pend[:free]
if not todo:
return
for s in todo:
rows = rows_for(s["id"])
_, kept, missing, attempt_id = batch.build(
rows, s["id"], s["model"], s["thinking"]
)
if kept == 0:
print(f" {s['id']}: no rows with patches, skipping"); continue
job = batch.submit(s["id"], s["model"], attempt_id)
s["status"] = "submitted"; s["job"] = job["name"]; s["built"] = kept; s["missing"] = missing
s["attempt_id"] = attempt_id
s["submitted_at"] = datetime.datetime.now(datetime.timezone.utc).isoformat()
print(f" {s['id']}: submitted {kept} rows ({missing} missing patches) job={job['name'].split('/')[-1]}")
save(st)
for s in pend:
if submitted == n or free == 0:
break
rows = rows_for(s["id"])
_, kept, missing, attempt_id = batch.build(
rows, s["id"], s["model"], s["thinking"]
)
if kept == 0:
print(f" {s['id']}: no rows with patches, skipping")
continue
job = batch.submit(s["id"], s["model"], attempt_id)
s["status"] = "submitted"
s["job"] = job["name"]
s["built"] = kept
s["missing"] = missing
s["attempt_id"] = attempt_id
s["submitted_at"] = datetime.datetime.now(datetime.timezone.utc).isoformat()
job_id = job["name"].split("/")[-1]
print(
f" {s['id']}: submitted {kept} rows "
f"({missing} missing patches) job={job_id}"
)
save(st)
submitted += 1
free -= 1
if submitted == n:
break
return True

if __name__ == "__main__":
cmd = sys.argv[1] if len(sys.argv) > 1 else "status"
if cmd == "status": cmd_status()
elif cmd == "harvest": cmd_harvest(); cmd_status()
elif cmd == "submit": cmd_submit(int(sys.argv[2]) if len(sys.argv)>2 else 1)
def _submit_limit(argv):
if len(argv) > 3:
raise ValueError("usage: session.py submit [non-negative integer]")
try:
limit = int(argv[2]) if len(argv) > 2 else 1
except ValueError as exc:
raise ValueError("submission limit must be a non-negative integer") from exc
if limit < 0:
raise ValueError("submission limit must be a non-negative integer")
return limit

def main(argv=None):
argv = sys.argv if argv is None else argv
cmd = argv[1] if len(argv) > 1 else "status"
if cmd == "status":
cmd_status()
elif cmd == "harvest":
cmd_harvest(); cmd_status()
elif cmd == "submit":
try:
limit = _submit_limit(argv)
except ValueError as e:
print(f"submit: {e}", file=sys.stderr)
return 2
allowed = cmd_submit(limit)
if not allowed:
return 3
elif cmd == "run":
cmd_harvest(); cmd_submit(MAX_INFLIGHT); cmd_status()
else: print(__doc__)
cmd_harvest(); allowed = cmd_submit(MAX_INFLIGHT); cmd_status()
if not allowed:
return 3
else:
print(__doc__)
return 2
return 0

if __name__ == "__main__":
raise SystemExit(main())
Loading