| 132 | |
| 133 | |
| 134 | def measure(binary, workers, config, template): |
| 135 | with tempfile.TemporaryDirectory(prefix="atomic-multi-session-") as temp: |
| 136 | root = pathlib.Path(temp) |
| 137 | f = PreparedFixture(binary, root, template) |
| 138 | try: |
| 139 | # Seed vault context before opening the owner. |
| 140 | f.start() |
| 141 | config = dict(config) |
| 142 | config["intent_ids"] = sorted(x["id"] for x in json.loads(f.run("intent", "list", "--json").stdout)) |
| 143 | assert len(config["intent_ids"]) == 32 |
| 144 | config["intent_id"] = config["intent_ids"][0] |
| 145 | assert config["intent_id"] in f.run("intent", "show", config["intent_id"]).stdout.decode() |
| 146 | for i in range(workers): |
| 147 | f.hook("session-start", {"session_id": f"multi-agent-{i}"}) |
| 148 | ctx = multiprocessing.get_context("spawn") |
| 149 | with ctx.Manager() as manager: |
| 150 | gate = manager.Barrier(workers + 1) |
| 151 | stop_lock = manager.Lock() |
| 152 | with concurrent.futures.ProcessPoolExecutor(max_workers=workers, mp_context=ctx) as pool: |
| 153 | pending = [pool.submit(worker, str(binary), str(root), f.endpoint, i, config, gate, stop_lock) for i in range(workers)] |
| 154 | gate.wait(timeout=60) |
| 155 | rows = [future.result(timeout=180) for future in pending] |
| 156 | result = {"workers": workers, "config": config, "sessions": rows, |
| 157 | "wall_s": max(r["wall_s"] for r in rows)} |
| 158 | result["verified"], result["verification_errors"] = [], [] |
| 159 | for row in rows: |
| 160 | try: |
| 161 | result["verified"].extend(verify(f, [row], config)) |
| 162 | except Exception as error: |
| 163 | result["verification_errors"].append({"session": row["session"], "error": str(error)}) |
| 164 | return result |
| 165 | finally: |
| 166 | logs = f.close() |
| 167 | if "result" in locals(): |
| 168 | result["owner_log"] = logs |
| 169 | |
| 170 | |
| 171 | def summarize(rows): |