From 66d69194bdd946e93a9fdae2a8cda78daf6294ef Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Mon, 7 Sep 2026 06:10:04 +0000 Subject: [PATCH] refactor(eval): consolidate duplicated harness logic after the review round Simplification pass over the branch. Behavior-preserving throughout; three reviewers, nine findings applied, two skipped. The review-artifact block in _run_hidden_oracle was unreachable. That function runs only in run_arm's non-review branch, while the directory it probes for is created only in the review branch, and each sandbox serves exactly one arm - so `review_artifact.parent.is_dir()` could never be true. It was added an hour earlier to make the hidden oracle resolve the moved artifact; the oracle never runs for review tasks, so the guard was dead on arrival. Deleting it also removes the duplication it had with the verify-command wiring. EXCLUDED_ERROR_KINDS is now one definition. runner.py and comparator_reuse.py each carried the same six-member frozenset, kept in sync by a comment. Only one direction is possible: runner already imports from comparator_reuse, so the reverse import fails at module-init with a circular-import error. That is now stated where the alias lives, so nobody tries it the other way. ensure_task_graph and prefetch_next_graph shared ten keyword parameters, passed through two call sites and forwarded whole between them. They now take a GraphBuildEnv, mirroring TaskCellContext, which already bundles per-cell state in this file. Its ready_keys() replaces an inline four-set union at the call site. Smaller consolidations: _sha256_file's hand-rolled chunk loop becomes hashlib.file_digest (3.11+, already used in runner_artifacts); _copy_owner_only reuses task_assets._write_all and COPY_CHUNK_BYTES instead of repeating the short-write retry; its stat-then-open existence check becomes the O_EXCL failure it was already relying on, which is atomic rather than merely narrow; and runner_environment reads the digest through comparator_reuse.current_runtime_digest instead of re-parsing the environment variable. Three test docstrings summarised the branch's own history ("the branch's core speedup", "the regression that produced fifteen runs") rather than the invariant under test. Rewritten to state the constraint, which is what survives the merge. Repaired the indentation left behind by the outage-streak edit and flattened the prefetch dispatch from three nested conditionals to one. Skipped: consolidating comparator_reuse._real_directory onto proposer_sandbox's same-named helper - they differ, the sandbox one rejects any symlink in the resolved path while this one checks only the leaf, so sharing it would tighten behavior rather than preserve it. That needs a decision about which policy the reuse path wants, not a simplification. 589 eval tests pass, ruff clean, 29 workflow contract tests pass. Unrelated and pre-existing: two test_model_gateway.py failures, and test_process_control.py::test_timeout_kills_term_ignoring_descendants_before_they_write, which is a TERM-to-KILL timing flake (passes 2 of 3 in isolation) in a file this branch does not touch. --- eval/tests/test_review_scoring.py | 6 +- eval/tests/test_runner_hardening.py | 2 +- eval/tests/test_workflow_bench.py | 52 ++++- eval/tests/test_workflow_bench_sessions.py | 2 +- eval/workflow_bench/comparator_reuse.py | 32 ++- eval/workflow_bench/evolve.py | 3 +- eval/workflow_bench/runner.py | 242 +++++++++------------ 7 files changed, 170 insertions(+), 169 deletions(-) diff --git a/eval/tests/test_review_scoring.py b/eval/tests/test_review_scoring.py index 65bac2a5f..eb105a57a 100644 --- a/eval/tests/test_review_scoring.py +++ b/eval/tests/test_review_scoring.py @@ -330,9 +330,9 @@ def test_clean_control_rewards_an_empty_approval_and_penalizes_noise(): def test_parse_review_output_names_the_actual_failure(tmp_path: Path): """One message per cause. - Folding these together is how a sandbox that made the artifact impossible - to write read for fifteen runs as an encoding fault: every cell reported - "not valid UTF-8 JSON" for a file the agent was never able to create. + Folding these together makes a sandbox that renders the artifact impossible + to write indistinguishable from an encoding fault: every cell reports "not + valid UTF-8 JSON" for a file the agent was never able to create. """ missing = tmp_path / "never-written.json" diff --git a/eval/tests/test_runner_hardening.py b/eval/tests/test_runner_hardening.py index 277c26b92..425109a94 100644 --- a/eval/tests/test_runner_hardening.py +++ b/eval/tests/test_runner_hardening.py @@ -1005,7 +1005,7 @@ def test_progress_line_reports_the_numbers_a_real_run_measured(): def test_review_artifact_is_mounted_as_a_writable_directory_outside_the_workspace(tmp_path): - """The regression that produced fifteen runs of empty evidence. + """A writable file inside a read-only directory is not a writable path. A writable FILE inside a read-only directory is not writable to anything that writes atomically. The Write tool creates `.tmp..` diff --git a/eval/tests/test_workflow_bench.py b/eval/tests/test_workflow_bench.py index b91ea780f..2bd18214b 100644 --- a/eval/tests/test_workflow_bench.py +++ b/eval/tests/test_workflow_bench.py @@ -13,6 +13,7 @@ import yaml from workflow_bench.runner import ( aggregate, + GraphBuildEnv, broken_incumbent_arms, build_parser, infra_error_record, @@ -596,16 +597,18 @@ def test_prefetch_next_graph_runs_ensure_on_a_background_thread(monkeypatch): task={"id": "review-b"}, binding={"repo_identity": "/repo", "resolved_sha": "bbb"}, graph_key=("/repo", "bbb"), - trees=Path("/tmp"), - task_asset_cache=None, - claude_bin="claude", - bwrap_bin="bwrap", - sandbox_backend="bwrap", - runtime_mounts=(), - clone_templates={}, - clone_template_errors={}, - graph_snapshots={}, - graph_snapshot_errors={}, + env=GraphBuildEnv( + trees=Path("/tmp"), + task_asset_cache=None, + claude_bin="claude", + bwrap_bin="bwrap", + sandbox_backend="bwrap", + runtime_mounts=(), + clone_templates={}, + clone_template_errors={}, + graph_snapshots={}, + graph_snapshot_errors={}, + ), cancel_event=cancel, ) assert job.key == ("/repo", "bbb") @@ -632,3 +635,32 @@ def test_a_reused_resolution_does_not_count_as_this_sweeps_health(): mixed = aggregate([record(resolved=True, reused=True), record(resolved=True)]) assert mixed["resolved_fresh"] == 1 assert broken_incumbent_arms({"t": {"review": mixed}}, {"review"}) == [] + + +def test_graph_build_env_ready_keys_covers_successes_and_failures(): + """A key that failed is attempted, not pending. + + next_graph_prefetch_target skips keys already in ready_keys. If a failed + build were omitted, the sweep would prefetch it again every iteration and + pay a full clone and offline index each time for a build that cannot + succeed. + """ + + env = GraphBuildEnv( + trees=Path("/tmp"), + task_asset_cache=None, + claude_bin="claude", + bwrap_bin="bwrap", + sandbox_backend="bwrap", + runtime_mounts=(), + clone_templates={("/repo", "aaa"): (Path("/tmp/a"), "aaa")}, + clone_template_errors={("/repo", "bbb"): OSError("clone failed")}, + graph_snapshots={("/repo", "ccc"): object()}, + graph_snapshot_errors={("/repo", "ddd"): OSError("index failed")}, + ) + assert env.ready_keys() == { + ("/repo", "aaa"), + ("/repo", "bbb"), + ("/repo", "ccc"), + ("/repo", "ddd"), + } diff --git a/eval/tests/test_workflow_bench_sessions.py b/eval/tests/test_workflow_bench_sessions.py index 8e5743519..267dba8b8 100644 --- a/eval/tests/test_workflow_bench_sessions.py +++ b/eval/tests/test_workflow_bench_sessions.py @@ -1381,7 +1381,7 @@ def test_copy_isolated_tree_does_not_share_git_objects_or_refs(tmp_path): def test_run_cell_uses_the_clone_template_instead_of_recloning(tmp_path, monkeypatch): - """The branch's core speedup, which had no coverage at all. + """run_cell must copy the template, never re-clone. run_cell takes the clone-template branch on essentially every multi-cell sweep: it copies a pre-sanitized template rather than paying `git clone diff --git a/eval/workflow_bench/comparator_reuse.py b/eval/workflow_bench/comparator_reuse.py index d39867e7c..3878f8fef 100644 --- a/eval/workflow_bench/comparator_reuse.py +++ b/eval/workflow_bench/comparator_reuse.py @@ -29,6 +29,7 @@ from .evolution import CANDIDATE_ARMS, EVIDENCE_MAX_AGE_DAYS from .proposer_sandbox import SandboxError from .runner_sessions import MAX_TRANSCRIPT_BYTES, PARENT_EVENT_STREAM_SOURCE from .runtime_mounts import CE_ARMS +from .task_assets import COPY_CHUNK_BYTES, _write_all REUSABLE_COMPARATOR_ARMS = frozenset( { @@ -378,34 +379,29 @@ def _regular_file(path: Path, *, label: str) -> Path: def _copy_owner_only(source: Path, destination: Path) -> None: - if destination.exists() or destination.is_symlink(): - raise SandboxError(f"reuse destination already exists: {destination}") - descriptor = os.open( - destination, - os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0), - 0o600, - ) + # O_CREAT|O_EXCL is the existence check, and unlike a stat beforehand it is + # atomic: a file appearing between check and open cannot slip through. + try: + descriptor = os.open( + destination, + os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0), + 0o600, + ) + except FileExistsError as exc: + raise SandboxError(f"reuse destination already exists: {destination}") from exc try: os.fchmod(descriptor, 0o600) with open(source, "rb") as handle: while True: - chunk = handle.read(1024 * 1024) + chunk = handle.read(COPY_CHUNK_BYTES) if not chunk: break - view = memoryview(chunk) - while view: - written = os.write(descriptor, view) - if written <= 0: - raise OSError(f"short write while copying {source}") - view = view[written:] + _write_all(descriptor, chunk) os.fsync(descriptor) finally: os.close(descriptor) def _sha256_file(path: Path) -> str: - digest = hashlib.sha256() with open(path, "rb") as handle: - for chunk in iter(lambda: handle.read(1024 * 1024), b""): - digest.update(chunk) - return digest.hexdigest() + return hashlib.file_digest(handle, "sha256").hexdigest() diff --git a/eval/workflow_bench/evolve.py b/eval/workflow_bench/evolve.py index e76c31218..db8f54df9 100644 --- a/eval/workflow_bench/evolve.py +++ b/eval/workflow_bench/evolve.py @@ -48,6 +48,7 @@ import yaml from . import runner from . import runner_sessions +from .comparator_reuse import current_runtime_digest from .model_gateway import ( ANTHROPIC_API_KEY_ENV, attach_openai_gateway, @@ -1036,7 +1037,7 @@ def runner_environment(args: argparse.Namespace) -> dict[str, str]: # process_control replaces the child environment wholesale, so a digest the # workflow exported reaches the runner only if it is forwarded here. Without # this the runner stamps no runtime_digest and the reuse lock never engages. - runtime_digest = os.environ.get("RUNTIME_DIGEST", "").strip() + runtime_digest = current_runtime_digest() if runtime_digest: env["RUNTIME_DIGEST"] = runtime_digest if args.auth_token: diff --git a/eval/workflow_bench/runner.py b/eval/workflow_bench/runner.py index 033b11b3b..70b1e476f 100644 --- a/eval/workflow_bench/runner.py +++ b/eval/workflow_bench/runner.py @@ -58,6 +58,7 @@ from typing import Any import yaml from .comparator_reuse import ( + REUSE_EXCLUDED_ERROR_KINDS, ComparatorReuseExpectation, TaskReuseBinding, current_runtime_digest, @@ -346,15 +347,6 @@ def _run_hidden_oracle( # credited patch is captured. oracle_mount = f"{SANDBOX_WORKSPACE}/{mount_name}" oracle_env[ORACLE_ENV_VAR] = str(stage_root) if host_unsafe else oracle_mount - review_artifact = review_output_path(sandbox, REVIEW_OUTPUT) - review_mounts: tuple[ReadOnlyMount, ...] = () - if review_artifact.parent.is_dir(): - oracle_env[REVIEW_OUTPUT_ENV_VAR] = sandbox.host_text( - f"{SANDBOX_REVIEW_OUTPUT}/{REVIEW_OUTPUT}" - ) - review_mounts = ( - ReadOnlyMount(source=review_artifact.parent, target=SANDBOX_REVIEW_OUTPUT), - ) passed, _output = _verification_outcome( run_verify( snapshot.command, @@ -363,9 +355,9 @@ def _run_hidden_oracle( command_prefix=sandbox.command_prefix_for( read_only_workspace=True, unshare_network=True, - extra_read_only_mounts=(*review_mounts,) + extra_read_only_mounts=() if host_unsafe - else (ReadOnlyMount(source=stage_root, target=oracle_mount), *review_mounts), + else (ReadOnlyMount(source=stage_root, target=oracle_mount),), ), env=oracle_env, require_pid_namespace=getattr(sandbox, "require_pid_namespace", True), @@ -768,9 +760,11 @@ CHURN_FIELDS = ("diff_files", "diff_insertions", "diff_deletions") # Rows where the session (or the harness) died carry no measured evidence and # must not skew efficiency medians or resolve denominators. verify-failed and # skill-not-invoked rows DO count: those sessions ran and spent real tokens. -EXCLUDED_ERROR_KINDS = frozenset( - {"session-error", "infra-error", "evidence-unverified", "cleanup-failure", "review-evidence-invalid", "cancelled"} -) +# One definition, in comparator_reuse: reuse eligibility and aggregate +# exclusion must never drift apart. The dependency only runs this way - +# comparator_reuse importing back from runner is a circular import. +EXCLUDED_ERROR_KINDS = REUSE_EXCLUDED_ERROR_KINDS + # A sustained upstream outage shows up as a run of session/infra/cleanup # failures. (cleanup-failure overwrites the primary error_kind, so a @@ -1862,56 +1856,79 @@ def next_graph_prefetch_target( return None +@dataclass(frozen=True) +class GraphBuildEnv: + """Per-sweep state every graph build shares, and the caches it fills. + + The four dicts are the sweep's memo of what has already been built, keyed by + (repo, sha). They are mutable by design and are written by both the sweep + thread and the prefetch thread, which is safe only because a build is + started for a key exactly once and joined before that key is read. + """ + + trees: Path + task_asset_cache: TaskAssetCache + claude_bin: Path | str + bwrap_bin: Path | str + sandbox_backend: str + runtime_mounts: Sequence[ReadOnlyMount] + clone_templates: dict[tuple[str, str], tuple[Path, str]] + clone_template_errors: dict[tuple[str, str], BaseException] + graph_snapshots: dict[tuple[str, str], SanitizedGraphSnapshot] + graph_snapshot_errors: dict[tuple[str, str], BaseException] + + def ready_keys(self) -> set[tuple[str, str]]: + """Keys whose build has already been attempted, successfully or not.""" + + return ( + set(self.clone_templates) + | set(self.clone_template_errors) + | set(self.graph_snapshots) + | set(self.graph_snapshot_errors) + ) + + def ensure_task_graph( *, task: Mapping[str, Any], repo: Path, task_sha: str, graph_key: tuple[str, str], - trees: Path, - task_asset_cache: TaskAssetCache, - claude_bin: Path | str, - bwrap_bin: Path | str, - sandbox_backend: str, - runtime_mounts: Sequence[ReadOnlyMount], - clone_templates: dict[tuple[str, str], tuple[Path, str]], - clone_template_errors: dict[tuple[str, str], BaseException], - graph_snapshots: dict[tuple[str, str], SanitizedGraphSnapshot], - graph_snapshot_errors: dict[tuple[str, str], BaseException], + env: GraphBuildEnv, ) -> None: """Build one SHA's sanitized clone template and graph. Idempotent per key.""" - if graph_key in graph_snapshots or graph_key in graph_snapshot_errors: + if graph_key in env.graph_snapshots or graph_key in env.graph_snapshot_errors: return try: validate_no_prebuilt_graph_assets(task) - if graph_key not in clone_templates and graph_key not in clone_template_errors: - template = make_worktree(repo, task_sha, trees) + if graph_key not in env.clone_templates and graph_key not in env.clone_template_errors: + template = make_worktree(repo, task_sha, env.trees) template_head = sanitize_clone_for_hidden_oracles(template) - clone_templates[graph_key] = (template, template_head) + env.clone_templates[graph_key] = (template, template_head) clone_template: Path | None = None template_head: str | None = None - if graph_key in clone_templates: - clone_template, template_head = clone_templates[graph_key] - if graph_key in clone_template_errors: - graph_snapshot_errors[graph_key] = clone_template_errors[graph_key] + if graph_key in env.clone_templates: + clone_template, template_head = env.clone_templates[graph_key] + if graph_key in env.clone_template_errors: + env.graph_snapshot_errors[graph_key] = env.clone_template_errors[graph_key] return - graph_snapshots[graph_key] = prepare_sanitized_graph( + env.graph_snapshots[graph_key] = prepare_sanitized_graph( task, repo=repo, resolved_sha=task_sha, - parent=trees, - cache=task_asset_cache, - claude_bin=claude_bin, - bwrap_bin=bwrap_bin, - sandbox_backend=sandbox_backend, - runtime_mounts=runtime_mounts, + parent=env.trees, + cache=env.task_asset_cache, + claude_bin=env.claude_bin, + bwrap_bin=env.bwrap_bin, + sandbox_backend=env.sandbox_backend, + runtime_mounts=env.runtime_mounts, clone_template=clone_template, sanitized_head=template_head, ) except (ManagedProcessError, OSError, SandboxError, RuntimeError, ValueError) as exc: - graph_snapshot_errors[graph_key] = exc - clone_template_errors.setdefault(graph_key, exc) + env.graph_snapshot_errors[graph_key] = exc + env.clone_template_errors.setdefault(graph_key, exc) @dataclass @@ -1930,16 +1947,7 @@ def prefetch_next_graph( task: Mapping[str, Any], binding: Mapping[str, Any], graph_key: tuple[str, str], - trees: Path, - task_asset_cache: TaskAssetCache, - claude_bin: Path | str, - bwrap_bin: Path | str, - sandbox_backend: str, - runtime_mounts: Sequence[ReadOnlyMount], - clone_templates: dict[tuple[str, str], tuple[Path, str]], - clone_template_errors: dict[tuple[str, str], BaseException], - graph_snapshots: dict[tuple[str, str], SanitizedGraphSnapshot], - graph_snapshot_errors: dict[tuple[str, str], BaseException], + env: GraphBuildEnv, cancel_event: threading.Event, ) -> GraphPrefetch: """Start clone+graph prep for the next unpaid SHA during paid sessions.""" @@ -1951,22 +1959,7 @@ def prefetch_next_graph( if cancel_event.is_set(): return print(f"[prefetch_next_graph] clone+graph for {task_sha}") - ensure_task_graph( - task=task, - repo=repo, - task_sha=task_sha, - graph_key=graph_key, - trees=trees, - task_asset_cache=task_asset_cache, - claude_bin=claude_bin, - bwrap_bin=bwrap_bin, - sandbox_backend=sandbox_backend, - runtime_mounts=runtime_mounts, - clone_templates=clone_templates, - clone_template_errors=clone_template_errors, - graph_snapshots=graph_snapshots, - graph_snapshot_errors=graph_snapshot_errors, - ) + ensure_task_graph(task=task, repo=repo, task_sha=task_sha, graph_key=graph_key, env=env) # copy_context, as the worker pool already does at _run_wave: a plain Thread # does not inherit ContextVars, so without this every run_managed inside the @@ -2091,10 +2084,18 @@ def _run_sweep( f"reuse-results {reuse_source}: {len(reusable_rows)} comparator " f"cell(s) match this sweep; candidate arms always run" ) - graph_snapshots: dict[tuple[str, str], SanitizedGraphSnapshot] = {} - graph_snapshot_errors: dict[tuple[str, str], BaseException] = {} - clone_templates: dict[tuple[str, str], tuple[Path, str]] = {} - clone_template_errors: dict[tuple[str, str], BaseException] = {} + graph_env = GraphBuildEnv( + trees=Path(trees), + task_asset_cache=task_asset_cache, + claude_bin=args.claude_bin, + bwrap_bin=bwrap_bin, + sandbox_backend=sandbox_backend, + runtime_mounts=runtime_mounts, + clone_templates={}, + clone_template_errors={}, + graph_snapshots={}, + graph_snapshot_errors={}, + ) graph_prefetch: GraphPrefetch | None = None sweep_rows = list(zip(tasks, task_bindings, oracle_snapshots, strict=True)) @@ -2139,26 +2140,13 @@ def _run_sweep( graph_key = (str(repo), task_sha) if graph_prefetch is not None and graph_prefetch.key == graph_key: _join_graph_prefetch() - if paid_cells and graph_key not in graph_snapshots and graph_key not in graph_snapshot_errors: + if paid_cells: ensure_task_graph( - task=task, - repo=repo, - task_sha=task_sha, - graph_key=graph_key, - trees=Path(trees), - task_asset_cache=task_asset_cache, - claude_bin=args.claude_bin, - bwrap_bin=bwrap_bin, - sandbox_backend=sandbox_backend, - runtime_mounts=runtime_mounts, - clone_templates=clone_templates, - clone_template_errors=clone_template_errors, - graph_snapshots=graph_snapshots, - graph_snapshot_errors=graph_snapshot_errors, + task=task, repo=repo, task_sha=task_sha, graph_key=graph_key, env=graph_env ) - graph_snapshot = graph_snapshots.get(graph_key) - graph_snapshot_error = graph_snapshot_errors.get(graph_key) - clone_template, template_head = clone_templates.get(graph_key, (None, None)) + graph_snapshot = graph_env.graph_snapshots.get(graph_key) + graph_snapshot_error = graph_env.graph_snapshot_errors.get(graph_key) + clone_template, template_head = graph_env.clone_templates.get(graph_key, (None, None)) cell_context = TaskCellContext( task=task, oracle_snapshot=oracle_snapshot, @@ -2209,56 +2197,40 @@ def _run_sweep( f"({started_cells}/{total_cells}, {(time.monotonic() - sweep_started) / 60:.0f}m elapsed)" ) keep(run_idx, arm, record) - if paid_cells: - if graph_prefetch is None and not cancel_event.is_set(): - ready_keys = ( - set(clone_templates) - | set(clone_template_errors) - | set(graph_snapshots) - | set(graph_snapshot_errors) + if paid_cells and graph_prefetch is None and not cancel_event.is_set(): + target = next_graph_prefetch_target( + [(later_task, later_binding) for later_task, later_binding, _ in sweep_rows[index + 1 :]], + arms=args.arms, + runs=args.runs, + reusable_rows=reusable_rows, + reuse_source=reuse_source, + ready_keys=graph_env.ready_keys(), + ) + if target is not None: + later_task, later_binding, later_key = target + graph_prefetch = prefetch_next_graph( + task=later_task, + binding=later_binding, + graph_key=later_key, + env=graph_env, + cancel_event=cancel_event, ) - target = next_graph_prefetch_target( - [(later_task, later_binding) for later_task, later_binding, _ in sweep_rows[index + 1 :]], - arms=args.arms, - runs=args.runs, - reusable_rows=reusable_rows, - reuse_source=reuse_source, - ready_keys=ready_keys, - ) - if target is not None: - later_task, later_binding, later_key = target - graph_prefetch = prefetch_next_graph( - task=later_task, - binding=later_binding, - graph_key=later_key, - trees=Path(trees), - task_asset_cache=task_asset_cache, - claude_bin=args.claude_bin, - bwrap_bin=bwrap_bin, - sandbox_backend=sandbox_backend, - runtime_mounts=runtime_mounts, - clone_templates=clone_templates, - clone_template_errors=clone_template_errors, - graph_snapshots=graph_snapshots, - graph_snapshot_errors=graph_snapshot_errors, - cancel_event=cancel_event, - ) - # A reused success is evidence the pipeline can produce a good - # row, so it resets the consecutive-failure count the same way a - # paid success does. Leaving reused rows out let a streak carry - # across them and trip on stale history. + # A reused success is evidence the pipeline can produce a good row, + # so it resets the consecutive-failure count the same way a paid + # success does. Leaving reused rows out let a streak carry across + # them and trip on stale history. for _run_idx, _arm, _record in reused_records: outage_streak = systemic_outage_streak(_record.get("error_kind"), outage_streak) outage_streak, outage_tripped = sweep_task_cells( - paid_cells, - workers=args.workers, - run=partial(run_cell, cell_context), - on_start=announce, - on_record=keep, - outage_streak=outage_streak, - outage_limit=args.outage_streak, - cancel_event=cancel_event, - ) + paid_cells, + workers=args.workers, + run=partial(run_cell, cell_context), + on_start=announce, + on_record=keep, + outage_streak=outage_streak, + outage_limit=args.outage_streak, + cancel_event=cancel_event, + ) results[task["id"]] = {a: aggregate(rs) for a, rs in per_arm.items() if rs} finally: _join_graph_prefetch()