MCPcopy Create free account
hub / github.com/atomicdotdev/atomic / worker

Function worker

tools/provenance-eval/multi_session.py:50–96  ·  view source on GitHub ↗
(binary, root, endpoint, number, config, gate, stop_lock)

Source from the content-addressed store, hash-verified

48
49
50def 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
99def verify(f, sessions, config):

Callers

nothing calls this directly

Calls 6

hookMethod · 0.95
runMethod · 0.95
closeMethod · 0.80
getMethod · 0.65
decodeMethod · 0.45
appendMethod · 0.45

Tested by

no test coverage detected