Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,8 @@
# -----------------------------------------------------------------------------------------------------------
"""Paged attention with small ring buffer sizes — stress test for ring rotation/reclamation.

Tests RUNTIME_ENV (PTO2_RING_TASK_WINDOW, PTO2_RING_HEAP, PTO2_RING_DEP_POOL),
Drives per-case ring sizing through ``config.runtime_env`` (ring_task_window /
ring_heap / ring_dep_pool) rather than the process-global PTO2_RING_* env, plus
INOUT tensors, bfloat16, and AIC+AIV mixed execution.
"""

Expand All @@ -29,11 +30,6 @@ class TestPagedAttentionRingbuffer(SceneTestCase):

RTOL = 1e-3
ATOL = 1e-3
RUNTIME_ENV = {
"PTO2_RING_TASK_WINDOW": "64",
"PTO2_RING_HEAP": "2621440",
"PTO2_RING_DEP_POOL": "256",
}

CALLABLE = {
"orchestration": {
Expand Down Expand Up @@ -73,7 +69,19 @@ class TestPagedAttentionRingbuffer(SceneTestCase):
{
"name": "ringbuffer_stress",
"platforms": ["a2a3"],
"config": {"aicpu_thread_num": 4, "block_dim": 24},
# ring_heap must be a power of 2; the historical RUNTIME_ENV value
# 2621440 (2.5 MiB) was silently rejected by the env parser, so the
# case actually ran with the 256 MiB default. 4 MiB keeps the
# small-ring stress intent with a valid size.
"config": {
"aicpu_thread_num": 4,
"block_dim": 24,
"runtime_env": {
"ring_task_window": 64,
"ring_heap": 4 * 1024 * 1024,
"ring_dep_pool": 256,
},
},
Comment thread
ChaoZheng109 marked this conversation as resolved.
"params": {
"batch": 32,
"num_heads": 16,
Expand Down
53 changes: 53 additions & 0 deletions examples/workers/l2/per_task_runtime_env/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
# `per_task_runtime_env/` — per-task ring sizing on one L2 Worker

Runs the same vector_add kernel three times on one L2 `Worker`, each with a
different `CallConfig.runtime_env` (ring buffer sizing). Ring sizing is a
**per-run** knob carried on `CallConfig` — not a process-wide env export.

## What it shows

`CallConfig.runtime_env` groups the three ring overrides as a distinct config
tier, separate from the top-level execution knobs (`block_dim`, …):

| field | unit | constraint |
| ----- | ---- | ---------- |
| `ring_task_window` | tasks | power of 2, >= 4 |
| `ring_heap` | bytes / ring | power of 2, >= 1024 |
| `ring_dep_pool` | entries | 4 .. INT32_MAX |

Precedence per value: **`runtime_env` field > `PTO2_RING_*` env var >
compile-time default**. A field left at 0 (or omitted) falls back to the env
var / default.

```python
cfg = CallConfig()
cfg.runtime_env.ring_task_window = 128
cfg.runtime_env.ring_heap = 8 * 1024 * 1024 # bytes per ring
cfg.runtime_env.ring_dep_pool = 256
worker.run(chip_handle, args, cfg)
```

The three runs (`small_ring`, `large_ring`, `env_or_default`) compute the same
vector add and all pass golden — only the ring footprint differs.

## Layout

```text
per_task_runtime_env/
main.py # 3 runs, one CallConfig.runtime_env each
test_per_task_runtime_env.py
```

The kernel is reused verbatim from the sibling `../vector_add/kernels` — this
example only varies the per-run ring configuration.

## Run

```bash
python examples/workers/l2/per_task_runtime_env/main.py -p a2a3sim -d 0
```

See [`../vector_add/main.py`](../vector_add/main.py) for the full L2 lifecycle
walk-through (kernel compile, `ChipCallable` assembly, device memory, readback).
For dispatching several L2 tasks with distinct ring sizes from one launch, see
[`../../l3/per_task_runtime_env/`](../../l3/per_task_runtime_env/).
9 changes: 9 additions & 0 deletions examples/workers/l2/per_task_runtime_env/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
# Copyright (c) PyPTO Contributors.
# This program is free software, you can redistribute it and/or modify it under the terms and conditions of
# CANN Open Software License Agreement Version 2.0 (the "License").
# Please refer to the License for details. You may not use this file except in compliance with the License.
# THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
# INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
# See LICENSE in the root of the software repository for the full text of the License.
# -----------------------------------------------------------------------------------------------------------
"""Package marker so ``test_*.py`` can do ``from .main import run``."""
189 changes: 189 additions & 0 deletions examples/workers/l2/per_task_runtime_env/main.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,189 @@
#!/usr/bin/env python3
# Copyright (c) PyPTO Contributors.
# This program is free software, you can redistribute it and/or modify it under the terms and conditions of
# CANN Open Software License Agreement Version 2.0 (the "License").
# Please refer to the License for details. You may not use this file except in compliance with the License.
# THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
# INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
# See LICENSE in the root of the software repository for the full text of the License.
# -----------------------------------------------------------------------------------------------------------
"""L2 Worker API demo — per-task ring sizing via ``CallConfig.runtime_env``.

Runs the same vector_add kernel three times on one L2 Worker, each time with a
different ``CallConfig.runtime_env`` (ring buffer sizing). Ring sizing is a
per-run knob carried on ``CallConfig`` — no process-wide ``PTO2_RING_*`` env
export needed, and each ``worker.run`` binds its ring buffers from the config
it was handed.

runtime_env fields (0 / unset => fall back to env var / compile default):
ring_task_window power of 2, >= 4
ring_heap bytes per ring, power of 2, >= 1024
ring_dep_pool 4 .. INT32_MAX
Precedence: runtime_env field > PTO2_RING_* env var > compile-time default.

See ../vector_add/main.py for the full L2 lifecycle walk-through; this example
reuses that kernel verbatim and only varies the per-run ring configuration.

Run:
python examples/workers/l2/per_task_runtime_env/main.py -p a2a3sim -d 0
"""

import argparse
import os
import sys
from typing import Optional

os.environ.setdefault("KMP_DUPLICATE_LIB_OK", "TRUE")

import torch # noqa: E402
from simpler.task_interface import (
ArgDirection,
CallConfig,
ChipCallable,
ChipStorageTaskArgs,
ContinuousTensor,
CoreCallable,
DataType,
)
from simpler.worker import Worker

from simpler_setup.kernel_compiler import KernelCompiler
from simpler_setup.pto_isa import ensure_pto_isa_root

HERE = os.path.dirname(os.path.abspath(__file__))
# Reuse the sibling vector_add kernel verbatim — this example only varies ring sizing.
VECTOR_ADD_KERNELS = os.path.join(HERE, "..", "vector_add", "kernels")

N_ROWS = 128
N_COLS = 128
N_ELEMS = N_ROWS * N_COLS
NBYTES = N_ELEMS * 4 # float32

# (label, runtime_env dict or None). None => no override; falls back to the
# PTO2_RING_* env var / compile-time default. Same kernel + same inputs run
# under every sizing, so all three produce identical (correct) output.
RING_CONFIGS = [
("small_ring", {"ring_task_window": 16, "ring_heap": 1 * 1024 * 1024, "ring_dep_pool": 64}),
("large_ring", {"ring_task_window": 128, "ring_heap": 8 * 1024 * 1024, "ring_dep_pool": 256}),
("env_or_default", None),
]


def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
parser.add_argument("-p", "--platform", required=True, choices=["a2a3sim", "a2a3"])
parser.add_argument("-d", "--device", type=int, default=0)
return parser.parse_args()


def build_chip_callable(platform: str) -> ChipCallable:
"""Compile the reused vector_add sources into a ChipCallable.

Identical to ../vector_add/main.py::build_chip_callable except the kernel
sources are read from the sibling vector_add example.
"""
kc = KernelCompiler(platform=platform)
runtime = "tensormap_and_ringbuffer"
pto_isa_root = ensure_pto_isa_root(clone_protocol="https")
include_dirs = kc.get_orchestration_include_dirs(runtime)

kernel_bytes = kc.compile_incore(
source_path=os.path.join(VECTOR_ADD_KERNELS, "aiv", "vector_add_kernel.cpp"),
core_type="aiv",
pto_isa_root=pto_isa_root,
extra_include_dirs=include_dirs,
)
if not platform.endswith("sim"):
from simpler_setup.elf_parser import extract_text_section # noqa: PLC0415

kernel_bytes = extract_text_section(kernel_bytes)

orch_bytes = kc.compile_orchestration(
runtime_name=runtime,
source_path=os.path.join(VECTOR_ADD_KERNELS, "orchestration", "vector_add_orch.cpp"),
)
core_callable = CoreCallable.build(
signature=[ArgDirection.IN, ArgDirection.IN, ArgDirection.OUT],
binary=kernel_bytes,
)
return ChipCallable.build(
signature=[ArgDirection.IN, ArgDirection.IN, ArgDirection.OUT],
func_name="vector_add_orchestration",
binary=orch_bytes,
children=[(0, core_callable)],
)


def _make_config(ring: Optional[dict]) -> CallConfig:
"""Build a CallConfig, attaching this run's ring sizing under runtime_env."""
cfg = CallConfig()
if ring is not None:
cfg.runtime_env.ring_task_window = ring["ring_task_window"]
cfg.runtime_env.ring_heap = ring["ring_heap"]
cfg.runtime_env.ring_dep_pool = ring["ring_dep_pool"]
return cfg


def _run_one(worker: Worker, chip_handle, label: str, ring: Optional[dict]) -> None:
"""One malloc → copy → run(config) → readback → verify cycle."""
torch.manual_seed(42)
host_a = torch.randn(N_ROWS, N_COLS, dtype=torch.float32)
host_b = torch.randn(N_ROWS, N_COLS, dtype=torch.float32)
expected = host_a + host_b
host_out = torch.zeros(N_ROWS, N_COLS, dtype=torch.float32)

dev_a = worker.malloc(NBYTES)
dev_b = worker.malloc(NBYTES)
dev_out = worker.malloc(NBYTES)
worker.copy_to(dev_a, host_a.data_ptr(), NBYTES)
worker.copy_to(dev_b, host_b.data_ptr(), NBYTES)

args = ChipStorageTaskArgs()
args.add_tensor(ContinuousTensor.make(dev_a, (N_ROWS, N_COLS), DataType.FLOAT32))
args.add_tensor(ContinuousTensor.make(dev_b, (N_ROWS, N_COLS), DataType.FLOAT32))
args.add_tensor(ContinuousTensor.make(dev_out, (N_ROWS, N_COLS), DataType.FLOAT32))

config = _make_config(ring)
print(f"[per_task_runtime_env] run '{label}': runtime_env={config.runtime_env!r}")
worker.run(chip_handle, args, config)

worker.copy_from(host_out.data_ptr(), dev_out, NBYTES)
worker.free(dev_a)
worker.free(dev_b)
worker.free(dev_out)

assert torch.allclose(host_out, expected, rtol=1e-5, atol=1e-5), f"{label} result mismatch"
print(f"[per_task_runtime_env] '{label}' golden check PASSED")


def run(platform: str, device_id: int) -> int:
"""Core logic — callable from both CLI and pytest."""
worker = Worker(
level=2,
platform=platform,
runtime="tensormap_and_ringbuffer",
device_id=device_id,
)

print(f"[per_task_runtime_env] compiling kernels for {platform}...")
chip_callable = build_chip_callable(platform)
chip_handle = worker.register(chip_callable)

print(f"[per_task_runtime_env] init worker (device={device_id})...")
worker.init()
try:
for label, ring in RING_CONFIGS:
_run_one(worker, chip_handle, label, ring)
finally:
worker.close()
print("[per_task_runtime_env] all ring configurations PASSED ✅")
return 0


def main() -> int:
args = parse_args()
return run(args.platform, args.device)


if __name__ == "__main__":
sys.exit(main())
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
# Copyright (c) PyPTO Contributors.
# This program is free software, you can redistribute it and/or modify it under the terms and conditions of
# CANN Open Software License Agreement Version 2.0 (the "License").
# Please refer to the License for details. You may not use this file except in compliance with the License.
# THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
# INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
# See LICENSE in the root of the software repository for the full text of the License.
# -----------------------------------------------------------------------------------------------------------
"""Hardware ST for examples/workers/l2/per_task_runtime_env."""

import pytest

from .main import run


@pytest.mark.platforms(["a2a3sim", "a2a3"])
@pytest.mark.runtime("tensormap_and_ringbuffer")
@pytest.mark.device_count(1)
def test_per_task_runtime_env(st_platform, st_device_ids):
rc = run(st_platform, int(st_device_ids[0]))
assert rc == 0
58 changes: 58 additions & 0 deletions examples/workers/l3/per_task_runtime_env/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
# `per_task_runtime_env/` — distinct ring sizes per L2 in one L3 launch

One L3 orchestration dispatches two L2 tasks, each binding its **own** ring
buffers via `CallConfig.runtime_env`. This is the headline use case the
per-task ring sizing enables: heterogeneous L2 tasks in a single launch that
each need a different ring footprint.

## What it shows

Before this knob, every L2 dispatched from one L3 shared the process-wide
`PTO2_RING_*` env and could not be sized independently. Now each
`submit_next_level` gets its own `CallConfig`:

```python
def orch_fn(orch, _args, _cfg):
for spec in L2_TASKS: # one entry per L2 task
cfg = CallConfig()
cfg.runtime_env.ring_task_window = spec["ring_task_window"]
cfg.runtime_env.ring_heap = spec["ring_heap"] # bytes per ring
cfg.runtime_env.ring_dep_pool = spec["ring_dep_pool"]
orch.submit_next_level(chip_handle, chip_args, cfg) # per-task config
```

The per-task config travels through the mailbox to the chip child, so each L2
binds its rings from its own values. The demo dispatches `l2_small`
(16 / 1 MiB / 64) and `l2_large` (128 / 8 MiB / 256); both run the same
vector_add and pass golden.

### Derive per-task config from the base, don't rebuild it

`runtime_env` lives on `CallConfig` alongside the diagnostics flags
(`enable_scope_stats`, `output_prefix`, …). The per-task config must preserve
any fields the harness injected on the orchestration's base config — otherwise
`--enable-scope-stats` and friends silently collect nothing for that L2. So
`_l2_config(base, spec)` copies those fields from `base` and overrides only
`runtime_env`, rather than starting from a blank `CallConfig()`.

## Layout

```text
per_task_runtime_env/
main.py # 2 submit_next_level, one runtime_env each
test_per_task_runtime_env.py
```

The kernel is reused verbatim from `../../l2/vector_add/kernels`.

## Run

```bash
python examples/workers/l3/per_task_runtime_env/main.py -p a2a3sim -d 0
```

The two L2 tasks run serially on one device. See
[`../multi_chip_dispatch/`](../multi_chip_dispatch/) for the multi-device DAG
primitives (`worker=i` pinning, `submit_sub`), and
[`../../l2/per_task_runtime_env/`](../../l2/per_task_runtime_env/) for the
single-L2 version.
9 changes: 9 additions & 0 deletions examples/workers/l3/per_task_runtime_env/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
# Copyright (c) PyPTO Contributors.
# This program is free software, you can redistribute it and/or modify it under the terms and conditions of
# CANN Open Software License Agreement Version 2.0 (the "License").
# Please refer to the License for details. You may not use this file except in compliance with the License.
# THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
# INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
# See LICENSE in the root of the software repository for the full text of the License.
# -----------------------------------------------------------------------------------------------------------
"""Package marker so ``test_*.py`` can do ``from .main import run``."""
Loading
Loading