Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
d65d719
CORE: Add ctx-scoped service allreduce
bwestheimer Aug 13, 2026
4f9d6ae
CORE: Add team cache identity fingerprint
bwestheimer Aug 13, 2026
d2cfedb
CORE: Add team cache container and FIFO eviction
bwestheimer Aug 13, 2026
4f58c2f
CORE: Wire team cache into team lifecycle
bwestheimer Aug 13, 2026
70df29c
TEST: Cover team cache dormant reuse lifecycle
bwestheimer Aug 13, 2026
5baedf1
CORE: Add cross-rank team cache vote
bwestheimer Aug 13, 2026
1a18e5f
TEST: Cover team cache agreement vote
bwestheimer Aug 13, 2026
0130c74
CI: Add team cache equivalence pass
bwestheimer Aug 13, 2026
f51d40b
CORE: Add refcounted team artifacts holder
bwestheimer Aug 12, 2026
73d1047
CORE: Add team artifacts accessor macros
bwestheimer Aug 12, 2026
499ea5f
TOPO: Materialize topo state before sharing
bwestheimer Aug 13, 2026
b963601
CORE: Add chained hash buckets to team cache
bwestheimer Aug 13, 2026
f56e691
TEST: Cover team cache chained buckets
bwestheimer Aug 13, 2026
e0a4dcd
CORE: Add derived-team caching
bwestheimer Aug 13, 2026
0d0b664
TEST: Cover derived-team caching
bwestheimer Aug 13, 2026
e6a3f42
DOCS: Document team cache DERIVED knob
bwestheimer Aug 13, 2026
eeed184
CORE: Add derived-team RESEAT on cid drift
bwestheimer Aug 13, 2026
8b4fccf
CL/: Implement update_id fanout for RESEAT
bwestheimer Aug 13, 2026
411ea66
TEST: Cover RESEAT path and document knob
bwestheimer Aug 13, 2026
3a54137
CORE: Add team cache LFU eviction policy
bwestheimer Aug 12, 2026
f4efa70
TEST: Cover LFU eviction and vote cookie
bwestheimer Aug 12, 2026
00acaa6
CORE: Return error when ctx service team absent
bwestheimer Sep 17, 2026
af0a357
CORE: Disable team cache when vote unavailable
bwestheimer Sep 17, 2026
ff45af1
CORE: Fail ctx create when team cache init fails
bwestheimer Sep 28, 2026
548f6fc
CORE: Track internal OOB ownership on teams
bwestheimer Sep 28, 2026
b79d520
CORE: Roll back team cache vote on failure
bwestheimer Sep 28, 2026
645f97f
CORE: Keep team id on failed cached teardown
bwestheimer Sep 28, 2026
8bdb2de
CORE: Progress pending cache teardown in ctx
bwestheimer Sep 28, 2026
24843e4
TEST: Keep poisoned cb box live across reuse
bwestheimer Sep 28, 2026
31f213c
CORE: Restrict RESEAT to external team ids
bwestheimer Sep 28, 2026
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
28 changes: 28 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,34 @@ 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"

# Second pass with UCC_TEAM_CACHE_RESEAT=y. The knob is off by default, so
# without this leg the reseat path (derived_reuse[drift] and the update_id
# fanout it drives) skips itself and ships untested.
echo "INFO: UCC team-cache correctness tests (RESEAT) ..."
# shellcheck disable=SC2046,SC2086 # MPI argument fragments intentionally word-split.
reseat_args="${cache_args} -x UCC_TEAM_CACHE_RESEAT=y "
mpirun $(mpi_params $PPN) $ucx_tls_no_cuda_ipc $reseat_args $EXE -c barrier -t world,half,odd_even
echo "INFO: UCC team-cache correctness tests (RESEAT) ... 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
147 changes: 147 additions & 0 deletions docs/user_guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -338,6 +338,153 @@ $ 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. |
| `UCC_TEAM_CACHE_DERIVED` | `y` | Reuse the shared artifacts of a still-live cached team when a create duplicates its membership, as `MPI_Comm_dup` does. No effect unless `UCC_TEAM_CACHE_ENABLE=y`. |
| `UCC_TEAM_CACHE_RESEAT` | `n` | Experimental: recover reuse under context-id drift by re-adopting a cached dormant derived team of identical membership but a different external id, re-seating its id and tag domain instead of rebuilding it. Requires `UCC_TEAM_CACHE_DERIVED=y`. |

### Derived teams

A create whose membership matches a team that is still live cannot re-adopt that
team, since the original is still in use - `MPI_Comm_dup` is the common case. With
`UCC_TEAM_CACHE_DERIVED` on, the new team is instead built as a *derived* team: it
draws its own team id, and therefore its own tag and sequence-number domain, but
borrows the live parent's membership map and topology instead of rebuilding them.
That skips the address exchange and the topology build, which is the expensive
part of a create.

The borrowed state is reference counted, so the parent and every derived team may
be destroyed in any order. A derived team is itself cacheable and can later be
re-adopted like any other dormant team. Under `UCC_TEAM_CACHE_AGREEMENT` the
derive decision is voted on like any other, and a member that cannot derive forces
all members to fall back to a full build.

Turn this off to make every such create an independent full build.

### Re-seating under context-id drift (experimental)

A cached team is keyed on its membership *and* its external id, which for MPI is
the communicator context id. Some workloads never reuse a context id: each
create/free cycle over the same ranks draws a fresh one, so the id drifts and
the exact-identity lookup misses every time even though a perfectly good dormant
team of that membership is sitting in the cache.

