Skip to content
Draft
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
50 changes: 50 additions & 0 deletions docs/nvtx-events.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
<!--
SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
SPDX-License-Identifier: Apache-2.0
-->

# Remote delivery NVTX events

Install the `profiling` extra. `KVCR_NVTX_LEVEL=off|low|medium` selects
library-side detail; the default is low when the optional binding is available.
Annotations do not depend on whether Nsight is attached. See [profiling](profiling.md)
for capture commands and the pin schema.

## Source transfers

`nixl.write.submit` is a same-thread synchronous scope. `nixl.write.posted`
reports the native posting result, including rejection and ambiguity.
`nixl.done_observed` records the first observed native DONE, including synchronous
DONE returned by posting. It is distinct from `nixl.write.released` (successful
handle release) and `source.write.completed` (logical source operation result).
A prior error or cancellation can make the logical result fail even after DONE.
The status on the release event reflects the transfer outcome, not a release
failure. Release failures emit `nixl.write.release_retry` at medium detail.

`nixl.write.error`, `source.write.cancel_requested`, and
`source.write.shutdown_unresolved` preserve the error, timeout/cancellation, and
unresolved ownership boundaries. Elapsed time does not prove DMA quiescence.
`source.write.refused` identifies a stale route before posting. No native numeric
error is inferred from exception text; `native_error_known=0` means unavailable.

Schema 2 events use `(instance_hi, instance_lo, trace_id)` for a local lifecycle.
Join schema 1 pin/waiter records using `(instance_hi, instance_lo, source_op_id)`.
Across workers use `(target_agent_hi, target_agent_lo,
target_incarnation_hi, target_incarnation_lo, op_handle)`; incarnation is taken
from existing control metadata, without extending the wire protocol. Zero
identity fields mean unavailable. Handles remain signed, including local fills.
Transfer IDs are local to the source instance. Unknown counts/bytes are `-1`.

`source_tier`/`destination_tier` describe memory ownership: 0 unknown, 1 framework,
2 KVCR-owned, 3 mixed. `source_memory`/`destination_memory` describe a verified
registration: 0 unknown, 1 DRAM, 2 VRAM, 3 FILE. Ownership alone does not identify
the physical memory tier; remote registration facts are left unknown when absent.

Request identities use a stable 128-bit BLAKE2b digest of the complete UTF-8
string. `request.context` carries a display prefix as an unsigned byte array,
explicit byte length (up to 256), and a truncation flag. Decode exactly `length`
bytes as UTF-8; a truncated prefix may end inside a character. Join by the digest,
never by the display prefix. Empty strings and unavailable IDs are distinct via
`request_known`. Request strings never become registered event names. Session,
parent-session and framework utilization remain explicitly unknown because the
current bindings do not provide them.
140 changes: 140 additions & 0 deletions docs/profiling.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,140 @@
<!--
SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
SPDX-License-Identifier: Apache-2.0
-->

# NVTX framework-pin tracing

This first instrumentation checkpoint traces framework pin requests used by the
remote-memory path. It separates the synchronous `request_pin` callback from the
asynchronous wait for its result. A shared pin has one lifetime and separate
associations to each waiting source operation.

Source NIXL submission, native completion, release and cancellation are also
traced; see the [source lifecycle reference](nvtx-events.md). Target processing,
router hints and caller polling remain separate work in this source checkpoint.

## Enable tracing

On Linux, install the optional Python bindings and structured-payload dependency:

```bash
uv sync --extra profiling
```

`KVCR_NVTX_LEVEL` is read when a KVCR remote-memory backend is constructed:

| Value | Behavior |
| --- | --- |
| `off` | No NVTX/NumPy import, pin trace objects, or payload construction. |
| `low` | Pin lifecycles and source NIXL writes. Default when the profiling dependencies are available. |
| `medium` | Low detail plus waiter-detachment and native release-retry events. |

