Skip to content
Closed
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
2 changes: 1 addition & 1 deletion docs/capability-survey.md
Original file line number Diff line number Diff line change
Expand Up @@ -217,7 +217,7 @@ values yield `PTO2_ERROR_ASYNC_COMPLETION_INVALID`.

| Engine | a2a3 | a5 | Status |
| ------ | ---- | -- | ------ |
| COUNTER (default) | registered | registered | **Shipped** — `async_notify_demo` runs onboard on both arches and `deferred_notify_demo` runs in sim on both, through the `st-onboard-*` / `st-sim-*` jobs in `ci.yml`. Routed by `@pytest.mark.platforms`, no `skipif` |
| COUNTER (default) | registered | registered | **Shipped** — `async_notify_demo` runs onboard on both arches and `deferred_notify_demo` runs in sim on both, through the `st-onboard-*` / `st-sim-*` jobs in `ci.yml`. Routed by `CASES[*]["platforms"]`, no `skipif` |
| SDMA | build macro forced ON; runtime opt-in | `option(... OFF)` | a2a3 **Shipped** (the "SDMA pytest (a2a3)" step in `ci.yml`); a5 not built |
| URMA | absent | full implementation | **Gated** — see below |
| ROCE, CCU | enum only | enum only | **Name only** |
Expand Down
4 changes: 4 additions & 0 deletions examples/a2a3/tensormap_and_ringbuffer/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,10 @@ completion rather than on task end. All but `deferred_notify_demo`, which runs
the simulator path, are **onboard-only**; every one except `prefetch_async_demo`
needs two dies.

These mechanism-focused examples use the scene-test lifecycle. For a complete
direct `Worker` communication-domain walkthrough from construction through
`close()`, see [`examples/workers/l3/allreduce/`](../../workers/l3/allreduce/).

| Example | Mechanism | Devices |
| ------- | --------- | ------- |
| [`prefetch_async_demo/`](prefetch_async_demo/) | `TPREFETCH_ASYNC` over the runtime-injected SDMA workspace, provisioned once by `Worker(enable_sdma=True)` and injected into every kernel's `GlobalContext`. | 1 |
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,173 +7,106 @@
# 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.
# -----------------------------------------------------------------------------------------------------------
"""Notification counter + deferred completion smoke test for onboard a2a3."""
"""Notification counter and deferred-completion smoke test for onboard a2a3."""

from __future__ import annotations

import argparse
import os

import pytest
import torch
from simpler.task_interface import (
ArgDirection,
CallConfig,
ChipCallable,
CommBufferSpec,
CoreCallable,
DataType,
TaskArgs,
Tensor,
TensorArgType,
)
from simpler.worker import Worker
from simpler.task_interface import ArgDirection as D
from simpler.task_interface import CommBufferSpec, DataType, TaskArgs, TensorArgType
from simpler.task_interface import Tensor as DeviceTensor

from simpler_setup.elf_parser import extract_text_section
from simpler_setup.kernel_compiler import KernelCompiler
from simpler_setup.pto_isa import ensure_pto_isa_root
from simpler_setup import SceneTestCase, TaskArgsBuilder, Tensor, scene_test
from simpler_setup.torch_interop import make_tensor_arg

HERE = os.path.dirname(os.path.abspath(__file__))
N = 128 * 128


def parse_device_range(spec: str) -> list[int]:
if "," in spec:
return [int(x) for x in spec.split(",") if x]
if "-" in spec:
lo, hi = (int(x) for x in spec.split("-"))
return list(range(lo, hi + 1))
return [int(spec)]


def build_chip_callable(platform: str) -> ChipCallable:
kc = KernelCompiler(platform=platform)
runtime = "tensormap_and_ringbuffer"
pto_isa_root = ensure_pto_isa_root()
include_dirs = kc.get_orchestration_include_dirs(runtime)
extra_includes = list(include_dirs) + [str(kc.project_root / "src" / "common")]