`UCC_TEAM_CACHE_RESEAT=y` recovers reuse in that case. When the exact lookup
misses, the cache is searched a second time for a dormant *derived* team of
identical membership, ignoring the external id. A match is re-adopted and
*re-seated*: its team id, and with it the tag and sequence-number domain of its
service team and of every CL and TL team beneath it, is moved to the caller's
new external id. Only derived teams are eligible, since only they hold borrowed
artifacts that make the move cheap relative to a full build.

Re-seating is the reason the id must be pushed all the way down. A team whose
core id changed while a TL team below it still addressed the retired id would
put two logically distinct teams in one tag domain, which is exactly the
aliasing this path exists to avoid.

This knob is experimental and off by default. It requires
`UCC_TEAM_CACHE_DERIVED=y`; with derived teams off there is nothing eligible to
re-seat. Under `UCC_TEAM_CACHE_AGREEMENT` the re-seat is voted on like any other
reuse, and the vote carries the candidate's instance cookie so that all members
re-seat the same team or none do.

Re-seating moves a team between two *external* ids only. A create that does not
pass `UCC_TEAM_PARAM_FIELD_ID` never re-seats, and a dormant derived team whose id
came from the internal pool is never a re-seat candidate: the pool owns that id,
and re-keying the team would orphan it. Such teams still take part in exact reuse.

### 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
3 changes: 3 additions & 0 deletions src/components/base/ucc_base_iface.h
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,8 @@ typedef struct ucc_base_team_iface {
ucc_status_t (*create_test)(ucc_base_team_t *team);
ucc_status_t (*destroy)(ucc_base_team_t *team);
ucc_get_coll_scores_fn_t get_scores;
/* Optional: re-seats the team id and tag domain in place */
void (*update_id)(ucc_base_team_t *team, uint16_t id);
} ucc_base_team_iface_t;

enum {
Expand Down Expand Up @@ -264,6 +266,7 @@ typedef struct ucc_base_coll_alg_info {
.super.team.create_test = ucc_##_f##_name##_team_create_test, \
.super.team.destroy = ucc_##_f##_name##_team_destroy, \
.super.team.get_scores = ucc_##_f##_name##_team_get_scores, \
.super.team.update_id = NULL, \
.super.coll.init = ucc_##_f##_name##_coll_init, \
.super.alg_info = {NULL}}; \
UCC_CONFIG_REGISTER_TABLE_ENTRY(&ucc_##_f##_name.super._f##lib_config, \
Expand Down
9 changes: 9 additions & 0 deletions src/components/cl/basic/cl_basic.c
Original file line number Diff line number Diff line change
Expand Up @@ -58,4 +58,13 @@ ucc_status_t ucc_cl_basic_coll_init(ucc_base_coll_args_t *coll_args,

ucc_status_t ucc_cl_basic_team_get_scores(ucc_base_team_t *cl_team,
ucc_coll_score_t **score);

void ucc_cl_basic_team_update_id(ucc_base_team_t *cl_team, uint16_t id);

UCC_CL_IFACE_DECLARE(basic, BASIC);

__attribute__((constructor)) static void cl_basic_iface_init(void)
{
/* UCC_BASE_IFACE_DECLARE leaves update_id NULL */
ucc_cl_basic.super.team.update_id = ucc_cl_basic_team_update_id;
}
15 changes: 15 additions & 0 deletions src/components/cl/basic/cl_basic_team.c
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,21 @@ UCC_CLASS_CLEANUP_FUNC(ucc_cl_basic_team_t)
UCC_CLASS_DEFINE_DELETE_FUNC(ucc_cl_basic_team_t, ucc_base_team_t);
UCC_CLASS_DEFINE(ucc_cl_basic_team_t, ucc_cl_team_t);

void ucc_cl_basic_team_update_id(ucc_base_team_t *cl_team, uint16_t id)
{
ucc_cl_basic_team_t *team = ucc_derived_of(cl_team, ucc_cl_basic_team_t);
ucc_tl_team_t *tl;
unsigned i;

/* Only UCP TLs expose scoll.update_id, so the others are skipped */
for (i = 0; i < team->n_tl_teams; i++) {
tl = team->tl_teams[i];
if (tl && UCC_TL_TEAM_IFACE(tl)->scoll.update_id) {
UCC_TL_TEAM_IFACE(tl)->scoll.update_id(&tl->super, id);
}
}
}

ucc_status_t ucc_cl_basic_team_destroy(ucc_base_team_t *cl_team)
{
ucc_cl_basic_team_t *team = ucc_derived_of(cl_team, ucc_cl_basic_team_t);
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
4 changes: 4 additions & 0 deletions src/components/cl/hier/cl_hier.c
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@
ucc_status_t ucc_cl_hier_get_lib_attr(const ucc_base_lib_t *lib,
ucc_base_lib_attr_t *base_attr);

void ucc_cl_hier_team_update_id(ucc_base_team_t *cl_team, uint16_t id);

ucc_status_t ucc_cl_hier_get_lib_properties(ucc_base_lib_properties_t *prop);

ucc_status_t ucc_cl_hier_get_context_attr(const ucc_base_context_t *context,
Expand Down Expand Up @@ -128,4 +130,6 @@ __attribute__((constructor)) static void cl_hier_iface_init(void)
ucc_cl_hier_bcast_algs;
ucc_cl_hier.super.alg_info[ucc_ilog2(UCC_COLL_TYPE_ALLGATHERV)] =
ucc_cl_hier_allgatherv_algs;
/* UCC_BASE_IFACE_DECLARE leaves update_id NULL */
ucc_cl_hier.super.team.update_id = ucc_cl_hier_team_update_id;
}
Loading
Loading