Without the optional dependencies, tracing is a no-op. An explicit request for
tracing with an unavailable backend warns once. An unsupported level warns once
and disables tracing. There is no `high` level in this checkpoint.

Tracing is independent of KVCR telemetry and is not gated on profiler attachment.
Consequently, low/medium payload construction also costs work outside a profiler.
Use `off` as the baseline when measuring overhead. Annotation failures are kept
out of KVCR's result and resource-ownership paths.

## Capture a real pin and transfer

The example starts two real KVCR/NIXL UCX agents in **one process**, with registered
host memory and a small framework binding that completes pins after a configured
delay. It checks transferred bytes and pin release. It requires Linux and NIXL;
it does not run a model, Dynamo, vLLM, or a GPU transfer.

```bash
KVCR_NVTX_LEVEL=low nsys profile \
--trace=nvtx,cuda --sample=none --cpuctxsw=none \
--output=kvcr-pin-success \
uv run --extra profiling python examples/nvtx_pin_capture.py \
--blocks 2 --delay-ms 20

nsys export --type=sqlite --include-json=true \
--output=kvcr-pin-success.sqlite kvcr-pin-success.nsys-rep
```

For controlled failed results and timeouts, use `--scenario failure` and
`--scenario timeout` with distinct report names. Timeout deliberately delays the
pin beyond the 10-second operation deadline. These are correctness examples,
not throughput benchmarks. Keep the reports outside the source checkout.

Extended payloads require a collector/viewer with support for them. The initial
probe used Python 3.12, `nvtx` 0.2.16, NumPy 2.5.3, and Nsight Systems 2025.3.2.
Verify decoding in your environment before relying on a capture. The CUDA runtime
package `nvidia-nvtx` does not supply the Python `nvtx` annotation API.

In the SQLite export, inspect `NVTX_EVENTS.jsonText` (enabled by `--include-json`)
and join its `textId` to `StringIds.id` for the event name. For example:

```sql
SELECT e.start, e.end, s.value AS event, e.jsonText
FROM NVTX_EVENTS AS e JOIN StringIds AS s ON s.id = e.textId
WHERE s.value LIKE 'source.pin.%'
ORDER BY e.start;
```

## Interpretation

All annotations use the `KVCR` domain and `framework_pin` category. Event names
are bounded static strings. IDs and counts are payloads, never registered names.

| Event | Meaning |
| --- | --- |
| `source.pin.framework` | Same-thread push/pop around the framework's `request_pin` callback, including exceptions. |
| `source.pin.registered` | A new physical request was accepted into KVCR's pending-pin state. |
| `source.pin.waiter` | Association between that pin and one source operation/target operation handle. |
| `source.pin.completed` | KVCR observed a usable/failed result, deadline, cancellation, or shutdown. Exactly one terminal observation per pin trace. |
| `source.pin.detached` | One source operation stopped waiting; other waiters may remain. Medium detail only. |

Registration-to-completion measures the observed asynchronous wait, including
delay until KVCR polls the framework. It is not the framework's internal execution
time. A cancellation marker reports KVCR's decision; it does not establish native
transfer quiescence or permission to reuse a buffer. Late results are still
discarded/released by the existing lifecycle and do not emit a second completion.

Schema version 1 uses a fresh structured NumPy payload for every event:

| Field | Interpretation |
| --- | --- |
| `instance_hi`, `instance_lo` | Two uint64 halves of a UUID assigned to this KVCR tracer instance. |
| `pin_id` | Independent uint64 sequence within that instance; distinguishes framework request-ID reuse. |
| `pin_request_known`, `pin_request_id` | Signed int64 framework request ID and availability flag. The flag is zero before the callback returns or when its Python integer is outside int64 range; the independent pin identity and lifecycle events are still recorded. Zero is a valid request ID when the flag is set. |
| `source_op_id`, `op_handle` | Signed int64 source and target operation handles. Zero denotes unavailable context in direct helper calls. The source path carries these even if the pin callback fails before registration. |
| `requested_blocks`, `completed_blocks` | Requested count and count of non-missing entries in an accepted result. `-1` means the completed count is unavailable. |
| `status`, `reason` | Bounded codes below. |
| `fw_dram_utilization_known` | Always zero in this checkpoint: no framework utilization source is exposed by the current bindings. |

