| 48 | |
| 49 | |
| 50 | def worker(binary, root, endpoint, number, config, gate, stop_lock): |
| 51 | session = f"multi-agent-{number}" |
| 52 | base.SESSION = session |
| 53 | f = base.Fixture(binary, pathlib.Path(root)) |
| 54 | f.endpoint = endpoint |
| 55 | reads, writes, stops, turn_times, failures = [], [], [], [], [] |
| 56 | gate.wait(timeout=60) |
| 57 | started = time.perf_counter() |
| 58 | try: |
| 59 | for turn in range(1, config["turns"] + 1): |
| 60 | turn_start = time.perf_counter() |
| 61 | f.hook("user-prompt-submit", {"session_id": session, "prompt": f"Review intent and source, turn {turn}"}) |
| 62 | for event in range(config["events"]): |
| 63 | for read in range(config["reads_per_event"]): |
| 64 | command = ["intent", "list", "--json"] if read % 2 == 0 else ["intent", "show", config["intent_id"]] |
| 65 | ts = time.perf_counter() |
| 66 | try: |
| 67 | result = f.run(*command) |
| 68 | text = result.stdout.decode() |
| 69 | if read % 2 == 0: |
| 70 | assert {x["id"] for x in json.loads(text)} == set(config["intent_ids"]), "intent list changed" |
| 71 | else: |
| 72 | assert config["intent_id"] in text, "intent content missing" |
| 73 | reads.append({"ms": (time.perf_counter() - ts) * 1000, "ok": True, "command": command}) |
| 74 | except Exception as error: |
| 75 | reads.append({"ms": (time.perf_counter() - ts) * 1000, "ok": False, "command": command, "error": str(error)}) |
| 76 | ts = time.perf_counter() |
| 77 | f.hook("post-tool", base.tool_payload(event, 1024, turn)) |
| 78 | writes.append((time.perf_counter() - ts) * 1000) |
| 79 | if config.get("think_ms", 0): |
| 80 | time.sleep(config["think_ms"] / 1000) |
| 81 | ts = time.perf_counter() |
| 82 | try: |
| 83 | with stop_lock if config["serialize_stop"] else contextlib.nullcontext(): |
| 84 | acquired = time.perf_counter() |
| 85 | f.hook("stop", {"session_id": session, "reason": "end_turn", "response": f"Review completed {turn}"}) |
| 86 | stops.append({"ms": (time.perf_counter() - ts) * 1000, "queue_ms": (acquired - ts) * 1000, "ok": True}) |
| 87 | except Exception as error: |
| 88 | stops.append({"ms": (time.perf_counter() - ts) * 1000, "ok": False, "error": str(error)}) |
| 89 | break |
| 90 | turn_times.append((time.perf_counter() - turn_start) * 1000) |
| 91 | except Exception as error: |
| 92 | failures.append(str(error)) |
| 93 | wall = time.perf_counter() - started |
| 94 | f.owner_log.close() |
| 95 | return {"session": session, "wall_s": wall, "reads": reads, "writes_ms": writes, |
| 96 | "stops": stops, "turn_ms": turn_times, "failures": failures} |
| 97 | |
| 98 | |
| 99 | def verify(f, sessions, config): |