Add host API for *_on_stream operations (#340)

* Add functional test for barrier_all_on_stream

* Add rocshmem_barrier_all_on_stream support for GDA and RO backends

Implements rocshmem_barrier_all_on_stream operation for
GPU Direct Access and Reverse Offload backends.

Previously, rocshmem_barrier_all_on_stream was only supported for IPC backend.

* Add functional test for rocshmem_broadcastmem_on_stream

* Add host-side rocshmem_broadcastmem_on_stream API

Implement stream-based broadcast collective operation

- Add rocshmem_broadcastmem_on_stream host API and kernel implementation
- Add functional test TeamBroadcastmemOnStreamTester with multi-stream
  support and correctness verification
- Use per-workgroup contexts to avoid contention across parallel streams

API:
rocshmem_broadcastmem_on_stream(team, dest, source, nelems, pe_root, stream)

* Add functional test for rocshmem_getmem_on_stream

* Add host-side rocshmem_getmem_on_stream API

Implement stream-based point-to-point RMA get operation

- Add rocshmem_getmem_on_stream host API and kernel implementation
- Support for asynchronous getmem operations on HIP streams
- Add backend support for GDA, RO, and IPC contexts
- Use work-group collective getmem for efficient memory transfer

API:
rocshmem_getmem_on_stream(dest, source, nelems, pe, stream)

(AI Assist)

* Add host-side rocshmem_putmem_on_stream API

- Add rocshmem_putmem_on_stream for asynchronous remote writes
- Support for concurrent RMA operations on HIP streams
- Add backend support for GDA, RO, and IPC contexts
- Use work-group device collective operation

API:
rocshmem_putmem_on_stream(dest, source, bytes, pe, stream)

(AI Assist)

* Add functional test for rocshmem_putmem_on_stream

* Add host-side rocshmem_putmem_signal_on_stream API

Enables asynchronous putmem operations with signaling on HIP streams.

The implementation includes:
- Kernel wrapper rocshmem_putmem_signal_kernel
- Host interface putmem_signal_on_stream method
- Context layer support across all backends (IPC, GDA, RO)
- Public API

Function signature:
void rocshmem_putmem_signal_on_stream(void *dest, const void *source,
                                      size_t bytes, uint64_t *sig_addr,
                                      uint64_t signal, int sig_op,
                                      int pe, hipStream_t stream);

* Add functional test for rocshmem_putmem_signal_on_stream

* Add host-side rocshmem_signal_wait_until_on_stream API

Enables asynchronous signal wait operations on HIP streams.

The implementation includes:
- Kernel wrapper rocshmem_signal_wait_until_kernel
- Host interface signal_wait_until_on_stream method
- Context layer support across all backends (IPC, GDA, RO)
- Native uint64_t support in wait_until API (generated from P2P_SYNC.py)

Function signature:
void rocshmem_signal_wait_until_on_stream(uint64_t *sig_addr, int cmp,
                                          uint64_t cmp_value,
                                          hipStream_t stream);

(AI Assist)

* Add functional test for rocshmem_signal_wait_until_on_stream

* Add documentation for stream API functions

This commit adds API documentation for the following host-side
stream functions:

- rocshmem_barrier_all_on_stream (collective routines)
- rocshmem_broadcastmem_on_stream (collective routines)
- rocshmem_getmem_on_stream (RMA operations)
- rocshmem_putmem_on_stream (RMA operations)
- rocshmem_putmem_signal_on_stream (signaling operations)
- rocshmem_signal_wait_until_on_stream (point-to-point sync)

The documentation includes function signatures, parameter descriptions,
and detailed explanations of asynchronous behavior and stream handling.

(AI Assist)

* Rename "bytes" -> "nelems"

* Add "_TEST_" to the variables used in tests

* Remove incorrect hipStreamDefault usage

hipStreamDefault is not a default stream. This is a flag.

If stream == nullptr, then just pass it to kernel. It will launch the kernel on the default stream

[ROCm/rocshmem commit: d0c8380650]
Cette révision appartient à :
Anatolii Rozanov
2025-12-09 15:55:46 +01:00
révisé par GitHub
Parent b9c172de16
révision f98c72d627
39 fichiers modifiés avec 2649 ajouts et 49 suppressions
+96
Voir le fichier
@@ -365,6 +365,102 @@ __host__ void rocshmem_alltoallmem_on_stream(rocshmem_team_t team, void *dest,
const void *source, size_t size,
hipStream_t stream);
/**
* @brief enqueues a broadcast collective operation on given stream.
*
* @param[in] team The team participating in the collective.
* @param[in] dest Destination address. Must be an address on the symmetric
* heap.
* @param[in] source Source address. Must be an address on the symmetric heap.
* @param[in] nelems Number of bytes to broadcast.
* @param[in] pe_root Root PE (relative to team) from which to broadcast.
* @param[in] stream HIP stream on which to enqueue the operation.
*
* @return void
*/
__host__ void rocshmem_broadcastmem_on_stream(rocshmem_team_t team, void *dest,
const void *source, size_t nelems,
int pe_root, hipStream_t stream);
/**
* @brief enqueues a getmem RMA operation on given stream.
*
* @param[in] dest Destination address. Must be an address on the symmetric
* heap.
* @param[in] source Source address. Must be an address on the symmetric heap.
* @param[in] nelems Size of the transfer in bytes.
* @param[in] pe PE of the remote process.
* @param[in] stream HIP stream on which to enqueue the operation.
*
* @return void
*/
__host__ void rocshmem_getmem_on_stream(void *dest, const void *source,
size_t nelems, int pe,
hipStream_t stream);
/**
* @brief enqueues a putmem RMA operation on given stream.
*
* @param[in] dest Destination address. Must be an address on the symmetric
* heap.
* @param[in] source Source address. Must be an address on the symmetric heap.
* @param[in] nelems Size of the transfer in bytes.
* @param[in] pe PE of the remote process.
* @param[in] stream HIP stream on which to enqueue the operation.
*
* @return void
*/
__host__ void rocshmem_putmem_on_stream(void *dest, const void *source,
size_t nelems, int pe,
hipStream_t stream);
/**
* @brief Perform a put operation with signal on a HIP stream.
*
* This routine initiates a remote memory transfer on a specified HIP stream.
* The source data is copied from the local PE to the remote PE's destination
* address. After the put operation completes, a signal operation is performed
* on a remote symmetric signal variable.
*
* @param[in] dest Destination address on the remote PE
* @param[in] source Source address on the local PE
* @param[in] nelems Size of the transfer in bytes
* @param[in] sig_addr Address of signal variable on the remote PE
* @param[in] signal Signal value to be written
* @param[in] sig_op Signal operation (ROCSHMEM_SIGNAL_SET or
* ROCSHMEM_SIGNAL_ADD)
* @param[in] pe PE number of the remote PE
* @param[in] stream HIP stream on which to enqueue the operation
*
* @return void
*/
__host__ void rocshmem_putmem_signal_on_stream(void *dest, const void *source,
size_t nelems,
uint64_t *sig_addr,
uint64_t signal, int sig_op,
int pe, hipStream_t stream);
/**
* @brief Wait on a signal variable until it satisfies the specified condition,
* with the operation enqueued on a HIP stream.
*
* This function blocks the calling thread until the signal variable at
* \p sig_addr satisfies the comparison condition (* \p sig_addr \p cmp
* \p cmp_value). The wait operation is executed asynchronously on the
* specified HIP stream.
*
* @param[in] sig_addr Address of the signal variable on the symmetric heap
* @param[in] cmp Comparison operator (e.g., ROCSHMEM_CMP_EQ,
* ROCSHMEM_CMP_GE, ROCSHMEM_CMP_NE, etc.)
* @param[in] cmp_value Value to compare against
* @param[in] stream HIP stream on which to enqueue the operation
*
* @return void
*/
__host__ void rocshmem_signal_wait_until_on_stream(uint64_t *sig_addr, int cmp,
uint64_t cmp_value,
hipStream_t stream);
/**
* @brief registers the arrival of a PE at a barrier.
* The caller is blocked until the synchronization is resolved.
+17
Voir le fichier
@@ -624,6 +624,23 @@ __global__ ATTR_NO_INLINE void rocshmem_alltoallmem_kernel(rocshmem_team_t team,
const void *source,
size_t size);
/**
* @brief kernel for performing a broadcast collective operation.
* Caller enqueues the kernel on given stream
*
* @param[in] team The team participating in the collective.
* @param[in] dest Destination address. Must be an address on the symmetric
* heap.
* @param[in] source Source address. Must be an address on the symmetric heap.
* @param[in] nelems Number of bytes to broadcast.
* @param[in] pe_root Root PE (relative to team) from which to broadcast.
*
* @return void
*/
__global__ ATTR_NO_INLINE void rocshmem_broadcastmem_kernel(
rocshmem_team_t team, void *dest, const void *source, size_t nelems,
int pe_root);
/**
* @brief perform a collective barrier between all PEs in the system.
* The caller is blocked until the barrier is resolved.
+46
Voir le fichier
@@ -576,6 +576,47 @@ __host__ size_t rocshmem_ulonglong_wait_until_some_vector(
unsigned long long *ivars, size_t nelems, size_t* indices, const int* status,
int cmp, unsigned long long val);
__device__ void rocshmem_uint64_wait_until(
uint64_t *ivars, int cmp, uint64_t val);
__device__ size_t rocshmem_uint64_wait_until_any(
uint64_t *ivars, size_t nelems, const int* status,
int cmp, uint64_t val);
__device__ void rocshmem_uint64_wait_until_all(
uint64_t *ivars, size_t nelems, const int* status,
int cmp, uint64_t val);
__device__ size_t rocshmem_uint64_wait_until_some(
uint64_t *ivars, size_t nelems, size_t* indices, const int* status,
int cmp, uint64_t val);
__device__ size_t rocshmem_uint64_wait_until_any_vector(
uint64_t *ivars, size_t nelems, const int* status,
int cmp, uint64_t val);
__device__ void rocshmem_uint64_wait_until_all_vector(
uint64_t *ivars, size_t nelems, const int* status,
int cmp, uint64_t val);
__device__ size_t rocshmem_uint64_wait_until_some_vector(
uint64_t *ivars, size_t nelems, size_t* indices, const int* status,
int cmp, uint64_t val);
__host__ void rocshmem_uint64_wait_until(
uint64_t *ivars, int cmp, uint64_t val);
__host__ size_t rocshmem_uint64_wait_until_any(
uint64_t *ivars, size_t nelems, const int* status,
int cmp, uint64_t val);
__host__ void rocshmem_uint64_wait_until_all(
uint64_t *ivars, size_t nelems, const int* status,
int cmp, uint64_t val);
__host__ size_t rocshmem_uint64_wait_until_some(
uint64_t *ivars, size_t nelems, size_t* indices, const int* status,
int cmp, uint64_t val);
__host__ size_t rocshmem_uint64_wait_until_any_vector(
uint64_t *ivars, size_t nelems, const int* status,
int cmp, uint64_t val);
__host__ void rocshmem_uint64_wait_until_all_vector(
uint64_t *ivars, size_t nelems, const int* status,
int cmp, uint64_t val);
__host__ size_t rocshmem_uint64_wait_until_some_vector(
uint64_t *ivars, size_t nelems, size_t* indices, const int* status,
int cmp, uint64_t val);
/**
* @name SHMEM_TEST
@@ -658,6 +699,11 @@ __device__ int rocshmem_ulonglong_test(
__host__ int rocshmem_ulonglong_test(
unsigned long long *ivars, int cmp, unsigned long long val);
__device__ int rocshmem_uint64_test(
uint64_t *ivars, int cmp, uint64_t val);
__host__ int rocshmem_uint64_test(
uint64_t *ivars, int cmp, uint64_t val);
} // namespace rocshmem
+31
Voir le fichier
@@ -1209,6 +1209,37 @@ __host__ void rocshmem_ctx_getmem_nbi(rocshmem_ctx_t ctx, void *dest,
__host__ void rocshmem_getmem_nbi(void *dest, const void *source,
size_t nelems, int pe);
/**
* @brief kernel for performing a getmem RMA operation.
* Caller enqueues the kernel on given stream
*
* @param[in] dest Destination address. Must be an address on the symmetric
* heap.
* @param[in] source Source address. Must be an address on the symmetric heap.
* @param[in] nelems Size of the transfer in bytes.
* @param[in] pe PE of the remote process.
*
* @return void
*/
__global__ ATTR_NO_INLINE void rocshmem_getmem_kernel(void *dest,
const void *source,
size_t nelems, int pe);
/**
* @brief kernel for performing a putmem RMA operation.
* Caller enqueues the kernel on given stream
*
* @param[in] dest Destination address. Must be an address on the symmetric
* heap.
* @param[in] source Source address. Must be an address on the symmetric heap.
* @param[in] nelems Size of the transfer in bytes.
* @param[in] pe PE of the remote process.
*
* @return void
*/
__global__ ATTR_NO_INLINE void rocshmem_putmem_kernel(void *dest,
const void *source,
size_t nelems, int pe);
} // namespace rocshmem
+29
Voir le fichier
@@ -619,6 +619,35 @@ __device__ ATTR_NO_INLINE uint64_t rocshmem_signal_fetch(const uint64_t *sig_add
__device__ ATTR_NO_INLINE uint64_t rocshmem_signal_fetch_wg(const uint64_t *sig_addr);
__device__ ATTR_NO_INLINE uint64_t rocshmem_signal_fetch_wave(const uint64_t *sig_addr);
/**
* @brief Kernel wrapper for putmem_signal operation on stream
*
* @param[in] dest Destination address on remote PE
* @param[in] source Source address on local PE
* @param[in] nelems Size of the transfer in bytes
* @param[in] sig_addr Address of signal variable on remote PE
* @param[in] signal Signal value to write
* @param[in] sig_op Signal operation (ROCSHMEM_SIGNAL_SET or
* ROCSHMEM_SIGNAL_ADD)
* @param[in] pe PE of the remote process
*
* @return void
*/
__global__ ATTR_NO_INLINE void rocshmem_putmem_signal_kernel(
void *dest, const void *source, size_t nelems, uint64_t *sig_addr,
uint64_t signal, int sig_op, int pe);
/**
* @brief Kernel wrapper for signal_wait_until operation on stream
*
* @param[in] sig_addr Address of signal variable on the symmetric heap
* @param[in] cmp Comparison operator
* @param[in] cmp_value Value to compare against
*
* @return void
*/
__global__ ATTR_NO_INLINE void rocshmem_signal_wait_until_kernel(
uint64_t *sig_addr, int cmp, uint64_t cmp_value);
} // namespace rocshmem