Join pin events by `(instance_hi, instance_lo, pin_id)`, not by framework request
ID alone. Operation handles alone are not globally unique across target agents.
This schema is a source-side pin checkpoint, not a cross-worker identity scheme.
No request, session, or parent-session identity is inferred from those handles.

Status codes: `1=pending`, `2=success`, `3=partial`, `4=failed`, `5=timeout`,
`6=cancelled`. Partial reports how many blocks were available, without asserting
why others were missing.

Reason codes: `0=unknown`, `1=none`, `2=callback_error`, `3=duplicate_request`,
`4=invalid_result`, `5=deadline`, `6=cancelled`, `7=shutdown`, `8=no_waiters`.
A framework result of `None` gives no failure cause: it is **not** evidence of
cache pressure. Callback errors identify the stage, not an underlying diagnosis.
Exception strings are not recorded.

## Integration gaps on the reviewed baseline

On `fb4f264`, `submit_hint()` accepts a `request_id`; the hint parser retains only
source endpoint and block hashes. The target pull retains the request ID, but its
`start_write` message does not transport request/session/parent-session context.
`KVCRBindings` exposes pin callbacks, not framework-cache utilization. Capturing
those inputs needs a verified integration source and potentially a separately
scoped API or wire change. KVCR local-cache occupancy is not a substitute.

The pin hooks can also be encountered by remote fetch. They do not provide full
fetch fan-out correlation. No production-overhead acceptance threshold has been
established by this example.