children = []
for func_id, rel in [
(0, "kernels/aiv/kernel_producer_notify.cpp"),
(1, "kernels/aiv/kernel_consumer.cpp"),
(2, "kernels/aiv/kernel_notify_wait.cpp"),
]:
kernel = kc.compile_incore(
source_path=os.path.join(HERE, rel),
core_type="aiv",
pto_isa_root=pto_isa_root,
extra_include_dirs=extra_includes,
)
if not platform.endswith("sim"):
kernel = extract_text_section(kernel)
children.append(
(
func_id,
CoreCallable.build(
signature=[ArgDirection.IN, ArgDirection.OUT, ArgDirection.OUT, ArgDirection.IN],
binary=kernel,
NRANKS = 2


def async_notify_orch_fn(orch, callables, task_args, config):
with orch.allocate_domain(
name="default",
workers=list(range(NRANKS)),
window_size=4 * 1024,
buffers=[CommBufferSpec(name="notify_counter", dtype="int32", count=1, nbytes=4)],
) as handle:
for rank in range(NRANKS):
domain = handle[rank]
args = TaskArgs()
args.add_tensor(make_tensor_arg(getattr(task_args, f"in_{rank}")), TensorArgType.INPUT)
args.add_tensor(make_tensor_arg(getattr(task_args, f"out_{rank}")), TensorArgType.OUTPUT_EXISTING)
args.add_tensor(make_tensor_arg(getattr(task_args, f"result_{rank}")), TensorArgType.OUTPUT_EXISTING)
args.add_tensor(
DeviceTensor.make(
data=domain.buffer_ptrs["notify_counter"],
shapes=(1,),
dtype=DataType.INT32,
child_memory=True,
),
TensorArgType.INPUT,
)
)

orch = kc.compile_orchestration(
runtime_name=runtime,
source_path=os.path.join(HERE, "kernels/orchestration/async_notify_orchestration.cpp"),
extra_include_dirs=[str(kc.project_root / "src" / "common")],
)
return ChipCallable.build(
signature=[ArgDirection.IN, ArgDirection.OUT, ArgDirection.OUT, ArgDirection.IN],
func_name="async_notify_orchestration",
binary=orch,
children=children,
)


def run(
platform: str = "a2a3",
device_ids: list[int] | None = None,
) -> int:
if device_ids is None:
device_ids = [0, 1]
nranks = len(device_ids)
if nranks != 2:
raise ValueError(f"async_notify_demo needs exactly 2 devices, got {device_ids}")

inp = [
torch.tensor([float(i % 251) / 10.0 for i in range(N)], dtype=torch.float32).share_memory_()
for _ in range(nranks)
]
out = [torch.zeros(N, dtype=torch.float32).share_memory_() for _ in range(nranks)]
result = [torch.zeros(N, dtype=torch.float32).share_memory_() for _ in range(nranks)]

chip_callable = build_chip_callable(platform)
worker = Worker(
level=3,
platform=platform,
runtime="tensormap_and_ringbuffer",
device_ids=device_ids,
num_sub_workers=0,
)
chip_handle = worker.register(chip_callable)
try:
worker.init()

def orch_fn(orch, _args, cfg):
with orch.allocate_domain(
name="default",
workers=list(range(nranks)),
window_size=4 * 1024,
buffers=[CommBufferSpec(name="notify_counter", dtype="int32", count=1, nbytes=4)],
) as handle:
for rank in range(nranks):
domain = handle[rank]
args = TaskArgs()
args.add_tensor(make_tensor_arg(inp[rank]), TensorArgType.INPUT)
args.add_tensor(make_tensor_arg(out[rank]), TensorArgType.OUTPUT_EXISTING)
args.add_tensor(make_tensor_arg(result[rank]), TensorArgType.OUTPUT_EXISTING)
args.add_tensor(
Tensor.make(
data=domain.buffer_ptrs["notify_counter"],
shapes=(1,),
dtype=DataType.INT32,
child_memory=True,
),
TensorArgType.INPUT,
args.add_scalar(domain.device_ctx)
orch.submit_next_level(callables.async_notify, args, config, worker=rank)


@scene_test(level=3, runtime="tensormap_and_ringbuffer")
class TestAsyncNotifyDemo(SceneTestCase):
CALLABLE = {
"orchestration": async_notify_orch_fn,
"callables": [
{
"name": "async_notify",
"orchestration": {
"source": "kernels/orchestration/async_notify_orchestration.cpp",
"function_name": "async_notify_orchestration",
"signature": [D.IN, D.OUT, D.OUT, D.IN],
},
"incores": [
{
"func_id": func_id,
"source": source,
"core_type": "aiv",
"signature": [D.IN, D.OUT, D.OUT, D.IN],
}
for func_id, source in enumerate(
[
"kernels/aiv/kernel_producer_notify.cpp",
"kernels/aiv/kernel_consumer.cpp",
"kernels/aiv/kernel_notify_wait.cpp",
]
)
args.add_scalar(domain.device_ctx)
orch.submit_next_level(chip_handle, args, cfg, worker=rank)

worker.run(orch_fn, args=None, config=CallConfig())

ok = True
for rank in range(nranks):
expected_out = inp[rank] * 2.0
expected_result = expected_out + 1.0
max_out = float(torch.max(torch.abs(out[rank] - expected_out)))
max_result = float(torch.max(torch.abs(result[rank] - expected_result)))
print(f"[async_notify_demo] rank {rank}: max_out={max_out:.3e} max_result={max_result:.3e}")
ok = ok and max_out <= 1e-3 and max_result <= 1e-3
return 0 if ok else 1
finally:
worker.close()


@pytest.mark.platforms(["a2a3"])
@pytest.mark.runtime("tensormap_and_ringbuffer")
@pytest.mark.device_count(2)
def test_async_notify_demo(st_device_ids, st_platform) -> None:
assert run(st_platform, [int(d) for d in st_device_ids]) == 0

],
}
],
}
CASES = [
{
"name": "notification_counter",
"platforms": ["a2a3"],
"config": {"device_count": NRANKS, "num_sub_workers": 0},
"params": {},
}
]
RTOL = 0.0
ATOL = 1e-3

def generate_args(self, params):
specs = []
for rank in range(NRANKS):
inp = torch.tensor([float(i % 251) / 10.0 for i in range(N)], dtype=torch.float32)
specs.extend(
[
Tensor(f"in_{rank}", inp),
Tensor(f"out_{rank}", torch.zeros(N, dtype=torch.float32)),
Tensor(f"result_{rank}", torch.zeros(N, dtype=torch.float32)),
]
)
return TaskArgsBuilder(*specs)

def main() -> int:
parser = argparse.ArgumentParser()
parser.add_argument("-p", "--platform", default="a2a3")
parser.add_argument("-d", "--device", default="0-1")
args = parser.parse_args()
return run(args.platform, parse_device_range(args.device))
def compute_golden(self, args, params):
for rank in range(NRANKS):
expected_out = getattr(args, f"in_{rank}") * 2.0
getattr(args, f"out_{rank}").copy_(expected_out)
getattr(args, f"result_{rank}").copy_(expected_out + 1.0)
Comment on lines +91 to +108

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Inspect the async notify kernels to confirm local vs peer data flow.
fd -t f -p 'a2a3/tensormap_and_ringbuffer/async_notify_demo/kernels' -x sed -n '1,200p' {}

Repository: hw-native-sys/simpler

Length of output: 10053


Make the per-rank input rank-dependent.

Each rank computes out = in * 2 only from in_{rank}, so a rank-swap can still pass if every rank uses the same input values. Use a rank-offset input, such as rank * 1000 + (i % 251), so the golden fails when data from the wrong rank is used locally.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In
`@examples/a2a3/tensormap_and_ringbuffer/async_notify_demo/test_async_notify_demo.py`
around lines 91 - 108, Update generate_args to make each rank’s input tensor
rank-dependent by incorporating rank into the per-element value, such as a
rank-based offset plus i % 251, while preserving the existing dtype and tensor
structure. Leave compute_golden unchanged so it continues deriving expected
outputs from each rank’s own input.



if __name__ == "__main__":
raise SystemExit(main())
SceneTestCase.run_module(__name__)
Loading
Loading