Files
kernbench2/src/kernbench/runtime_api/bench_runner.py
T
ywkang 998cc85762 Add PE-level IPCQ collective infra + unified ccl_allreduce bench (ADR-0023)
Major changes:

PE-level IPCQ infrastructure:
- New PE_IPCQ component: ring-buffer control plane with 4-direction
  neighbor mapping, head/tail pointers, backpressure (poll/sleep).
- PE_DMA extended with vc_comm channel for IPCQ outbound/inbound DMA,
  including in-flight data snapshot (D9) and op_log recording at
  outbound time for Phase 2 replay correctness.
- IpcqDmaToken piggyback model: data + metadata travel together,
  atomic visibility at receiver (invariant I6).
- Credit return fast path: bottleneck-BW latency, no fabric vc_comm.

Phase 2 data execution (ADR-0020 integration):
- op_log extended: DmaWriteCmd now captures src_space/src_addr for
  Phase 2 dma_write copy; ipcq_copy ops recorded at outbound time.
- DataExecutor replays dma_write + ipcq_copy in t_start order.
- Engine._flush_data_phase: incremental cursor-based replay after
  each engine.wait() so host reads see post-Phase-2 data.
- KernelRunner Phase 1 writes disabled when op_log is active to
  prevent stale data from corrupting the MemoryStore snapshot.

TLContext / kernel API:
- tl.send(dir, src=TensorHandle), tl.recv(dir, shape, dtype),
  tl.recv_async, tl.wait(RecvFuture), copy_to_dst mode.
- TensorHandle operator overloading (add/sub/mul/div) via thread-local
  active TLContext → MathCmd dispatch through PE_MATH.
- PE-local scratch allocator for math output handles.
- tl.load returns space="hbm" handles for correct Phase 2 addressing.
- Additional math functions: maximum, minimum, fma, clamp, softmax, cdiv.

Unified ccl_allreduce bench (PyTorch-compat host code):
- Single benches/ccl_allreduce.py with run() + worker(rank, ws, torch)
  split matching real PyTorch DDP worker pattern.
- torch.distributed facade: init_process_group, get_world_size,
  get_rank, get_backend, all_reduce, barrier — only real PyTorch names.
- AhbmCCLBackend: eager install_ipcq at init, all_reduce dispatches
  kernel via tensor shard metadata (n_elem from shards[0].nbytes).
- world_size derived from topology spec (sips × cubes × pes_per_cube)
  with optional algorithm-level override in ccl.yaml.

Tensor API (PyTorch-compat surface):
- Tensor.numpy(): gather-aware (all shards via VA-based addressing).
- Tensor.copy_(source): scatter from host tensor into sharded target.
- RuntimeContext.from_numpy(arr): host-side staging tensor.
- Tensor.data property fixed to use numpy() (was shards[0]-only).

Algorithm modules moved to src/kernbench/ccl/algorithms/:
- ring_allreduce, mesh_allreduce, tree_allreduce, hello_send.
- Each module exports kernel_args(world_size, n_elem) helper.
- ccl.yaml module paths updated to kernbench.ccl.algorithms.*.

Dead code removed:
- 7 per-variant bench files (ccl_allreduce_{tcm,hbm,sram}, etc.).
- _run_ccl_bench greenlet-per-SIP scheduler.
- benches.loader.is_ccl_bench + run_rank detection.
- benches/ccl/ directory.

Tests:
- New test_ccl_allreduce_matrix.py: 7 parametrized cases
  (ring×3 buffers, ring 8/16, mesh 4, tree 7).
- New test_runtime_api_tensor.py: copy_/numpy/from_numpy unit tests.
- Existing tests updated for new import paths + world_size_override.

Docs:
- Korean ccl-author-guide.md and ADR-0023 paths updated.
- New English versions: ccl-author-guide.en.md, ADR-0023.en.md.

502 tests pass.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-04-12 19:36:59 -07:00

96 lines
3.2 KiB
Python

from __future__ import annotations
from collections.abc import Callable
from enum import Enum
from typing import Any
from kernbench.common.types import Completion, SimEngine, Trace
from .context import RuntimeContext
from .types import BenchResult, DeviceSelector
class CompletionPolicy(str, Enum):
LAST_SUBMITTED = "last_submitted"
LAST_COMPLETED = "last_completed" # requires trace/timestamps or engine support; stub for now
ALL_OK_FAIL_FAST = "all_ok_fail_fast"
BenchFn = Callable[[RuntimeContext], Any]
EngineFactory = Callable[[object, DeviceSelector], SimEngine]
def run_bench(
*,
topology: object,
bench_fn: BenchFn,
device: DeviceSelector,
engine_factory: EngineFactory,
correlation_id: str = "bench0",
completion_policy: CompletionPolicy = CompletionPolicy.LAST_SUBMITTED,
) -> BenchResult:
"""Minimal bench runner.
- topology: compiled topology object (opaque to runtime here)
- bench_fn: callable ``run(torch)`` receiving a RuntimeContext
- device: DeviceSelector ("all" or "sip:<N>")
- engine_factory: builds sim_engine for given topology & device
- completion_policy: how to determine overall completion/result
"""
engine = engine_factory(topology, device)
# Extract spec from TopologyHandle or TopologyGraph
topo_obj = getattr(topology, "topology_obj", topology)
spec = getattr(topo_obj, "spec", None)
ctx = RuntimeContext(
engine=engine, target_device=device,
correlation_id=correlation_id, spec=spec,
)
bench_fn(ctx)
ctx.wait_all()
collected_traces = ctx._traces or None
handles = ctx.handles()
if not handles:
return BenchResult(
completion=Completion(
ok=False, error_code="NO_REQUESTS", error_message="Bench submitted no requests"
),
correlation_id=correlation_id,
trace=None,
traces=collected_traces,
engine=engine,
)
if completion_policy == CompletionPolicy.LAST_SUBMITTED:
last = handles[-1]
completion, trace = engine.get_completion(last)
return BenchResult(
completion=completion, correlation_id=correlation_id,
trace=trace, traces=collected_traces, engine=engine,
)
if completion_policy == CompletionPolicy.ALL_OK_FAIL_FAST:
last_trace: Trace | None = None
for h in handles:
c, t = engine.get_completion(h)
last_trace = t if t is not None else last_trace
if not c.ok:
return BenchResult(
completion=c, correlation_id=correlation_id,
trace=last_trace, traces=collected_traces, engine=engine,
)
return BenchResult(
completion=Completion(ok=True), correlation_id=correlation_id,
trace=last_trace, traces=collected_traces, engine=engine,
)
# LAST_COMPLETED placeholder (needs engine support for timing). Fall back.
last = handles[-1]
completion, trace = engine.get_completion(last)
return BenchResult(
completion=completion, correlation_id=correlation_id,
trace=trace, traces=collected_traces, engine=engine,
)