References: [NVTX Python best practices](https://nvidia.github.io/NVTX/python/best_practices.html)
and [extended payload examples](https://nvidia.github.io/NVTX/python/annotation_attributes.html).
195 changes: 195 additions & 0 deletions examples/nvtx_pin_capture.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,195 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
"""Capture framework pinning through two real KVCR/NIXL agents on Linux.

The framework binding is deliberately small: it returns registered host-memory
blocks after a configurable delay. This validates the pin instrumentation, not
Dynamo/vLLM integration, GPU transfers, or production performance.
"""

import argparse
import ctypes
import json
import socket
import time
from contextlib import ExitStack

from kvcr import KVCR, KVCRBindings
from kvcr.config import KVCRBackendConfigs, KVCRConfig, RemoteFWDramOptions
from kvcr.control_channels import ZmqPeerControlChannel
from kvcr.types import BlockKey, MemoryRef, PinRequestId, RegionDescriptor


def control_channel():
with socket.socket() as sock:
sock.bind(("127.0.0.1", 0))
port = sock.getsockname()[1]
return ZmqPeerControlChannel("127.0.0.1", port, "127.0.0.1")


class Framework:
def __init__(self, name, keys, delay, fail):
self.name, self.keys, self.delay, self.fail = name, keys, delay, fail
self.pending = {}
self.requests, self.releases, self.cancelled = 0, [], []

def request_pin(self, keys):
request = PinRequestId(self.requests)
self.requests += 1
self.pending[request] = (time.monotonic() + self.delay, tuple(keys))
return request

def poll_pin_results(self):
ready = []
for request, (deadline, keys) in list(self.pending.items()):
if time.monotonic() < deadline:
continue
del self.pending[request]
result = (
None
if self.fail
else (
f"pin-{request}",
{
key: [
MemoryRef(
end_point_name=self.name,
element_index=self.keys.index(key),
)
]
for key in keys
},
)
)
ready.append((request, result))
return ready

def release_pin(self, handle):
self.releases.append(handle)
return True

def cancel_pin_request(self, request):
self.cancelled.append(request)
self.pending.pop(request, None)


def run(blocks, delay_ms, scenario):
size = 4096
keys = tuple(BlockKey(f"block-{i}".encode()) for i in range(blocks))
expected = b"".join(bytes([i % 255 + 1]) * size for i in range(blocks))
source_memory = ctypes.create_string_buffer(expected, len(expected))
target_memory = ctypes.create_string_buffer(len(expected))
source_framework = Framework(
"nvtx-source",
keys,
20.0 if scenario == "timeout" else delay_ms / 1000,
scenario == "failure",
)
target_framework = Framework("nvtx-target", keys, 0, False)
channels = [control_channel(), control_channel()]
with ExitStack() as stack:
workers = []
for name, memory, framework, control in zip(
("nvtx-source", "nvtx-target"),
(source_memory, target_memory),
(source_framework, target_framework),
channels,
):
worker = KVCR(
KVCRConfig(
nixl_agent_name=name,
nixl_listen_port=0,
pool_layouts=[("", size)],
operation_timeout_ms=10_000,
abandon_timeout_ms=20_000,
),
KVCRBindings(
framework.request_pin,
framework.poll_pin_results,
framework.release_pin,
cancel_pin_request=framework.cancel_pin_request,
framework_control=control,
),
KVCRBackendConfigs(
framework_regions=[
RegionDescriptor(
addr=ctypes.addressof(memory),
size=size,
count=blocks,
)
],
remote_fw_dram=RemoteFWDramOptions(eager_ctrl_connect=False),
),
)
stack.callback(worker.close)
workers.append(worker)
source, target = workers
target.submit_hint(
{
"protocol_version": "0.1",
"actions": [
{
"action_type": "kv.fetch",
"action_version": "1.0",
"payload": {
"source_control_endpoint": channels[0].endpoint,
"block_hashes": list(range(blocks)),
},
}
],
},
request_id="nvtx-pin-example",
)
operation = target.deliver(
{
key: [MemoryRef(end_point_name="nvtx-target", element_index=i)]
for i, key in enumerate(keys)
},
request_id="nvtx-pin-example",
)
deadline = time.monotonic() + 30
results = {}
while time.monotonic() < deadline:
source.poll_completed()
results.update(target.poll_completed())
if operation in results and not source_framework.pending:
# Let source-side release follow target notification processing.
if scenario != "success" or source_framework.releases:
break
time.sleep(0.001)
assert operation in results, "delivery did not complete"
assert source_framework.requests == 1, "framework pin path was not exercised"
success = all(item.success for item in results[operation].values())
assert success == (scenario == "success"), results
if success:
assert target_memory.raw == expected, "transferred bytes differ"
assert source_framework.releases == ["pin-0"], "pin was not released once"
if scenario == "timeout":
assert source_framework.cancelled == [0]
print(
json.dumps(
{
"scenario": scenario,
"blocks": blocks,
"bytes": len(expected),
"success": success,
"pin_requests": source_framework.requests,
"pin_releases": source_framework.releases,
"cancelled": source_framework.cancelled,
"operation": operation,
}
)
)


if __name__ == "__main__":
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--blocks", type=int, default=2)
parser.add_argument("--delay-ms", type=float, default=20)
parser.add_argument(
"--scenario", choices=("success", "failure", "timeout"), default="success"
)
args = parser.parse_args()
if args.blocks < 1 or args.delay_ms < 0:
parser.error("blocks must be positive and delay-ms nonnegative")
run(args.blocks, args.delay_ms, args.scenario)
3 changes: 3 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,9 @@ dependencies = [
# cursor operations. It is loaded on first use, so importing kvcr without it
# works; running a KVCR-Service or claiming one of its pools does not.

[project.optional-dependencies]
profiling = ["nvtx>=0.2.16,<0.3", "numpy>=1.26,<3"]

[dependency-groups]
dev = [
"pytest>=8,<9",
Expand Down
Loading
Loading