Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
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
19 changes: 19 additions & 0 deletions .ci/scripts/run_tests_ucc_mpi.sh
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,25 @@ for MT in "" "-T"; do
echo "INFO: UCC MPI unit tests (CL/HIER+2step bcast) ... DONE"
done

# Team-cache correctness tests: exercise dormant reuse and eviction under
# pressure. UCC_TEAM_CACHE_MAX_SIZE=2 keeps the cache tiny so the overlapping
# subcommunicator test forces divergent per-rank eviction, the case the
# cross-rank agreement must reconcile.
echo "INFO: UCC team-cache correctness tests (world,half,odd_even) ..."
# shellcheck disable=SC2046,SC2086 # MPI argument fragments intentionally word-split.
cache_args=" -x UCC_TEAM_CACHE_ENABLE=y -x UCC_TEAM_CACHE_MAX_SIZE=2 -x UCC_TEAM_CACHE_CORRECTNESS_TESTS=y "
mpirun $(mpi_params $PPN) $ucx_tls_no_cuda_ipc $cache_args $EXE -c barrier -t world,half,odd_even
echo "INFO: UCC team-cache correctness tests (world,half,odd_even) ... DONE"

# Team-cache enabled-vs-disabled equivalence pass. The script runs ucc_test_mpi
# with cache on, then off, using its built-in per-collective correctness checks
# as the equivalence oracle.
echo "INFO: UCC team-cache enabled-vs-disabled equivalence test (8 ranks) ..."
# shellcheck disable=SC2086
MPIRUN="$(command -v mpirun)" EXE="${EXE%% *}" \
bash "${SCRIPT_DIR}/../../test/mpi/run_cache_equivalence.sh" 8
echo "INFO: UCC team-cache enabled-vs-disabled equivalence test (8 ranks) ... DONE"

end=`date +%s`

echo Tests took $((end - start)) seconds
95 changes: 95 additions & 0 deletions docs/user_guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -338,6 +338,101 @@ $ UCC_COLL_TRACE=INFO srun ./c/mpi/collective/osu_allreduce -i 1 -x 0 -d cuda -m
[1678205653.810705] [node_name:903 :0] ucc_coll.c:255 UCC_COLL INFO coll_init: Barrier; CL_BASIC {TL_UCP}, team_id 32768
```

## Team Cache (Experimental)

UCC supports an optional per-context communicator team cache that retains
`ucc_team_t` objects after `ucc_team_destroy` so that a subsequent
`ucc_team_create_post` with identical membership can re-adopt the same built
team instead of rebuilding it from scratch.

### Configuration knobs

| Environment variable | Default | Description |
|---|---|---|
| `UCC_TEAM_CACHE_ENABLE` | `n` | Enable the team cache. Off by default; opt-in. |
| `UCC_TEAM_CACHE_MAX_SIZE` | `128` | Maximum number of teams retained in the cache. Also clamped by `UCC_TEAM_IDS_POOL_SIZE`. |
| `UCC_TEAM_CACHE_EVICTION` | `fifo` | Eviction policy when the cache is full. `none`: never evict (new teams stay uncached). `fifo`: evict the oldest dormant entry (default). |
| `UCC_TEAM_CACHE_DISABLE_LINEAR_CHECK` | `n` | Trust the 64-bit membership hash alone in lookup, skipping the exact rank-array compare. Faster but unsafe on hash collision. |
| `UCC_TEAM_CACHE_DUMP_STATS` | `n` | Log hit/miss/eviction counters at context destroy. |
| `UCC_TEAM_CACHE_AGREEMENT` | `y` | Agree on the reuse decision across the members of every cacheable team create. Makes reuse safe for overlapping team scopes, at the cost of one small allreduce per create. |

### Cross-rank agreement

Each rank classifies a create as a cache hit or a miss from its own cache
contents, and those contents can diverge - for example when an eviction happens
on one rank only. Without agreement, the members of a single create could then
disagree on whether to re-adopt a dormant team or build a fresh one, and a create
where some ranks re-adopt while others rebuild does not make progress.

`UCC_TEAM_CACHE_AGREEMENT` (on by default) reconciles that with a small
`UCC_OP_BAND` allreduce over the members before any rank skips the address
exchange. Reuse happens only when every member independently classified the
create the same way; otherwise all members fall back to a fresh build. This makes
reuse safe even when team scopes overlap.

Disable the agreement only when team scopes never overlap - that is, when no rank
belongs to two simultaneously created teams with the same membership - and the
per-create allreduce is measurably too expensive. Applications that build only
disjoint or strictly nested communicators, such as a fixed set of row/column
communicators recreated over and over, satisfy that condition. Single-rank teams
never vote, since they cannot diverge.

If the vote itself fails (as opposed to being lost), `ucc_team_create_test`
returns the error and the handle is terminal. A handle the create allocated may
still be passed to `ucc_team_destroy`, which releases it. A handle that named a
cached team has already been handed back to the cache; `ucc_team_destroy` rejects
it and the caller must simply drop it. Either way, do not call
`ucc_team_create_test` on it again.

### Team-cache settings must be identical on every rank

> **The team-cache settings above are not per-rank tunables. A rank whose
> settings differ from its peers' can hang the job, not merely lose reuse.**

When caching and agreement are both on, a cacheable multi-rank create posts a
member-scoped allreduce (the *agreement vote*) so that every member reaches the
same reuse-vs-rebuild decision. A rank only enters that vote if all of the
following hold on that rank:

- `UCC_TEAM_CACHE_ENABLE=y`
- `UCC_TEAM_CACHE_AGREEMENT=y`
- the team is cacheable (no optional behavioral fields in `ucc_team_params_t`)
- the team has more than one member, and
- the caller passed `UCC_TEAM_PARAM_FIELD_EP_MAP`.

A rank that fails any of these skips the vote entirely and proceeds to build its
team directly. Its peers, meanwhile, have posted an allreduce that now has no
matching contribution from that rank and will never complete: the create hangs.

In practice this means:

- Set the team-cache variables in the launcher environment so every rank
inherits the same values (`mpirun -x UCC_TEAM_CACHE_ENABLE=y ...` or
`ucc.conf`). Do not set them from a per-rank wrapper script or from a rank
conditional.
- If a middleware creates some teams with `EP_MAP` and others without, that is
safe only when the choice is the same on every rank for a given team, which
it is for MPI communicators.
- If you must disable caching for part of a job, disable it for the whole job.

Setting `UCC_TEAM_CACHE_AGREEMENT=n` uniformly on every rank removes the vote
and with it this hazard, but it is only safe when communicator scopes never
overlap (see the table above).

### Requirements

- `UCC_TEAM_PARAM_FIELD_EP_MAP` must be set in `ucc_team_params_t` for a team
to be cacheable (it provides the membership the cache keys on).
- Teams with optional behavioral parameters (`ORDERING`, `OUTSTANDING_COLLS`,
`SYNC_TYPE`, `P2P_CONN`, `MEM_PARAMS`) are not cached because those
parameters are not part of the identity.

### Usage example

```bash
UCC_TEAM_CACHE_ENABLE=y UCC_TEAM_CACHE_MAX_SIZE=64 mpirun -np 8 ./my_app
```

## Known Issues

- For the CUDA and NCCL TL CUDA device dependent data structures are created when UCC
Expand Down
2 changes: 2 additions & 0 deletions src/Makefile.am
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ noinst_HEADERS = \
core/ucc_lib.h \
core/ucc_context.h \
core/ucc_team.h \
core/ucc_team_cache.h \
core/ucc_ee.h \
core/ucc_progress_queue.h \
core/ucc_service_coll.h \
Expand Down Expand Up @@ -115,6 +116,7 @@ libucc_la_SOURCES = \
core/ucc_version.c \
core/ucc_context.c \
core/ucc_team.c \
core/ucc_team_cache.c \
core/ucc_ee.c \
core/ucc_coll.c \
core/ucc_progress_queue.c \
Expand Down
7 changes: 4 additions & 3 deletions src/components/cl/hier/allgatherv/allgatherv.c
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ static inline ucc_status_t find_leader_rank(ucc_base_team_t *team,
ucc_assert(team_rank < UCC_CL_TEAM_SIZE(cl_team));
ucc_assert(SBGP_EXISTS(cl_team, NODE_LEADERS));

status = ucc_topo_get_node_leaders(core_team->topo, &node_leaders);
status = ucc_topo_get_node_leaders(UCC_TEAM_TOPO(core_team), &node_leaders);
if (UCC_OK != status) {
cl_error(team->context->lib, "Could not get node leaders");
return status;
Expand All @@ -69,7 +69,8 @@ static inline ucc_status_t find_leader_rank(ucc_base_team_t *team,
dst buffer is contiguous */
static inline ucc_status_t is_block_ordered(ucc_cl_hier_team_t *cl_team, int *ordered)
{
ucc_topo_t *topo = cl_team->super.super.params.team->topo;
ucc_topo_t *topo =
UCC_TEAM_TOPO(cl_team->super.super.params.team);
ucc_sbgp_t *all_nodes = NULL;
int is_block_ordered = 1;
int n_nodes;
Expand Down Expand Up @@ -129,7 +130,7 @@ UCC_CL_HIER_PROFILE_FUNC(ucc_status_t, ucc_cl_hier_allgatherv_init,
ucc_rank_t node_sbgp_size = SBGP_SIZE(cl_team, NODE);
ucc_rank_t leader_sbgp_size = SBGP_SIZE(cl_team, NODE_LEADERS);
ucc_rank_t team_size = UCC_CL_TEAM_SIZE(cl_team);
ucc_topo_t *topo = team->params.team->topo;
ucc_topo_t *topo = UCC_TEAM_TOPO(team->params.team);
ucc_aint_t *node_disps = NULL;
ucc_count_t *node_counts = NULL;
ucc_aint_t *leader_disps = NULL;
Expand Down
3 changes: 2 additions & 1 deletion src/components/cl/hier/allgatherv/unpack.c
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,8 @@ ucc_status_t ucc_cl_hier_allgatherv_unpack_start(ucc_coll_task_t *task)
ucc_rank_t *node_leaders = NULL;
ucc_sbgp_t *all_nodes = NULL;
ucc_sbgp_t *node_leaders_sbgp = NULL;
ucc_topo_t *topo = task->team->params.team->topo;
ucc_topo_t *topo = UCC_TEAM_TOPO(
task->team->params.team);
ucc_ee_executor_t *exec;
ucc_status_t status;
ucc_rank_t i;
Expand Down
2 changes: 1 addition & 1 deletion src/components/cl/hier/allreduce/allreduce_split_rail.c
Original file line number Diff line number Diff line change
Expand Up @@ -296,7 +296,7 @@ UCC_CL_HIER_PROFILE_FUNC(ucc_status_t, ucc_cl_hier_allreduce_split_rail_init,
return UCC_ERR_NOT_SUPPORTED;
}

if (!ucc_topo_isoppn(team->params.team->topo)) {
if (!ucc_topo_isoppn(UCC_TEAM_TOPO(team->params.team))) {
cl_debug(team->context->lib, "split_rail algorithm does not support "
"teams with non-uniform ppn across nodes");
return UCC_ERR_NOT_SUPPORTED;
Expand Down
4 changes: 2 additions & 2 deletions src/components/cl/hier/alltoallv/alltoallv.c
Original file line number Diff line number Diff line change
Expand Up @@ -52,8 +52,8 @@ static ucc_status_t ucc_cl_hier_alltoallv_finalize(ucc_coll_task_t *task)
for (_i = 0; _i < (_sbgp)->group_size; _i++) { \
_scount = ((_type *)(_coll_args)->args.src.info_v.counts)[_i]; \
_rcount = ((_type *)(_coll_args)->args.dst.info_v.counts)[_i]; \
_is_local = \
ucc_rank_on_local_node(_i, (_team)->params.team->topo); \
_is_local = ucc_rank_on_local_node( \
_i, UCC_TEAM_TOPO((_team)->params.team)); \
if ((_scount * _sdt_size > (_node_thresh)) && _is_local) { \
((_type *)_sc_full)[_i] = 0; \
} else { \
Expand Down
2 changes: 1 addition & 1 deletion src/components/cl/hier/bcast/bcast_2step.c
Original file line number Diff line number Diff line change
Expand Up @@ -146,7 +146,7 @@ ucc_cl_hier_bcast_2step_init_schedule(ucc_base_coll_args_t *coll_args,
if (SBGP_ENABLED(cl_team, NODE)) {
args.args.root = root_on_local_node
? find_root_node_rank(root, cl_team)
: core_team->topo->node_leader_rank_id;
: UCC_TEAM_TOPO(core_team)->node_leader_rank_id;
status =
ucc_coll_init(SCORE_MAP(cl_team, NODE), &args, &tasks[n_tasks]);
if (ucc_unlikely(UCC_OK != status)) {
Expand Down
7 changes: 4 additions & 3 deletions src/components/cl/hier/cl_hier_team.c
Original file line number Diff line number Diff line change
Expand Up @@ -47,13 +47,13 @@ UCC_CLASS_INIT_FUNC(ucc_cl_hier_team_t, ucc_base_context_t *cl_context,
ucc_tl_lib_t *tl_lib;
ucc_base_lib_attr_t attr;

if (!params->team->topo) {
if (!UCC_TEAM_TOPO(params->team)) {
cl_debug(cl_context->lib,
"can't create hier team without topology data");
return UCC_ERR_INVALID_PARAM;
}

if (ucc_topo_is_single_node(params->team->topo)) {
if (ucc_topo_is_single_node(UCC_TEAM_TOPO(params->team))) {
cl_debug(cl_context->lib, "skipping single node team");
return UCC_ERR_INVALID_PARAM;
}
Expand All @@ -65,7 +65,8 @@ UCC_CLASS_INIT_FUNC(ucc_cl_hier_team_t, ucc_base_context_t *cl_context,
for (i = 0; i < UCC_HIER_SBGP_LAST; i++) {
hs = &self->sbgps[i];
if (hs->state == UCC_HIER_SBGP_ENABLED) {
hs->sbgp = ucc_topo_get_sbgp(params->team->topo, hs->sbgp_type);
hs->sbgp = ucc_topo_get_sbgp(UCC_TEAM_TOPO(params->team),
hs->sbgp_type);
if (hs->sbgp->status != UCC_SBGP_ENABLED) {
/* SBGP of that type either not exists or the calling process
* is not part of subgroup
Expand Down
2 changes: 1 addition & 1 deletion src/components/tl/cuda/tl_cuda_team.c
Original file line number Diff line number Diff line change
Expand Up @@ -338,7 +338,7 @@ ucc_status_t ucc_tl_cuda_team_create_test(ucc_base_team_t *tl_team)
team->scratch.rem[i] = NULL;
}

if (!ucc_topo_has_device_info(UCC_TL_CORE_TEAM(team)->topo)) {
if (!ucc_topo_has_device_info(UCC_TEAM_TOPO(UCC_TL_CORE_TEAM(team)))) {
tl_debug(tl_team->context->lib,
"not all ranks have visible GPU device info; "
"skipping TL/CUDA team creation");
Expand Down
4 changes: 2 additions & 2 deletions src/components/tl/cuda/tl_cuda_team_topo.c
Original file line number Diff line number Diff line change
Expand Up @@ -342,7 +342,7 @@ static ucc_status_t
ucc_tl_cuda_team_topo_init_matrix(const ucc_tl_cuda_team_t *team,
ucc_rank_t *matrix)
{
ucc_topo_t *topo = UCC_TL_CORE_TEAM(team)->topo;
ucc_topo_t *topo = UCC_TEAM_TOPO(UCC_TL_CORE_TEAM(team));
ucc_proc_info_t *procs = topo->topo->procs;
ucc_device_id_t *dev_ids = topo->device_map.device_ids;
int size = UCC_TL_TEAM_SIZE(team);
Expand Down Expand Up @@ -406,7 +406,7 @@ ucc_status_t ucc_tl_cuda_team_topo_create(const ucc_tl_team_t *cuda_team,
* connectivity. This handles NVSwitch, fabric clique, and direct NVLink
* connections consistently and avoids rescanning the matrix for zeros. */
{
ucc_topo_t *utopo = UCC_TL_CORE_TEAM(team)->topo;
ucc_topo_t *utopo = UCC_TEAM_TOPO(UCC_TL_CORE_TEAM(team));
ucc_sbgp_t *node_sg = ucc_topo_get_sbgp(utopo, UCC_SBGP_NODE);
topo->is_fully_connected =
ucc_topo_is_nvlink_fully_connected(utopo, node_sg);
Expand Down
2 changes: 1 addition & 1 deletion src/components/tl/mlx5/tl_mlx5_team.c
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ static ucc_status_t ucc_tl_mlx5_topo_init(ucc_tl_mlx5_team_t *team)
ucc_subset_t subset;
ucc_status_t status;

status = ucc_ep_map_create_nested(&UCC_TL_CORE_TEAM(team)->ctx_map,
status = ucc_ep_map_create_nested(&UCC_TEAM_CTX_MAP(UCC_TL_CORE_TEAM(team)),
&UCC_TL_TEAM_MAP(team), &team->ctx_map);
if (UCC_OK != status) {
tl_debug(UCC_TL_TEAM_LIB(team), "failed to create ctx map");
Expand Down
3 changes: 2 additions & 1 deletion src/components/tl/sharp/tl_sharp_team.c
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,8 @@ UCC_CLASS_INIT_FUNC(ucc_tl_sharp_team_t, ucc_base_context_t *tl_context,
set.map = UCC_TL_TEAM_MAP(self);

if (UCC_TL_SHARP_TEAM_LIB(self)->cfg.use_internal_oob) {
status = ucc_ep_map_create_nested(&UCC_TL_CORE_TEAM(self)->ctx_map,
status =
ucc_ep_map_create_nested(&UCC_TEAM_CTX_MAP(UCC_TL_CORE_TEAM(self)),
&UCC_TL_TEAM_MAP(self),
&self->oob_ctx.subset.map);
if (status != UCC_OK) {
Expand Down
2 changes: 1 addition & 1 deletion src/components/tl/ucc_tl.c
Original file line number Diff line number Diff line change
Expand Up @@ -188,7 +188,7 @@ static ucc_status_t ucc_tl_is_reachable(const ucc_base_team_params_t *params,
for (i = 0; i < params->size; i++) {
rank = ucc_ep_map_eval(params->map, i);
if (use_ctx) {
rank = ucc_ep_map_eval(core_team->ctx_map, rank);
rank = ucc_ep_map_eval(UCC_TEAM_CTX_MAP(core_team), rank);
}
addr_header = UCC_ADDR_STORAGE_RANK_HEADER(addr_storage, rank);
for (j = 0; j < addr_header->n_components; j++) {
Expand Down
1 change: 1 addition & 0 deletions src/components/tl/ucp/tl_ucp_tag.h
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@
#define UCC_TL_UCP_MAX_COLL_TAG (UCC_TL_UCP_MAX_TAG - UCC_TL_UCP_RESERVED_TAGS)
#define UCC_TL_UCP_SERVICE_TAG (UCC_TL_UCP_MAX_COLL_TAG + 1)
#define UCC_TL_UCP_ACTIVE_SET_TAG (UCC_TL_UCP_MAX_COLL_TAG + 2)
/* Tags MAX_COLL_TAG+3..+7 are unused; team-cache voting reuses SERVICE_TAG */
#define UCC_TL_UCP_MAX_SENDER UCC_MASK(UCC_TL_UCP_SENDER_BITS)
#define UCC_TL_UCP_MAX_ID UCC_MASK(UCC_TL_UCP_ID_BITS)

Expand Down
3 changes: 2 additions & 1 deletion src/components/tl/ucp/tl_ucp_team.c
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,8 @@ static inline ucc_status_t ucc_tl_ucp_get_topo(ucc_tl_ucp_team_t *team)
return UCC_OK;
}

status = ucc_ep_map_create_nested(&UCC_TL_CORE_TEAM(team)->ctx_map,
/* The nested map aliases ctx_map, kept at a fixed address by the holder */
status = ucc_ep_map_create_nested(&UCC_TEAM_CTX_MAP(UCC_TL_CORE_TEAM(team)),
&UCC_TL_TEAM_MAP(team),
&team->ctx_map);
if (UCC_OK != status) {
Expand Down
43 changes: 43 additions & 0 deletions src/components/topo/ucc_topo.c
Original file line number Diff line number Diff line change
Expand Up @@ -311,6 +311,49 @@ void ucc_topo_cleanup(ucc_topo_t *topo)
}
}

#define UCC_TOPO_PREP_FATAL(_status) ((_status) == UCC_ERR_NO_MEMORY)

ucc_status_t ucc_topo_prepare_shared(ucc_topo_t *topo)
{
ucc_sbgp_t *sbgps;
ucc_rank_t *node_leaders;
ucc_status_t status;
int n_sbgps, i;

if (!topo) {
return UCC_OK;
}

/* Each sbgp records a terminal status, so a failure is never retried */
for (i = 0; i < UCC_SBGP_LAST; i++) {
(void)ucc_topo_get_sbgp(topo, (ucc_sbgp_type_t)i);
}

/* The all_* arrays stay NULL and retryable on failure, so report it */
status = ucc_topo_get_all_sockets(topo, &sbgps, &n_sbgps);
if (UCC_TOPO_PREP_FATAL(status)) {
return status;
}
status = ucc_topo_get_all_numas(topo, &sbgps, &n_sbgps);
if (UCC_TOPO_PREP_FATAL(status)) {
return status;
}
status = ucc_topo_get_all_nodes(topo, &sbgps, &n_sbgps);
if (UCC_TOPO_PREP_FATAL(status)) {
return status;
}

/* The node leaders map is only defined for multi-node teams */
if (topo->topo->nnodes > 1) {
status = ucc_topo_get_node_leaders(topo, &node_leaders);
if (UCC_TOPO_PREP_FATAL(status)) {
return status;
}
}

return UCC_OK;
}

ucc_sbgp_t *ucc_topo_get_sbgp(ucc_topo_t *topo, ucc_sbgp_type_t type)
{
if (topo->sbgps[type].status == UCC_SBGP_NOT_INIT) {
Expand Down
3 changes: 3 additions & 0 deletions src/components/topo/ucc_topo.h
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,9 @@ ucc_status_t ucc_topo_init(

void ucc_topo_cleanup(ucc_topo_t *subset_topo);

/* Materializes the lazily built topo fields so it can be shared read-only */
ucc_status_t ucc_topo_prepare_shared(ucc_topo_t *topo);

ucc_sbgp_t *ucc_topo_get_sbgp(ucc_topo_t *topo, ucc_sbgp_type_t type);

int ucc_topo_is_single_node(ucc_topo_t *topo);
Expand Down
Loading
Loading