diff --git a/cpp/monoprop/Evolution.cpp b/cpp/monoprop/Evolution.cpp index 47bf3234..296082cf 100644 --- a/cpp/monoprop/Evolution.cpp +++ b/cpp/monoprop/Evolution.cpp @@ -14,6 +14,7 @@ #include "monoprop/Evolution.h" +#include #include #include #include @@ -60,8 +61,7 @@ auto combine_endpoint_contrib(const EndpointContrib &a, const EndpointContrib &b struct FlatExchangeBuffers { VecD send_buffer; VecD recv_buffer; - std::vector recv_counts; - std::vector recv_displs; + LayerExchangeLayout layout; }; auto &acquire_flat_exchange_buffers() { @@ -72,22 +72,18 @@ auto &acquire_flat_exchange_buffers() { return scratch.buffers; } -void resize_flat_exchange_buffers(const LayerExchangeLayout &layout, FlatExchangeBuffers &buffers) { - // Recv size isn't known until counts are exchanged; keep it at 1 element so data() stays non-null. - const size_t send_alloc = layout.total_count == 0 ? 1 : layout.total_count; - buffers.send_buffer.resize(send_alloc); - buffers.recv_buffer.resize(1); - buffers.recv_counts.clear(); - buffers.recv_displs.clear(); +// A property of the communicator, not the layer: all ranks participate even at local total_count 0. +auto layer_exchange_participates(const mpi::Comm &comm) -> bool { + return mpi::size(comm) != 1; } -auto active_evolution_exchange_layout(const LayerTraversal &layer, const mpi::Comm &comm) - -> const LayerExchangeLayout * { - if (mpi::size(comm) == 1) { - return nullptr; - } - // All ranks must participate even at local total_count 0, else MPI_Alltoallv deadlocks. - return &layer.evolution_exchange_layout(); +// Derives both sides at once: the count matrix is symmetric, so the recv layout is the send layout. +auto derive_layer_exchange(const LayerTraversal &layer, const mpi::Comm &comm, int scale, LayerExchangeLayout &layout) + -> void { + const auto my_rank = static_cast(mpi::rank(comm)); + const char *what = scale == 1 ? "Layer exchange" : "Layer derivative exchange"; + detail::derive_exchange_layout(layer.cross_rank(), my_rank, scale, layout, what); + mpi::check_exchange_layout_width(layout.counts, comm); } // The completed alltoallv payload as an apply pass sees it: peer `rank`'s entries start at @@ -105,21 +101,19 @@ struct CrossRankExchangeHandle { [[no_unique_address]] mpi::Ticket ticket; }; -inline auto begin_flat_exchange(const LayerExchangeLayout &layout, FlatExchangeBuffers &buffers, const mpi::Comm &comm) - -> CrossRankExchangeHandle { +inline auto begin_flat_exchange(FlatExchangeBuffers &buffers, const mpi::Comm &comm) -> CrossRankExchangeHandle { + const LayerExchangeLayout &layout = buffers.layout; CrossRankExchangeHandle handle; handle.layout = &layout; handle.buffers = &buffers; - const auto &recv = mpi::resolve_recv(layout.counts, comm, layout.recv_cache); - buffers.recv_counts = recv.counts; - buffers.recv_displs = recv.displs; - buffers.recv_buffer.resize(recv.total == 0 ? 1 : static_cast(recv.total)); + buffers.recv_buffer.resize(layout.total_count == 0 ? 1 : layout.total_count); + // Same arrays on both sides: MPI reads recvcounts/recvdispls and never writes them. handle.ticket = mpi::post_flat_alltoallv({.send = buffers.send_buffer.data(), .send_counts = layout.counts.data(), .send_displs = layout.displs.data(), .recv = buffers.recv_buffer.data(), - .recv_counts = buffers.recv_counts.data(), - .recv_displs = buffers.recv_displs.data()}, + .recv_counts = layout.counts.data(), + .recv_displs = layout.displs.data()}, mpi::size(comm), comm); return handle; @@ -136,19 +130,18 @@ struct InFlightExchange { bool active = false; }; -// Callers run the participation guard first (see active_evolution_exchange_layout) so the layout is -// never materialized at a single rank. template -inline auto begin_layer_exchange(const LayerExchangeLayout &layout, const mpi::Comm &comm, Pack pack) +inline auto begin_layer_exchange(const LayerTraversal &layer, int scale, const mpi::Comm &comm, Pack pack) -> InFlightExchange { InFlightExchange in_flight; in_flight.my_rank = mpi::rank(comm); in_flight.active = true; auto &buffers = acquire_flat_exchange_buffers(); - resize_flat_exchange_buffers(layout, buffers); - pack(in_flight.my_rank, layout, buffers.send_buffer); - in_flight.handle = begin_flat_exchange(layout, buffers, comm); + derive_layer_exchange(layer, comm, scale, buffers.layout); + buffers.send_buffer.resize(buffers.layout.total_count == 0 ? 1 : buffers.layout.total_count); + pack(in_flight.my_rank, buffers.layout, buffers.send_buffer); + in_flight.handle = begin_flat_exchange(buffers, comm); return in_flight; } @@ -167,13 +160,12 @@ inline auto finish_layer_exchange(InFlightExchange &in_flight, Apply apply) } wait_flat_exchange(in_flight.handle); return apply(ExchangePayload{.recv_buffer = in_flight.handle.buffers->recv_buffer, - .recv_displs = in_flight.handle.buffers->recv_displs, + .recv_displs = in_flight.handle.buffers->layout.displs, .my_rank = in_flight.my_rank}); } -// Per-thread pre-cos snapshot buffers, reused across layers to avoid a malloc/free per layer. Passed whole -// to the pack and apply passes: each reads two of the four vectors, and re-splitting them at every hand-off -// is what lets a send buffer be mistaken for a recv one. +// Per-thread pre-cos snapshot buffers, reused across layers, indexed by occupied position. Passed +// whole rather than re-split at each hand-off, which is how a send buffer becomes a recv one. struct DerivativeSnapshotScratch { std::vector sin_send_state; std::vector sin_send_op; @@ -192,23 +184,20 @@ void pack_cross_rank_derivative_payload_impl(const DerivativeSnapshotScratch &sn int my_rank, const LayerExchangeLayout &layout, VecD &send_buffer) { - const size_t num_ranks = layer.cross_rank_rank_count(); - for (size_t rank = 0; rank < num_ranks; ++rank) { - if (static_cast(rank) == my_rank) { - continue; - } - const size_t end = layer.cross_rank_sin_send_size(rank); - if (end == 0) { - continue; - } - const auto base = static_cast(layout.displs[rank]); - const auto &bs = snap.sin_send_state[rank]; - const auto &bh = snap.sin_send_op[rank]; - layer.for_each_cross_rank_sin_send_range(rank, 0, end, [&send_buffer, &base, &bs, &bh](size_t k, size_t /*i*/) { - send_buffer[base + 2 * k] = bs[k]; - send_buffer[base + 2 * k + 1] = bh[k]; + layer.for_each_occupied_slot( + [my_rank, &layout, &snap, &send_buffer](size_t pos, size_t rank, const detail::CrossRankSlotView &slot) { + if (static_cast(rank) == my_rank) { + return; + } + // The send layout stays dense in the world; only the snapshot is indexed by position. + const auto base = static_cast(layout.displs[rank]); + const auto &bs = snap.sin_send_state[pos]; + const auto &bh = snap.sin_send_op[pos]; + for (size_t k = 0; k < slot.sin_send_count; ++k) { + send_buffer[base + 2 * k] = bs[k]; + send_buffer[base + 2 * k + 1] = bh[k]; + } }); - } } // Remote endpoint pass: own pre-cos values come from the sin_recv snapshots, partner values from the @@ -219,25 +208,18 @@ auto apply_cross_rank_derivative_exchange_impl(VecD &state, const DerivativeSnapshotScratch &snap, const TrigValues &trig, const ExchangePayload &payload) -> EndpointContrib { - const size_t num_ranks = layer.cross_rank_rank_count(); EndpointContrib local{}; - for (size_t rank = 0; rank < num_ranks; ++rank) { - if (static_cast(rank) == payload.my_rank) { - continue; - } - const size_t end = layer.cross_rank_sin_recv_size(rank); - if (end == 0) { - continue; - } - const auto *rv = payload.recv_buffer.data() + payload.recv_displs[rank]; - const auto &ds = snap.sin_recv_state[rank]; - const auto &dh = snap.sin_recv_op[rank]; - layer.for_each_cross_rank_sin_recv_range( - rank, - 0, - end, - [&trig, &ds, &dh, &rv, &local, &op, &state](size_t k, size_t i, int phi_signed) { - const auto phi = static_cast(phi_signed); + layer.for_each_occupied_slot( + [&payload, &snap, &trig, &op, &state, &local](size_t pos, size_t rank, const detail::CrossRankSlotView &slot) { + if (static_cast(rank) == payload.my_rank) { + return; + } + const auto *rv = payload.recv_buffer.data() + payload.recv_displs[rank]; + const auto &ds = snap.sin_recv_state[pos]; + const auto &dh = snap.sin_recv_op[pos]; + for (size_t k = 0; k < slot.sin_send_count; ++k) { + const size_t i = detail::slot_sin_recv_index(slot, k); + const auto phi = static_cast(detail::slot_sin_recv_phase(slot, k)); // Inverse-rotation write-back (−sin): un-evolves state/op for the next reverse layer. const double ps = -trig.sin_val * phi; const double s_old = ds[k]; @@ -248,8 +230,8 @@ auto apply_cross_rank_derivative_exchange_impl(VecD &state, local.sin_terms += phi * s_old * h_p; op[i] = (h_old * trig.cos_val) + (ps * h_p); state[i] = (s_old * trig.cos_val) + (ps * s_p); - }); - } + } + }); return local; } @@ -258,11 +240,13 @@ inline auto begin_cross_rank_derivative_exchange(const DerivativeSnapshotScratch const LayerTraversal &layer, const mpi::Comm &comm) -> InFlightExchange { // Single-rank (or no peer participating): nothing to exchange — the self slot covers everything. - if (active_evolution_exchange_layout(layer, comm) == nullptr) { + if (!layer_exchange_participates(comm)) { return {}; } // Safe to fire before the cos pass: pack reads pre-cos snapshots and the transfer touches only buffers. - return begin_layer_exchange(layer.derivative_exchange_layout(), + // Scale 2: each rotation endpoint carries both the op and the state payload. + return begin_layer_exchange(layer, + 2, comm, [&snap, &layer](int my_rank, const LayerExchangeLayout &layout, VecD &send_buffer) { pack_cross_rank_derivative_payload_impl(snap, layer, my_rank, layout, send_buffer); @@ -286,20 +270,16 @@ void pack_cross_rank_evolution_payload_impl(VecD &op, int my_rank, const LayerExchangeLayout &layout, VecD &send_buffer) { - const size_t num_ranks = layer.cross_rank_rank_count(); - for (size_t rank = 0; rank < num_ranks; ++rank) { - if (static_cast(rank) == my_rank) { - continue; - } - const size_t end = layer.cross_rank_sin_send_size(rank); - if (end == 0) { - continue; - } - const auto base = static_cast(layout.displs[rank]); - layer.for_each_cross_rank_sin_send_range(rank, 0, end, [&send_buffer, &base, &op](size_t k, size_t i) { - send_buffer[base + k] = op[i]; + layer.for_each_occupied_slot( + [my_rank, &layout, &send_buffer, &op](size_t rank, const detail::CrossRankSlotView &slot) { + if (static_cast(rank) == my_rank) { + return; + } + const auto base = static_cast(layout.displs[rank]); + for (size_t k = 0; k < slot.sin_send_count; ++k) { + send_buffer[base + k] = op[detail::slot_sin_send_index(slot, k)]; + } }); - } } void apply_cross_rank_evolution_exchange_impl(VecD &op, @@ -307,33 +287,26 @@ void apply_cross_rank_evolution_exchange_impl(VecD &op, double sin_val, const ExchangePayload &payload) { // op[i] is already cos-scaled, so only the sine term is added; rv[k] is the partner's pre-cos value. - const size_t num_ranks = layer.cross_rank_rank_count(); - for (size_t rank = 0; rank < num_ranks; ++rank) { + layer.for_each_occupied_slot([&payload, sin_val, &op](size_t rank, const detail::CrossRankSlotView &slot) { if (static_cast(rank) == payload.my_rank) { - continue; - } - const size_t end = layer.cross_rank_sin_recv_size(rank); - if (end == 0) { - continue; + return; } const auto *rv = payload.recv_buffer.data() + payload.recv_displs[rank]; - layer.for_each_cross_rank_sin_recv_range(rank, - 0, - end, - [&op, &sin_val, &rv](size_t k, size_t i, int phi_signed) { - op[i] += sin_val * static_cast(phi_signed) * rv[k]; - }); - } + for (size_t k = 0; k < slot.sin_send_count; ++k) { + const size_t i = detail::slot_sin_recv_index(slot, k); + op[i] += sin_val * static_cast(detail::slot_sin_recv_phase(slot, k)) * rv[k]; + } + }); } inline auto begin_cross_rank_evolution_exchange(VecD &op, const LayerTraversal &layer, const mpi::Comm &comm) -> InFlightExchange { - const auto *layout = active_evolution_exchange_layout(layer, comm); - if (layout == nullptr) { + if (!layer_exchange_participates(comm)) { return {}; } return begin_layer_exchange( - *layout, + layer, + 1, comm, [&op, &layer](int my_rank, const LayerExchangeLayout &active_layout, VecD &send_buffer) { pack_cross_rank_evolution_payload_impl(op, layer, my_rank, active_layout, send_buffer); @@ -352,22 +325,20 @@ inline auto finish_cross_rank_evolution_exchange(VecD &op, // Snapshot-free self-slot endpoint pass: sin_recv entries k and k+pairs are the two endpoints of one // rotation, so reading both (pre-cos recovered from the post-cos slots) before writing either avoids the // read-after-write hazard. -auto apply_self_slot_derivative_paired(VecD &state, - VecD &op, - const LayerTraversal &layer, - size_t my_rank, - const TrigValues &trig) -> EndpointContrib { - const size_t self_d_count = layer.cross_rank_sin_recv_size(my_rank); +auto apply_self_slot_derivative_paired(VecD &state, VecD &op, const LayerTraversal &layer, const TrigValues &trig) + -> EndpointContrib { + const auto slot = layer.cross_rank_self_slot(); + const size_t self_d_count = slot.sin_send_count; if (self_d_count == 0) { return {}; } const auto pairs = self_d_count / 2; EndpointContrib local{}; for (size_t k = 0; k < pairs; ++k) { - const size_t i1 = layer.cross_rank_sin_recv_index_at(my_rank, k); - const double phi1 = static_cast(layer.cross_rank_sin_recv_phase_at(my_rank, k)); - const size_t i2 = layer.cross_rank_sin_recv_index_at(my_rank, k + pairs); - const auto phi2 = static_cast(layer.cross_rank_sin_recv_phase_at(my_rank, k + pairs)); + const size_t i1 = detail::slot_sin_recv_index(slot, k); + const auto phi1 = static_cast(detail::slot_sin_recv_phase(slot, k)); + const size_t i2 = detail::slot_sin_recv_index(slot, k + pairs); + const auto phi2 = static_cast(detail::slot_sin_recv_phase(slot, k + pairs)); // Recover pre-cos values. const double s1 = state[i1] * trig.sec_val; const double h1 = op[i1] * trig.cos_val; @@ -386,52 +357,59 @@ auto apply_self_slot_derivative_paired(VecD &state, return local; } +// Emptied, not filled: the self slot recovers live, and a previous layer's snapshot must not be read. +void clear_slot_snapshot(DerivativeSnapshotScratch &snap, size_t pos) { + snap.sin_send_state[pos].clear(); + snap.sin_send_op[pos].clear(); + snap.sin_recv_state[pos].clear(); + snap.sin_recv_op[pos].clear(); +} + +void fill_slot_snapshot(DerivativeSnapshotScratch &snap, + size_t pos, + const VecD &state, + const VecD &op, + const detail::CrossRankSlotView &slot) { + const size_t count = slot.sin_send_count; + auto &bs = snap.sin_send_state[pos]; + auto &bh = snap.sin_send_op[pos]; + auto &ds = snap.sin_recv_state[pos]; + auto &dh = snap.sin_recv_op[pos]; + bs.resize(count); + bh.resize(count); + ds.resize(count); + dh.resize(count); + for (size_t k = 0; k < count; ++k) { + const size_t bi = detail::slot_sin_send_index(slot, k); + bs[k] = state[bi]; + bh[k] = op[bi]; + const size_t di = detail::slot_sin_recv_index(slot, k); + ds[k] = state[di]; + dh[k] = op[di]; + } +} + // Snapshot pre-cos (state, op) at every remote rank's sin_send/sin_recv endpoints before the cos pass // clobbers them. The self slot needs none — it recovers live — so it is cleared. void snapshot_remote_endpoints(const VecD &state, const VecD &op, const LayerTraversal &layer, size_t my_rank, - size_t R, DerivativeSnapshotScratch &snap) { - snap.sin_send_state.resize(R); - snap.sin_send_op.resize(R); - snap.sin_recv_state.resize(R); - snap.sin_recv_op.resize(R); - for (size_t r = 0; r < R; ++r) { - if (r == my_rank) { - snap.sin_send_state[r].clear(); - snap.sin_send_op[r].clear(); - snap.sin_recv_state[r].clear(); - snap.sin_recv_op[r].clear(); - continue; - } - const size_t bc = layer.cross_rank_sin_send_size(r); - snap.sin_send_state[r].resize(bc); - snap.sin_send_op[r].resize(bc); - if (bc > 0) { - auto &bs = snap.sin_send_state[r]; - auto &bh = snap.sin_send_op[r]; - layer.for_each_cross_rank_sin_send_range(r, 0, bc, [&bs, &bh, &state, &op](size_t k, size_t i) { - bs[k] = state[i]; - bh[k] = op[i]; - }); - } - const size_t dc = layer.cross_rank_sin_recv_size(r); - snap.sin_recv_state[r].resize(dc); - snap.sin_recv_op[r].resize(dc); - if (dc > 0) { - auto &ds = snap.sin_recv_state[r]; - auto &dh = snap.sin_recv_op[r]; - layer.for_each_cross_rank_sin_recv_range(r, - 0, - dc, - [&ds, &dh, &state, &op](size_t k, size_t i, int /*phi*/) { - ds[k] = state[i]; - dh[k] = op[i]; - }); - } - } + // Grow-only: shrinking would free the allocations this scratch exists to reuse. + const size_t occupied = layer.occupied_slot_count(); + snap.sin_send_state.resize(std::max(snap.sin_send_state.size(), occupied)); + snap.sin_send_op.resize(std::max(snap.sin_send_op.size(), occupied)); + snap.sin_recv_state.resize(std::max(snap.sin_recv_state.size(), occupied)); + snap.sin_recv_op.resize(std::max(snap.sin_recv_op.size(), occupied)); + layer.for_each_occupied_slot( + [&snap, my_rank, &state, &op](size_t pos, size_t r, const detail::CrossRankSlotView &slot) { + if (r == my_rank) { + clear_slot_snapshot(snap, pos); + return; + } + fill_slot_snapshot(snap, pos, state, op, slot); + }); } } // namespace @@ -449,7 +427,7 @@ auto state_operator_derivative_local(VecD &state, const size_t R = layer.cross_rank_rank_count(); auto &snap = derivative_snapshot_scratch(); - snapshot_remote_endpoints(state, op, layer, my_rank, R, snap); + snapshot_remote_endpoints(state, op, layer, my_rank, snap); // No-op at single rank; the transfer touches only buffers, so the cos pass below may mutate state/op. auto in_flight = begin_cross_rank_derivative_exchange(snap, layer, comm); @@ -459,7 +437,7 @@ auto state_operator_derivative_local(VecD &state, EndpointContrib ep; if (my_rank < R) { - ep = apply_self_slot_derivative_paired(state, op, layer, my_rank, trig); + ep = apply_self_slot_derivative_paired(state, op, layer, trig); } const auto remote = finish_cross_rank_derivative_exchange(state, op, layer, snap, trig, in_flight); ep = combine_endpoint_contrib(ep, remote); @@ -478,19 +456,15 @@ auto evolve_step_traversal_impl(VecD &op, const double sin_val = std::sin(2 * param); auto *const op_data = op.data(); - const int my_rank_int = mpi::rank(comm); - const auto my_rank = static_cast(my_rank_int); - // Snapshot my_rank's own sin_send values before the cos pass; runs unconditionally (the remote pack - // skips my_rank) so single-rank works. - const size_t self_b_count = (my_rank < layer.cross_rank_rank_count()) ? layer.cross_rank_sin_send_size(my_rank) : 0; + // This rank's own sin_send values before the cos pass. Unconditional, since the remote pack skips + // the self slot, so single-rank works. + const auto self_slot = layer.cross_rank_self_slot(); + const size_t self_b_count = self_slot.sin_send_count; VecD self_b_snapshot; self_b_snapshot.resize(self_b_count); - if (self_b_count > 0) { - auto &snap = self_b_snapshot; - layer.for_each_cross_rank_sin_send_range(my_rank, 0, self_b_count, [&snap, &op](size_t k, size_t i) { - snap[k] = op[i]; - }); + for (size_t k = 0; k < self_b_count; ++k) { + self_b_snapshot[k] = op[detail::slot_sin_send_index(self_slot, k)]; } // Pack + start the exchange before the cos scan so partner values are pre-cos and the transfer overlaps. @@ -499,15 +473,9 @@ auto evolve_step_traversal_impl(VecD &op, finish_cross_rank_evolution_exchange(op, layer, sin_val, in_flight); // Self-slot sin_recv entries: op[i] is already cos-scaled, so only the sine term is added. - if (self_b_count > 0) { - const size_t self_d_count = layer.cross_rank_sin_recv_size(my_rank); - layer.for_each_cross_rank_sin_recv_range(my_rank, - 0, - self_d_count, - [&op, &sin_val, &self_b_snapshot](size_t k, size_t i, int phi_signed) { - op[i] += - sin_val * static_cast(phi_signed) * self_b_snapshot[k]; - }); + for (size_t k = 0; k < self_b_count; ++k) { + const size_t i = detail::slot_sin_recv_index(self_slot, k); + op[i] += sin_val * static_cast(detail::slot_sin_recv_phase(self_slot, k)) * self_b_snapshot[k]; } } diff --git a/cpp/monoprop/MPGraph.cpp b/cpp/monoprop/MPGraph.cpp index a425582d..f9391a60 100644 --- a/cpp/monoprop/MPGraph.cpp +++ b/cpp/monoprop/MPGraph.cpp @@ -51,7 +51,11 @@ auto layer_storage_memory_usage(const LayerCore &storage) -> GraphMemoryBreakdow GraphMemoryBreakdown breakdown; breakdown.layer_storage_object_bytes = sizeof(LayerCore); breakdown.cross_rank_bytes = detail::cross_rank_storage_bytes(storage.cross_rank); - breakdown.exchange_layout_bytes = detail::layer_exchange_layout_storage_bytes(storage.evolution_exchange_layout); + breakdown.slot_record_bytes = detail::cross_rank_slot_record_bytes(storage.cross_rank); + breakdown.layer_cores = 1; + breakdown.slot_records = storage.cross_rank.rank_count(); + breakdown.occupied_slots = detail::cross_rank_occupied_slots(storage.cross_rank); + breakdown.cross_rank_endpoints = detail::cross_rank_endpoint_count(storage.cross_rank); return breakdown; } diff --git a/cpp/monoprop/detail/evolution/layer_build/Engine.h b/cpp/monoprop/detail/evolution/layer_build/Engine.h index f19bddf7..1c831995 100644 --- a/cpp/monoprop/detail/evolution/layer_build/Engine.h +++ b/cpp/monoprop/detail/evolution/layer_build/Engine.h @@ -175,7 +175,7 @@ struct GraphSink { append_inserted_endpoints(cos_all, combined_size, op); *out_cos = std::move(cos_all); } - return build_layer_storage_unified(std::move(partners), my_rank); + return build_layer_storage_unified(partners, my_rank); } }; diff --git a/cpp/monoprop/detail/graph/MPGraphLayers.h b/cpp/monoprop/detail/graph/MPGraphLayers.h index 715aedda..9699cea0 100644 --- a/cpp/monoprop/detail/graph/MPGraphLayers.h +++ b/cpp/monoprop/detail/graph/MPGraphLayers.h @@ -28,8 +28,7 @@ namespace monoprop { // recompute (nullopt) — cosine rebuilt from the generator's inverted-index columns at replay. // pruned (has value) — cosine pre-filtered to a backward-reachable subset, stored explicitly; an // empty stored list is still pruned (replay as nothing, do not recompute). -// Cores are shared and immutable in value only: their eval-time caches (recv_cache, the lazy derivative -// layout) are filled through const handles, so evaluating two aliasing propagators concurrently is a race. +// Cores are shared and immutable: they hold no eval-time cache. // Cross-rank data is always read verbatim; only the cosine set is ever filtered. struct LayerTraversal final { @@ -56,46 +55,48 @@ struct LayerTraversal final { auto cross_rank_sin_recv_size(size_t rank) const -> size_t { return core_->cross_rank.sin_recv_size(rank); } auto cross_rank_in_count(size_t rank) const -> size_t { return core_->cross_rank.in_count(rank); } - // Random access into the D list, for the paired self-slot derivative fetches d[k], d[k+P]. - auto cross_rank_sin_recv_index_at(size_t rank, size_t idx) const -> size_t { - return detail::cross_rank_sin_recv_index(core_->cross_rank, rank, idx); + auto cross_rank_self_slot() const -> detail::CrossRankSlotView { + return detail::cross_rank_self_slot(core_->cross_rank); } - auto cross_rank_sin_recv_phase_at(size_t rank, size_t idx) const -> int { - return detail::cross_rank_sin_recv_phase(core_->cross_rank, rank, idx); + + // Every slot carrying traffic, ascending, each with its offset. Use this for a sweep. + template + auto for_each_occupied_slot(Func &&func) const -> void { + detail::for_each_occupied_slot(core_->cross_rank, std::forward(func)); } + // The size an array indexed by occupied position needs. + auto occupied_slot_count() const -> size_t { return detail::cross_rank_occupied_slots(core_->cross_rank); } + + // Resolved once outside the loop: per endpoint would make per-slot work per-term. template auto for_each_cross_rank_sin_send_range(size_t rank, size_t begin, size_t end, Func &&func) const -> void { + const auto slot = detail::cross_rank_slot(core_->cross_rank, rank); for (size_t idx = begin; idx < end; ++idx) { - func(idx, detail::cross_rank_sin_send_index(core_->cross_rank, rank, idx)); + func(idx, detail::slot_sin_send_index(slot, idx)); } } template auto for_each_cross_rank_sin_recv_range(size_t rank, size_t begin, size_t end, Func &&func) const -> void { + const auto slot = detail::cross_rank_slot(core_->cross_rank, rank); for (size_t idx = begin; idx < end; ++idx) { - func(idx, - detail::cross_rank_sin_recv_index(core_->cross_rank, rank, idx), - detail::cross_rank_sin_recv_phase(core_->cross_rank, rank, idx)); + func(idx, detail::slot_sin_recv_index(slot, idx), detail::slot_sin_recv_phase(slot, idx)); } } - auto evolution_exchange_layout() const -> const LayerExchangeLayout & { return core_->evolution_exchange_layout; } - auto derivative_exchange_layout() const -> const LayerExchangeLayout & { - return core_->derivative_exchange_layout(); - } + // Both sides of the exchange layout derive from these. + auto cross_rank() const -> const PackedCrossRankStorage & { return core_->cross_rank; } auto param_index() const -> size_t { return core_->param_index; } auto gen_coeff() const -> double { return core_->gen_coeff; } auto gate_index() const -> size_t { return core_->gate_index; } // Rotations (Givens cycles) = sum of per-rank in-counts (one in-entry per rotation). sin_recv_size - // would double-count self-rank rotations (in+out). + // would double-count self-rank rotations (in+out); empty slots contribute nothing to either. auto total_cycles() const -> size_t { size_t count = 0; - for (size_t rank = 0; rank < cross_rank_rank_count(); ++rank) { - count += cross_rank_in_count(rank); - } + for_each_occupied_slot([&count](size_t, const detail::CrossRankSlotView &slot) { count += slot.in_count; }); return count; } @@ -103,9 +104,8 @@ struct LayerTraversal final { // indices = num_cos_inds() - total_rotation_endpoints(). auto total_rotation_endpoints() const -> size_t { size_t count = 0; - for (size_t rank = 0; rank < cross_rank_rank_count(); ++rank) { - count += cross_rank_sin_recv_size(rank); - } + for_each_occupied_slot( + [&count](size_t, const detail::CrossRankSlotView &slot) { count += slot.sin_send_count; }); return count; } diff --git a/cpp/monoprop/detail/graph/MPGraphViews.h b/cpp/monoprop/detail/graph/MPGraphViews.h index 0af7d7a2..80df2672 100644 --- a/cpp/monoprop/detail/graph/MPGraphViews.h +++ b/cpp/monoprop/detail/graph/MPGraphViews.h @@ -40,6 +40,13 @@ struct GraphMemoryBreakdown final { size_t cross_rank_bytes = 0; size_t exchange_layout_bytes = 0; + // Diagnostics, outside total_bytes(): counts, or a subset of a byte field, not additions to it. + size_t slot_record_bytes = 0; + size_t layer_cores = 0; + size_t slot_records = 0; // slot_records / layer_cores is the flat world P + size_t occupied_slots = 0; + size_t cross_rank_endpoints = 0; + auto total_bytes() const -> size_t { return layer_descriptor_bytes + layer_storage_object_bytes + cos_data_bytes + cross_rank_bytes + exchange_layout_bytes; @@ -52,6 +59,11 @@ struct GraphMemoryBreakdown final { cos_data_bytes += o.cos_data_bytes; cross_rank_bytes += o.cross_rank_bytes; exchange_layout_bytes += o.exchange_layout_bytes; + slot_record_bytes += o.slot_record_bytes; + layer_cores += o.layer_cores; + slot_records += o.slot_records; + occupied_slots += o.occupied_slots; + cross_rank_endpoints += o.cross_rank_endpoints; return *this; } }; diff --git a/cpp/monoprop/detail/graph_encoding/MPGraphEncoding.cpp b/cpp/monoprop/detail/graph_encoding/MPGraphEncoding.cpp index 2e8da60f..b7ccc437 100644 --- a/cpp/monoprop/detail/graph_encoding/MPGraphEncoding.cpp +++ b/cpp/monoprop/detail/graph_encoding/MPGraphEncoding.cpp @@ -14,9 +14,11 @@ #include "monoprop/detail/graph_encoding/MPGraphEncodingStorage.h" +#include #include #include #include +#include #include #include #include @@ -32,32 +34,14 @@ auto checked_mpi_int(size_t value, const char *what) -> int { return static_cast(value); } -auto build_layer_exchange_layout(const std::vector &send_counts, int scale, const char *what) - -> LayerExchangeLayout { - const std::string count_label = std::format("{} count", what); - const std::string displacement_label = std::format("{} displacement", what); - - LayerExchangeLayout layout; - layout.counts.resize(send_counts.size()); - layout.displs.resize(send_counts.size()); - size_t total = 0; - for (size_t r = 0; r < send_counts.size(); ++r) { - const size_t count = static_cast(scale) * send_counts[r]; - layout.counts[r] = checked_mpi_int(count, count_label.c_str()); - layout.displs[r] = checked_mpi_int(total, displacement_label.c_str()); - total += count; - } - layout.total_count = total; - return layout; -} - -auto build_derivative_exchange_layout(const LayerExchangeLayout &evolution) -> LayerExchangeLayout { - std::vector send_counts; - send_counts.reserve(evolution.counts.size()); - for (const int count : evolution.counts) { - send_counts.push_back(static_cast(count)); +// The u32 slot id bounds the flat world; a narrowing conversion, so checked rather than cast. +auto checked_world_slot(size_t rank) -> uint32_t { + if (rank > static_cast(std::numeric_limits::max())) { + throw std::overflow_error(std::format("World slot {} exceeds the {} the occupied-slot record can hold.", + rank, + std::numeric_limits::max())); } - return build_layer_exchange_layout(send_counts, 2, "Layer derivative exchange"); + return static_cast(rank); } auto checked_term_index(size_t value, const char *what) -> TermIndex { @@ -100,21 +84,44 @@ auto packed_phase_storage_bytes(const PackedPhaseStorage &storage) -> size_t { auto build_packed_cross_rank_storage(const std::vector &data) -> PackedCrossRankStorage { PackedCrossRankStorage storage; const size_t num_ranks = data.size(); - storage.ranges.resize(num_ranks); + storage.world_size = num_ranks; size_t total_b = 0; - size_t total_d = 0; for (size_t rank = 0; rank < num_ranks; ++rank) { const auto &partner = data[rank]; - auto &range = storage.ranges[rank]; - range.sin_send_offset = total_b; - range.sin_send_count = static_cast(partner.sin_send_indices.size()); - range.sin_recv_offset = total_d; - range.sin_recv_count = static_cast(partner.sin_recv_entries.size()); - range.in_count = static_cast(partner.in_count); + // One count and one offset serve both B and D: a skew mis-derives Q and reads the wrong + // endpoint rather than throwing. + if (partner.sin_send_indices.size() != partner.sin_recv_entries.size()) { + throw CrossRankSlotLayoutError(std::format( + "Cross-rank slot {} has {} send endpoints against {} recv endpoints; B and D are the same set.", + rank, + partner.sin_send_indices.size(), + partner.sin_recv_entries.size())); + } + // The in-block bounds a block within B, and Q is an unsigned subtraction, so a boundary past + // the end wraps near 2^64 instead of going negative and every D read runs off the array. + if (partner.in_count > partner.sin_send_indices.size()) { + throw CrossRankSlotLayoutError( + std::format("Cross-rank slot {} declares an in-block of {} inside {} endpoints; the in-block is a " + "boundary within the endpoint list, not an addition to it.", + rank, + partner.in_count, + partner.sin_send_indices.size())); + } + // A slot with no traffic gets no record; ascending rank order leaves `occupied` sorted. + if (partner.sin_send_indices.empty()) { + continue; + } + // Checked, not cast: readers rebuild offsets from the stored counts, so a truncated slot + // shifts every later slot's window. + storage.occupied.push_back( + {.slot = checked_world_slot(rank), + .sin_send_count = checked_term_index(partner.sin_send_indices.size(), "Cross-rank slot endpoint count"), + .in_count = checked_term_index(partner.in_count, "Cross-rank slot in-block size")}); total_b += partner.sin_send_indices.size(); - total_d += partner.sin_recv_entries.size(); } + storage.occupied.shrink_to_fit(); // push_back overshoots, and this array is the thing being shrunk + const size_t total_d = total_b; bool uses_binary_phases = true; for (const auto &partner : data) { @@ -128,10 +135,12 @@ auto build_packed_cross_rank_storage(const std::vector &da storage.sin_send_indices.resize(total_b); storage.sin_recv_phases = make_packed_phase_storage(total_d, uses_binary_phases); - for (size_t rank = 0; rank < num_ranks; ++rank) { - const auto &partner = data[rank]; - const size_t b_off = storage.ranges[rank].sin_send_offset; - const size_t d_off = storage.ranges[rank].sin_recv_offset; + // The order the offsets accumulate in, so readers reconstruct this exact prefix. + size_t offset = 0; + for (const auto &entry : storage.occupied) { + const auto &partner = data[entry.slot]; + const size_t b_off = offset; + const size_t d_off = b_off; // equal counts, so equal prefix sums for (size_t k = 0; k < partner.sin_send_indices.size(); ++k) { storage.sin_send_indices[b_off + k] = checked_term_index(partner.sin_send_indices[k], "Cross-rank B index"); @@ -142,61 +151,102 @@ auto build_packed_cross_rank_storage(const std::vector &da (void)i; store_packed_phase(storage.sin_recv_phases, d_off + k, phi, "Cross-rank D phase"); } + offset += entry.sin_send_count; } return storage; } +auto resolve_self_slot(PackedCrossRankStorage &storage, size_t my_rank) -> void { + storage.self_pos = kNoSelfSlot; + storage.self_offset = 0; + size_t offset = 0; + for (size_t pos = 0; pos < storage.occupied.size(); ++pos) { + const auto &entry = storage.occupied[pos]; + if (entry.slot == my_rank) { + storage.self_pos = pos; + storage.self_offset = offset; + return; + } + offset += entry.sin_send_count; + } +} + auto cross_rank_storage_bytes(const PackedCrossRankStorage &storage) -> size_t { - size_t bytes = - storage.ranges.capacity() * sizeof(CrossRankPartnerRange) + packed_phase_storage_bytes(storage.sin_recv_phases); + size_t bytes = cross_rank_slot_record_bytes(storage) + packed_phase_storage_bytes(storage.sin_recv_phases); bytes += storage.sin_send_indices.capacity() * sizeof(TermIndex); return bytes; } -auto layer_exchange_layout_storage_bytes(const LayerExchangeLayout &layout) -> size_t { - return layout.counts.capacity() * sizeof(int) + layout.displs.capacity() * sizeof(int); +auto cross_rank_slot_record_bytes(const PackedCrossRankStorage &storage) -> size_t { + return storage.occupied.capacity() * sizeof(CrossRankOccupiedSlot); } -auto build_layer_storage_unified(std::vector all_partners, size_t my_rank) - -> std::shared_ptr { - auto storage = std::make_shared(); - - { - std::vector send_counts; - send_counts.reserve(all_partners.size()); - for (size_t r = 0; r < all_partners.size(); ++r) { - send_counts.push_back((r == my_rank) ? size_t{0} : all_partners[r].sin_send_indices.size()); - } - storage->evolution_exchange_layout = build_layer_exchange_layout(send_counts, 1); +auto cross_rank_occupied_slots(const PackedCrossRankStorage &storage) -> size_t { + return storage.occupied.size(); +} - // The derivative layout (2x) is allocated lazily on first gradient read, but validated here: an - // overflow must throw during build_graph, not from inside the gradient collective window, where - // peers are already blocked in mpi::resolve_recv's count round -> a distributed hang, not an error. - static_cast(build_derivative_exchange_layout(storage->evolution_exchange_layout)); +auto cross_rank_endpoint_count(const PackedCrossRankStorage &storage) -> size_t { + size_t count = 0; + for (const auto &entry : storage.occupied) { + count += entry.sin_send_count; } + return count; +} - storage->cross_rank = build_packed_cross_rank_storage(std::move(all_partners)); +namespace { +// Formats the label only on the throwing path: this runs per posted exchange, not once at build. +auto checked_exchange_int(size_t value, const char *what, const char *field) -> int { + if (value > static_cast(std::numeric_limits::max())) { + return checked_mpi_int(value, std::format("{} {}", what, field).c_str()); + } + return static_cast(value); +} +} // namespace + +auto derive_exchange_layout(const PackedCrossRankStorage &cross_rank, + size_t my_rank, + int scale, + LayerExchangeLayout &out, + const char *what) -> void { + const size_t num_ranks = cross_rank.rank_count(); + // assign(), not resize(): `out` is reused across layers and an empty slot must read zero. + out.counts.assign(num_ranks, 0); + out.displs.resize(num_ranks); + + // Scatter over the occupied slots: a dense probe would binary-search every possible partner to + // fill an array that is mostly zeros. + for_each_occupied_slot(cross_rank, [my_rank, &out, scale, what](size_t slot, const CrossRankSlotView &view) { + if (slot == my_rank) { + return; // excluded from the transfer and handled locally, as the stored layout did + } + out.counts[slot] = checked_exchange_int(static_cast(scale) * view.sin_send_count, what, "count"); + }); - // Both are indexed by the same rank space. - if (storage->evolution_exchange_layout.counts.size() != storage->cross_rank.rank_count()) { - throw ExchangeLayoutRankMismatch( - std::format("Layer exchange layout covers {} ranks but cross-rank storage has {}.", - storage->evolution_exchange_layout.counts.size(), - storage->cross_rank.rank_count())); + // The prefix sum stays dense: MPI_Alltoallv wants a valid displacement for every rank. + size_t total = 0; + for (size_t r = 0; r < num_ranks; ++r) { + out.displs[r] = checked_exchange_int(total, what, "displacement"); + total += static_cast(out.counts[r]); } - return storage; + out.total_count = total; } -} // namespace monoprop::detail +auto build_layer_storage_unified(const std::vector &all_partners, size_t my_rank) + -> std::shared_ptr { + auto storage = std::make_shared(); -namespace monoprop { + storage->cross_rank = build_packed_cross_rank_storage(all_partners); + resolve_self_slot(storage->cross_rank, my_rank); -auto LayerCore::derivative_exchange_layout() const -> const LayerExchangeLayout & { - if (!derivative_exchange_layout_cache_) { - derivative_exchange_layout_cache_ = detail::build_derivative_exchange_layout(evolution_exchange_layout); + { + // Result discarded: an int overflow must throw here, not inside a committed exchange. Scale 2 + // bounds scale 1. + LayerExchangeLayout scratch; + derive_exchange_layout(storage->cross_rank, my_rank, 2, scratch, "Layer derivative exchange"); } - return *derivative_exchange_layout_cache_; + + return storage; } -} // namespace monoprop +} // namespace monoprop::detail diff --git a/cpp/monoprop/detail/graph_encoding/MPGraphEncodingStorage.h b/cpp/monoprop/detail/graph_encoding/MPGraphEncodingStorage.h index c1ed79c9..976b8e03 100644 --- a/cpp/monoprop/detail/graph_encoding/MPGraphEncodingStorage.h +++ b/cpp/monoprop/detail/graph_encoding/MPGraphEncodingStorage.h @@ -14,21 +14,18 @@ #pragma once +#include #include #include #include #include #include +#include #include #include "monoprop/detail/graph_encoding/MPGraphEncodingTypes.h" namespace monoprop::detail { -// The layer exchange layout and the packed cross-rank storage disagree on the rank count. -class ExchangeLayoutRankMismatch : public std::logic_error { -public: - using std::logic_error::logic_error; -}; // The ceiling has to track the TermIndex width, not a fixed 32-bit limit. auto checked_term_index(size_t value, const char *what) -> TermIndex; @@ -75,33 +72,128 @@ inline auto store_packed_phase(PackedPhaseStorage &storage, size_t idx, int phas } } +// Invariants of how the layer was built, so a violation is a bug here rather than bad input. +class CrossRankSlotLayoutError : public std::logic_error { +public: + using std::logic_error::logic_error; +}; + auto build_packed_cross_rank_storage(const std::vector &data) -> PackedCrossRankStorage; -inline auto cross_rank_sin_send_index(const PackedCrossRankStorage &storage, size_t rank, size_t idx) -> size_t { - const size_t offset = storage.ranges[rank].sin_send_offset + idx; - return static_cast(storage.sin_send_indices[offset]); +// Record where my_rank's own slot sits, so the gradient's self-slot reads are O(1). Once per layer. +auto resolve_self_slot(PackedCrossRankStorage &storage, size_t my_rank) -> void; + +// One world slot's position in the flat B/D arrays. The accessors below take this rather than a slot +// id, so walking a slot's endpoints pays the lookup once instead of per endpoint. +struct CrossRankSlotView final { + const TermIndex *sin_send_indices = nullptr; // B, already advanced to this slot's offset + const PackedPhaseStorage *sin_recv_phases = nullptr; + size_t phase_offset = 0; + size_t sin_send_count = 0; + size_t in_count = 0; +}; + +namespace slot_detail { +inline auto view_at(const PackedCrossRankStorage &storage, const CrossRankOccupiedSlot &entry, size_t offset) + -> CrossRankSlotView { + return CrossRankSlotView{.sin_send_indices = storage.sin_send_indices.data() + offset, + .sin_recv_phases = &storage.sin_recv_phases, + .phase_offset = offset, + .sin_send_count = entry.sin_send_count, + .in_count = entry.in_count}; +} +} // namespace slot_detail + +// Every occupied slot in ascending order with its B/D offset. func(slot_id, view), or +// func(occupied_pos, slot_id, view) -- the position is handed out because a counter at the call site +// would skew past those loops' self-slot `return`s. O(occupied) for the whole sweep. +template +auto for_each_occupied_slot(const PackedCrossRankStorage &storage, Func &&func) -> void { + size_t offset = 0; + for (size_t pos = 0; pos < storage.occupied.size(); ++pos) { + const auto &entry = storage.occupied[pos]; + if constexpr (std::is_invocable_v) { + func(pos, static_cast(entry.slot), slot_detail::view_at(storage, entry, offset)); + } + else { + func(static_cast(entry.slot), slot_detail::view_at(storage, entry, offset)); + } + offset += entry.sin_send_count; + } +} + +// O(1), for the innermost gradient loop; an all-zero view when this rank's slot carries no traffic. +inline auto cross_rank_self_slot(const PackedCrossRankStorage &storage) -> CrossRankSlotView { + if (storage.self_pos == kNoSelfSlot) { + return CrossRankSlotView{.sin_recv_phases = &storage.sin_recv_phases}; + } + return slot_detail::view_at(storage, storage.occupied[storage.self_pos], storage.self_offset); } -// Invariant B=[in(P)]++[out(Q)], D=[out(Q)]++[in(P)] (P=in_count, Q=sin_recv_count-P): +// Arbitrary slot, O(occupied) because the offset is a prefix. Diagnostic and test use. +inline auto cross_rank_slot(const PackedCrossRankStorage &storage, size_t rank) -> CrossRankSlotView { + size_t offset = 0; + for (const auto &entry : storage.occupied) { + if (entry.slot == rank) { + return slot_detail::view_at(storage, entry, offset); + } + if (entry.slot > rank) { + break; + } + offset += entry.sin_send_count; + } + return CrossRankSlotView{.sin_recv_phases = &storage.sin_recv_phases}; +} + +inline auto slot_sin_send_index(const CrossRankSlotView &slot, size_t idx) -> size_t { + return static_cast(slot.sin_send_indices[idx]); +} + +// Invariant B=[in(P)]++[out(Q)], D=[out(Q)]++[in(P)] (P=in_count, Q=sin_send_count-P): // D[idx] = (idx size_t { + // Unsigned subtraction, kept safe by the precondition build_packed_cross_rank_storage enforces. + assert(slot.in_count <= slot.sin_send_count && "in-block cannot exceed the slot's endpoint count"); + const size_t out_count = slot.sin_send_count - slot.in_count; // Q + const size_t sin_send_local = (idx < out_count) ? (slot.in_count + idx) : (idx - out_count); + return slot_sin_send_index(slot, sin_send_local); +} + +inline auto slot_sin_recv_phase(const CrossRankSlotView &slot, size_t idx) -> int { + return packed_phase_at(*slot.sin_recv_phases, slot.phase_offset + idx); +} + +inline auto cross_rank_sin_send_index(const PackedCrossRankStorage &storage, size_t rank, size_t idx) -> size_t { + return slot_sin_send_index(cross_rank_slot(storage, rank), idx); +} + inline auto cross_rank_sin_recv_index(const PackedCrossRankStorage &storage, size_t rank, size_t idx) -> size_t { - const auto &range = storage.ranges[rank]; - const size_t in_count = range.in_count; // P - const size_t out_count = range.sin_recv_count - in_count; // Q - const size_t sin_send_local = (idx < out_count) ? (in_count + idx) : (idx - out_count); - return cross_rank_sin_send_index(storage, rank, sin_send_local); + return slot_sin_recv_index(cross_rank_slot(storage, rank), idx); } inline auto cross_rank_sin_recv_phase(const PackedCrossRankStorage &storage, size_t rank, size_t idx) -> int { - return packed_phase_at(storage.sin_recv_phases, storage.ranges[rank].sin_recv_offset + idx); + return slot_sin_recv_phase(cross_rank_slot(storage, rank), idx); } auto cross_rank_storage_bytes(const PackedCrossRankStorage &storage) -> size_t; -auto layer_exchange_layout_storage_bytes(const LayerExchangeLayout &layout) -> size_t; +// The slot-proportional slice of cross_rank_storage_bytes: one record per stored world slot. +auto cross_rank_slot_record_bytes(const PackedCrossRankStorage &storage) -> size_t; + +auto cross_rank_occupied_slots(const PackedCrossRankStorage &storage) -> size_t; + +// The B array's length, and an upper bound on the occupied slot count. +auto cross_rank_endpoint_count(const PackedCrossRankStorage &storage) -> size_t; + +// Derive a layer's send layout into caller-owned scratch; `out` is resized, not reallocated, on reuse. +auto derive_exchange_layout(const PackedCrossRankStorage &cross_rank, + size_t my_rank, + int scale, + LayerExchangeLayout &out, + const char *what = "Layer exchange") -> void; // Local cycles fold into the self-rank slot (my_rank); the exchange layout zeroes counts[my_rank] so // MPI_Alltoallv skips it (replay does a local copy). -auto build_layer_storage_unified(std::vector all_partners, size_t my_rank) +auto build_layer_storage_unified(const std::vector &all_partners, size_t my_rank) -> std::shared_ptr; } // namespace monoprop::detail diff --git a/cpp/monoprop/detail/graph_encoding/MPGraphEncodingTypes.h b/cpp/monoprop/detail/graph_encoding/MPGraphEncodingTypes.h index 20be7ace..9f91de6a 100644 --- a/cpp/monoprop/detail/graph_encoding/MPGraphEncodingTypes.h +++ b/cpp/monoprop/detail/graph_encoding/MPGraphEncodingTypes.h @@ -14,26 +14,26 @@ #pragma once +#include #include #include #include #include +#include #include #include #include #include "monoprop/TypeAliases.h" -#include "monoprop/detail/mpi/RecvLayout.h" namespace monoprop { +// counts/displs for one alltoallv, dense int[P] as MPI requires. Materialised into per-thread scratch +// for the exchange being posted, never retained per layer, and it describes the recv side too. struct LayerExchangeLayout final { std::vector counts; std::vector displs; size_t total_count = 0; - - // Cached recv counts/displs (see mpi::resolve_recv); mutable — filled through const handles at eval time. - mutable mpi::RecvLayoutCache recv_cache; }; } // namespace monoprop @@ -42,15 +42,6 @@ namespace monoprop::detail { auto checked_mpi_int(size_t value, const char *what) -> int; -// Per-rank MPI counts = send_counts[r] * scale, with prefix-sum displacements. send_counts is full-width -// (size_t) so checked_mpi_int catches the narrowing to MPI's int. -auto build_layer_exchange_layout(const std::vector &send_counts, int scale, const char *what = "Layer exchange") - -> LayerExchangeLayout; - -// The derivative layout is the evolution layout at 2x (each rotation endpoint carries both the op and -// state payload). -auto build_derivative_exchange_layout(const LayerExchangeLayout &evolution) -> LayerExchangeLayout; - } // namespace monoprop::detail namespace monoprop { @@ -116,38 +107,63 @@ struct CrossRankPartnerData { // Default-init: every element is overwritten before any read. DefaultInitVector sin_send_indices; DefaultInitVector> sin_recv_entries; - // Size of the in-block (P). + // In-block size, a boundary within sin_send_indices. size_t in_count = 0; }; -struct CrossRankPartnerRange final { - size_t sin_send_offset = 0; // into sin_send_indices; cumulative across ranks, so size_t (may exceed 2^32) - TermIndex sin_send_count = - 0; // == sin_recv_count (both endpoints); TermIndex-wide so one rank/layer can exceed 2^32 - size_t sin_recv_offset = 0; // into sin_recv_phases; cumulative across ranks, so size_t (see sin_send_offset) - TermIndex sin_recv_count = 0; +// One world slot carrying traffic. Offsets come from a running prefix, which keeps the record at 12 B. +struct CrossRankOccupiedSlot final { + uint32_t slot = 0; // flat world slot id -- what the dense array encoded by position + TermIndex sin_send_count = 0; TermIndex in_count = 0; }; +// A width-independent rule, so it is checked on both TermIndex widths. +inline constexpr size_t kOccupiedSlotIdField = std::max(sizeof(uint32_t), alignof(TermIndex)); +static_assert(sizeof(CrossRankOccupiedSlot) == kOccupiedSlotIdField + 2 * sizeof(TermIndex), + "the occupied-slot record is the graph's P coefficient: it must carry no padding beyond the " + "alignment slot its u32 id already occupies."); +static_assert(alignof(CrossRankOccupiedSlot) == alignof(TermIndex), + "the record aligns to its widest member; if that stops holding the size rule above is " + "measuring something else."); + +// self_pos sentinel: this rank's own slot carries no traffic in this layer. +inline constexpr size_t kNoSelfSlot = static_cast(-1); struct PackedCrossRankStorage final { - std::vector ranges; // size == R + std::vector occupied; std::vector sin_send_indices; PackedPhaseStorage sin_recv_phases; // one phased entry per D index, sign baked in - auto rank_count() const -> size_t { return ranges.size(); } - auto sin_send_size(size_t rank) const -> size_t { return ranges[rank].sin_send_count; } - auto sin_recv_size(size_t rank) const -> size_t { return ranges[rank].sin_recv_count; } - auto in_count(size_t rank) const -> size_t { return ranges[rank].in_count; } + // P. No longer recoverable from the array length, and MPI_Alltoallv still wants dense counts. + size_t world_size = 0; + + // Read in the innermost gradient loop, so resolved once at build. + size_t self_pos = kNoSelfSlot; + size_t self_offset = 0; + + auto rank_count() const -> size_t { return world_size; } + + // O(log occupied); a sweep should use for_each_occupied_slot. + auto find(size_t rank) const -> const CrossRankOccupiedSlot * { + const auto it = std::ranges::lower_bound(occupied, rank, {}, [](const CrossRankOccupiedSlot &e) { + return static_cast(e.slot); + }); + return (it == occupied.end() || it->slot != rank) ? nullptr : std::to_address(it); + } + + auto sin_send_size(size_t rank) const -> size_t { + const auto *e = find(rank); + return e == nullptr ? 0 : e->sin_send_count; + } + auto sin_recv_size(size_t rank) const -> size_t { return sin_send_size(rank); } + auto in_count(size_t rank) const -> size_t { + const auto *e = find(rank); + return e == nullptr ? 0 : e->in_count; + } }; struct LayerCore final { PackedCrossRankStorage cross_rank; - LayerExchangeLayout evolution_exchange_layout; - - auto derivative_exchange_layout() const -> const LayerExchangeLayout &; - - // A copied core must not inherit the source's cache: it is eval-time state, not data. - auto reset_derivative_exchange_layout() -> void { derivative_exchange_layout_cache_.reset(); } // Per-layer recompute metadata: generator_words = this layer's generator G as backing words; // scaled_count = fold truncation bound = operator size after this layer's partner inserts. @@ -159,9 +175,6 @@ struct LayerCore final { double gen_coeff = 0.0; // Shared by all layers of one multi-term gate; absolute across build_graph calls (parameter_mapping). size_t gate_index = 0; - -private: - mutable std::optional derivative_exchange_layout_cache_; }; } // namespace monoprop diff --git a/cpp/monoprop/detail/monomial_propagator/MonomialPropagator.inl b/cpp/monoprop/detail/monomial_propagator/MonomialPropagator.inl index 117590ec..6e73befd 100644 --- a/cpp/monoprop/detail/monomial_propagator/MonomialPropagator.inl +++ b/cpp/monoprop/detail/monomial_propagator/MonomialPropagator.inl @@ -389,31 +389,23 @@ auto MonomialPropagator::graph_data() const -> std::vector // Always empty: local cycles are folded into cross_rank[my_rank]. std::vector local_cyc_data; - std::vector b_data, d_data; - b_data.reserve(rank_count); - d_data.reserve(rank_count); - for (size_t rank = 0; rank < rank_count; ++rank) { - VecZ sin_send_indices(traversal.cross_rank_sin_send_size(rank)); - VecI b_phases(traversal.cross_rank_sin_send_size(rank), 0); - VecZ d_indices(traversal.cross_rank_sin_recv_size(rank)); - VecI sin_recv_phases(traversal.cross_rank_sin_recv_size(rank)); - - traversal.for_each_cross_rank_sin_send_range( - rank, - 0, - traversal.cross_rank_sin_send_size(rank), - [&](size_t logical_idx, size_t value_idx) { sin_send_indices[logical_idx] = value_idx; }); - traversal.for_each_cross_rank_sin_recv_range(rank, - 0, - traversal.cross_rank_sin_recv_size(rank), - [&](size_t logical_idx, size_t value_idx, int phase) { - d_indices[logical_idx] = value_idx; - sin_recv_phases[logical_idx] = phase; - }); - - b_data.emplace_back(std::move(sin_send_indices), std::move(b_phases)); - d_data.emplace_back(std::move(d_indices), std::move(sin_recv_phases)); - } + // The exported shape stays dense (callers index by rank), but it is filled by scattering the + // occupied slots rather than by asking every possible slot how much it holds. + std::vector b_data(rank_count), d_data(rank_count); + traversal.for_each_occupied_slot([&](size_t rank, const detail::CrossRankSlotView &slot) { + const size_t count = slot.sin_send_count; + VecZ sin_send_indices(count); + VecI b_phases(count, 0); + VecZ d_indices(count); + VecI sin_recv_phases(count); + for (size_t k = 0; k < count; ++k) { + sin_send_indices[k] = detail::slot_sin_send_index(slot, k); + d_indices[k] = detail::slot_sin_recv_index(slot, k); + sin_recv_phases[k] = detail::slot_sin_recv_phase(slot, k); + } + b_data[rank] = CrossRankData{std::move(sin_send_indices), std::move(b_phases)}; + d_data[rank] = CrossRankData{std::move(d_indices), std::move(sin_recv_phases)}; + }); // Same two-way read as cos_index_count_(): a pared layer's stored set is authoritative, and // recomputing the fold over it would report the indices the pare removed. VecZ cos_inds; @@ -819,8 +811,6 @@ auto MonomialPropagator::set_parameter_mapping(const VecZ ¶meter_m auto relabel = [this](size_t layer, size_t new_param_index) { auto &target = graph_.get_layer(layer); auto new_core = std::make_shared(target.core()); - // Drop the inherited eval-time derivative layout: it must not depend on a prior gradient run. - new_core->reset_derivative_exchange_layout(); new_core->param_index = new_param_index; if (const CosMask *pruned = target.pruned_cos()) { target = Layer(std::move(new_core), *pruned); diff --git a/cpp/monoprop/detail/mpi/CMakeLists.txt b/cpp/monoprop/detail/mpi/CMakeLists.txt index c8226a41..5b9969b0 100644 --- a/cpp/monoprop/detail/mpi/CMakeLists.txt +++ b/cpp/monoprop/detail/mpi/CMakeLists.txt @@ -12,7 +12,6 @@ target_sources( "MPICompat.h" "MPIUtils.h" "PartitionBarrier.h" - "RecvLayout.h" "ShmComm.h" ) diff --git a/cpp/monoprop/detail/mpi/Exchange.h b/cpp/monoprop/detail/mpi/Exchange.h index 2bd0f7ae..9ec32530 100644 --- a/cpp/monoprop/detail/mpi/Exchange.h +++ b/cpp/monoprop/detail/mpi/Exchange.h @@ -19,16 +19,14 @@ #include #include "monoprop/detail/mpi/MPICompat.h" -#include "monoprop/detail/mpi/RecvLayout.h" // Keeps #ifdef monoprop_ENABLE_MPI out of the consumers; non-MPI builds get self-copy stubs. namespace monoprop::mpi { -// Resolve the recv side of a send-count vector, reusing `cache` when comm size is unchanged: a -// replayed graph's send pattern is fixed, so a hit removes one blocking count round-trip per layer -// per evaluation. -auto resolve_recv(std::span send_counts, const Comm &comm, RecvLayoutCache &cache) -> const RecvLayout &; +// Precondition, not a diagnostic: MPI_Alltoallv reads one count and one displacement per rank +// whatever the span holds, so a layout built for a differently sized communicator reads out of bounds. +auto check_exchange_layout_width(std::span send_counts, const Comm &comm) -> void; // Idempotent completion handle for a posted payload transfer; move-only, so a request is waited on // exactly once. wait() is a no-op on the blocking path and in non-MPI builds. Owns its request: the @@ -97,7 +95,7 @@ inline auto post_flat_alltoallv(const FlatAlltoallvArgs &args, int num_ranks, return Ticket{}; } (void)num_ranks; - auto request = MPI_REQUEST_NULL; + MPI_Request request = MPI_REQUEST_NULL; MPI_Ialltoallv(args.send, args.send_counts, args.send_displs, diff --git a/cpp/monoprop/detail/mpi/HybridComm.h b/cpp/monoprop/detail/mpi/HybridComm.h index 62efb62e..47a79880 100644 --- a/cpp/monoprop/detail/mpi/HybridComm.h +++ b/cpp/monoprop/detail/mpi/HybridComm.h @@ -16,6 +16,7 @@ #include #include +#include #include #include #include @@ -65,14 +66,31 @@ class HybridComm { } // Size all (R,S)-fixed scratch once so per-call paths never allocate; staging grows on demand. const size_t rss = static_cast(r_) * static_cast(s_) * static_cast(s_); + const size_t p = static_cast(r_) * static_cast(s_); counts_send_.resize(rss); counts_recv_.resize(rss); mpi_send_counts_.resize(static_cast(r_)); mpi_recv_counts_.resize(static_cast(r_)); mpi_send_displs_.resize(static_cast(r_)); mpi_recv_displs_.resize(static_cast(r_)); - pack_off_.resize(rss); - scatter_off_.resize(rss); + pack_off_.resize(static_cast(s_) * p); // same S*R*S elements as before, indexed (u, g) + base_send_.resize(p); + base_recv_.resize(p); + col_sum_.resize(p); + recv_col_.resize(p); + // Rounded up to whole 64-byte lines: each row has its own writer, so a shared line false-shares. + counts_stride_ = round_up_(p, kIntsPerLine); + rows_stride_ = round_up_(static_cast(r_), kLongsPerLine); + // One spare line each: the allocator guarantees alignof(T), not 64, so row 0 must be realigned. + counts_matrix_store_.assign(static_cast(s_) * counts_stride_ + kIntsPerLine, 0); + counts_matrix_ = align_to_line_(counts_matrix_store_.data()); + rows_store_.assign(static_cast(s_) * rows_stride_ + kLongsPerLine, 0LL); + rows_ = align_to_line_(rows_store_.data()); + assert(reinterpret_cast(counts_matrix_) % kLineBytes == 0); + assert(reinterpret_cast(rows_) % kLineBytes == 0); + assert(counts_matrix_ + static_cast(s_) * counts_stride_ + <= counts_matrix_store_.data() + counts_matrix_store_.size()); + assert(rows_ + static_cast(s_) * rows_stride_ <= rows_store_.data() + rows_store_.size()); } HybridComm(const HybridComm &) = delete; @@ -122,12 +140,10 @@ class HybridComm { auto reset() -> void { barrier_.reset(); } private: + // Counts travel through counts_matrix_ / rows_, not through here. struct alignas(64) Slot { const std::byte *ptr = nullptr; // byte view of this partition's send buffer; null until published - const int *counts = nullptr; - const int *send_counts = nullptr; const int *send_displs = nullptr; - const int *recv_counts = nullptr; const double *vec = nullptr; double f64 = 0.0; uint64_t u64 = 0; @@ -156,11 +172,10 @@ class HybridComm { // recv_counts[g] = amount global partition g sends to this partition. 2 barriers + one S*S-int MPI_Alltoall. auto alltoall_counts_impl_(int local_partition, const int *send_counts /*[P]*/, int *recv_counts /*[P]*/) -> void { - const size_t u = static_cast(local_partition); - slots_[u].counts = send_counts; + publish_counts_row_(local_partition, send_counts); sync(); if (local_partition == 0) { - pack_count_matrix_(&Slot::counts); + pack_count_matrix_(); MPI_Alltoall(counts_send_.data(), s_ * s_, MPI_INT, counts_recv_.data(), s_ * s_, MPI_INT, parent_); } sync(); @@ -168,14 +183,12 @@ class HybridComm { const int t = local_partition; for (int a = 0; a < r_; ++a) { for (int su = 0; su < s_; ++su) { - const size_t idx = (static_cast(a) * static_cast(s_) * static_cast(s_)) - + (static_cast(t) * static_cast(s_)) + static_cast(su); - recv_counts[a * s_ + su] = counts_recv_[idx]; + recv_counts[a * s_ + su] = counts_recv_[counts_idx_(a, t, su)]; } } // No trailing barrier: past the last sync only counts_recv_ is read, and partition 0 cannot rewrite // it until a future call's second barrier — unreachable until every extractor here has arrived. - // send_counts was consumed before the last sync, so peers may free/reuse it on return. + // send_counts is consumed before the first sync, so peers may free or reuse it on return. } // Flat variable all-to-all over caller-owned buffers; see AlltoallvArgs for the conventions. @@ -183,14 +196,16 @@ class HybridComm { const size_t u = static_cast(local_partition); Slot &me = slots_[u]; me.ptr = args.send; - me.send_counts = args.send_counts; me.send_displs = args.send_displs; - me.recv_counts = args.recv_counts; + publish_counts_row_(local_partition, args.send_counts); + publish_recv_rows_(local_partition, args.recv_counts); sync(); // B1 // B2: partition 0 sizes/reallocates staging; must finish before any partition packs into stage_send_. if (local_partition == 0) { - size_staging_(args.elem); + size_staging_send_(args.elem); + fill_recv_col_from_rows_(); + size_staging_recv_(args.elem); } sync(); // B2 @@ -213,21 +228,24 @@ class HybridComm { sync(); // B4 // Scatter each global source's contiguous run from stage_recv_ to recv_displs[g] (all legs, incl. - // self-rank, go through staging). Block starts come from scatter_off_: no peer slot is read past B4. + // self-rank, go through staging). Walks (a, su) in accumulation order, so `cur` re-derives the + // block starts from base_recv_ and this partition's own counts. std::byte *dst = args.recv; const int t = local_partition; for (int a = 0; a < r_; ++a) { + size_t cur = base_recv_[static_cast(a) * static_cast(s_) + static_cast(t)]; for (int su = 0; su < s_; ++su) { const int g = a * s_ + su; const int cnt = args.recv_counts[g]; if (cnt != 0) { std::memcpy(dst + static_cast(args.recv_displs[g]) * args.elem, - stage_recv_.data() + scatter_off_[block_idx_(a, t, su)] * args.elem, + stage_recv_.data() + cur * args.elem, static_cast(cnt) * args.elem); } + cur += static_cast(cnt); } } - // No trailing barrier (see alltoall_counts_impl_): past B4 only stage_recv_/scatter_off_/own buffers. + // No trailing barrier: base_recv_ is rewritten only in a later verb's B1→B2 window. } // Fused count-resolve + payload alltoallv: folds the standalone count exchange into this verb's @@ -242,19 +260,17 @@ class HybridComm { // Typed here but byte-addressed in the slot: pack_send_ copies by (displ, count) in elements and // never reconstructs T, so the slot stays type-erased for the untyped alltoallv_impl_ above. me.ptr = reinterpret_cast(args.send); - me.send_counts = args.send_counts; me.send_displs = args.send_displs; - // recv_counts is an output here — deliberately not published; the count Alltoall resolves it. + // Count row only: the recv counts do not exist until the count Alltoall in B1→B2. + publish_counts_row_(local_partition, args.send_counts); sync(); // B1 if (local_partition == 0) { - pack_count_matrix_(&Slot::send_counts); + pack_count_matrix_(); MPI_Alltoall(counts_send_.data(), s_ * s_, MPI_INT, counts_recv_.data(), s_ * s_, MPI_INT, parent_); - // Size staging from counts_recv_: recv of partition t from (rank, su) sits at rank*S*S + t*S + su. - size_staging_impl_(elem, [this](int t, int rank, int su) { - return counts_recv_[(static_cast(rank) * static_cast(s_) * static_cast(s_)) - + (static_cast(t) * static_cast(s_)) + static_cast(su)]; - }); + size_staging_send_(elem); + fill_recv_col_from_counts_recv_(); + size_staging_recv_(elem); } sync(); // B2 @@ -263,9 +279,7 @@ class HybridComm { for (int a = 0; a < r_; ++a) { for (int su = 0; su < s_; ++su) { const int g = a * s_ + su; - const int c = - counts_recv_[(static_cast(a) * static_cast(s_) * static_cast(s_)) - + (static_cast(t) * static_cast(s_)) + static_cast(su)]; + const int c = counts_recv_[counts_idx_(a, t, su)]; args.recv_counts[g] = c; args.recv_displs[g] = checked_mpi_count(total, "Recv displacement"); total += c; @@ -291,14 +305,16 @@ class HybridComm { std::byte *dst = reinterpret_cast(args.recv.data()); // after the resize: it may reallocate for (int a = 0; a < r_; ++a) { + size_t cur = base_recv_[static_cast(a) * static_cast(s_) + static_cast(t)]; for (int su = 0; su < s_; ++su) { const int g = a * s_ + su; const int cnt = args.recv_counts[g]; if (cnt != 0) { std::memcpy(dst + static_cast(args.recv_displs[g]) * elem, - stage_recv_.data() + scatter_off_[block_idx_(a, t, su)] * elem, + stage_recv_.data() + cur * elem, static_cast(cnt) * elem); } + cur += static_cast(cnt); } } // No trailing barrier: same discipline as alltoallv_impl_. @@ -372,25 +388,64 @@ class HybridComm { // No trailing barrier: red_vec_ is rewritten only inside a future verb's barriered phases. } - // Flat index of the (rank, dest partition, source partition) block in the R*S*S offset/count tables. - auto block_idx_(int b, int t, int u) const -> size_t { + static constexpr size_t kLineBytes = 64; + static constexpr size_t kIntsPerLine = kLineBytes / sizeof(int); + static constexpr size_t kLongsPerLine = kLineBytes / sizeof(long long); + + static auto round_up_(size_t n, size_t m) -> size_t { return ((n + m - 1) / m) * m; } + + template + static auto align_to_line_(T *p) -> T * { + static_assert(kLineBytes % sizeof(T) == 0); + const auto addr = reinterpret_cast(p); + const size_t pad = (kLineBytes - static_cast(addr % kLineBytes)) % kLineBytes; + return p + pad / sizeof(T); + } + + // (rank, dest partition, source partition) in the R*S*S count matrices: the count message's wire layout. + auto counts_idx_(int b, int t, int su) const -> size_t { return (static_cast(b) * static_cast(s_) + static_cast(t)) * static_cast(s_) - + static_cast(u); + + static_cast(su); } - // Pack the S*S count matrix per dest rank into counts_send_, dest-partition-major (t) then - // source-partition-minor (su), ready for the one S*S-int MPI_Alltoall. `counts` selects which slot - // field to read, the only difference between the standalone count exchange (Slot::counts) and the - // fused resolve (Slot::send_counts). - // - // Partition 0 only, and only inside a barriered window: it reads every peer partition's published - // count pointer. Every element of counts_send_ is written here and MPI_Alltoall fills counts_recv_ - // fully, so neither buffer is pre-zeroed. - auto pack_count_matrix_(const int *Slot::*counts) -> void { - for (int b = 0; b < r_; ++b) { - for (int t = 0; t < s_; ++t) { - for (int su = 0; su < s_; ++su) { - counts_send_[block_idx_(b, t, su)] = (slots_[static_cast(su)].*counts)[b * s_ + t]; + // [S x P] payload offset table: row u is source partition u's staging starts, one per destination g. + auto pack_idx_(int u, int g) const -> size_t { + return static_cast(u) * static_cast(r_) * static_cast(s_) + static_cast(g); + } + + auto counts_row_(int u) -> int * { return counts_matrix_ + static_cast(u) * counts_stride_; } + auto row_recv_(int u) -> long long * { return rows_ + static_cast(u) * rows_stride_; } + + // Phase P0: partition u writes only its own row, so the pre-barrier write needs no barrier. The one + // cross-partition reader is partition 0 inside B1→B2, and no peer reaches verb k+1's P0 without + // passing verb k's B2, so the rewrite always follows the read. + auto publish_counts_row_(int local_partition, const int *send_counts) -> void { + std::memcpy(counts_row_(local_partition), + send_counts, + static_cast(r_) * static_cast(s_) * sizeof(int)); + } + + // long long: it sums S int counts. + auto publish_recv_rows_(int local_partition, const int *recv_counts) -> void { + long long *rr = row_recv_(local_partition); + for (int a = 0; a < r_; ++a) { + long long sum = 0; + for (int su = 0; su < s_; ++su) { + sum += recv_counts[a * s_ + su]; + } + rr[a] = sum; + } + } + + // Transpose the published count rows into counts_send_, dest-major then source-minor, for the one + // S*S-int MPI_Alltoall. Partition 0 only, inside a barriered window. Source partition outer, so the + // peer-owned side streams. Every element is written here, so no pre-zeroing. + auto pack_count_matrix_() -> void { + for (int su = 0; su < s_; ++su) { + const int *row = counts_row_(su); + for (int b = 0; b < r_; ++b) { + for (int t = 0; t < s_; ++t) { + counts_send_[counts_idx_(b, t, su)] = row[b * s_ + t]; } } } @@ -403,80 +458,118 @@ class HybridComm { } } - // Partition 0 only; recv counts come from partition t's published recv_counts. - auto size_staging_(size_t elem) -> void { - size_staging_impl_(elem, [this](int t, int rank, int su) { - return slots_[static_cast(t)].recv_counts[rank * s_ + su]; - }); - } - - // recv_count(t, rank, su) yields the count partition t on this rank receives from (rank, source partition su). - template - auto size_staging_impl_(size_t elem, RecvCountFn recv_count) -> void { + // Partition 0's send-side staging sizing, between B1 and B2, in two sweeps of counts_matrix_. Wire + // block order is destination major, source minor: a per-destination base plus a per-source prefix. + auto size_staging_send_(size_t elem) -> void { + const size_t p = static_cast(r_) * static_cast(s_); + // Pass A: the column sums W over source partitions, u outer so both sides sweep in address order. + std::fill(col_sum_.begin(), col_sum_.end(), 0LL); + for (int u = 0; u < s_; ++u) { + const int *row = counts_row_(u); + for (size_t g = 0; g < p; ++g) { + col_sum_[g] += row[g]; + } + } for (int b = 0; b < r_; ++b) { + const long long *col = col_sum_.data() + static_cast(b) * static_cast(s_); long long send_sum = 0; - long long recv_sum = 0; for (int t = 0; t < s_; ++t) { - for (int su = 0; su < s_; ++su) { - send_sum += slots_[static_cast(su)].send_counts[b * s_ + t]; - recv_sum += recv_count(t, b, su); - } + send_sum += col[t]; } mpi_send_counts_[static_cast(b)] = checked_mpi_count(send_sum, "Per-rank send count"); - mpi_recv_counts_[static_cast(b)] = checked_mpi_count(recv_sum, "Per-rank recv count"); } long long send_running = 0; - long long recv_running = 0; for (int b = 0; b < r_; ++b) { mpi_send_displs_[static_cast(b)] = checked_mpi_count(send_running, "Send displacement"); - mpi_recv_displs_[static_cast(b)] = checked_mpi_count(recv_running, "Recv displacement"); send_running += mpi_send_counts_[static_cast(b)]; - recv_running += mpi_recv_counts_[static_cast(b)]; } const size_t total_send = static_cast(checked_mpi_count(send_running, "Total send count")); - const size_t total_recv = static_cast(checked_mpi_count(recv_running, "Total recv count")); - // Grow-only, no zero-fill: pack_send_'s blocks tile [0, total_send) exactly and MPI_Alltoallv - // fills every live byte of stage_recv_, so stale bytes past a prior high-water mark are never read. - grow_(stage_send_, total_send * elem); - grow_(stage_recv_, total_recv * elem); - // Precompute each (rank b, dest partition t, source partition u) block start in elements, dest-major/ - // source-minor to match staging: pack/scatter become O(1) lookups that read no peer slot past B4 - // (re-summing peer count matrices instead would be O(R*S^3)). for (int b = 0; b < r_; ++b) { size_t cur = static_cast(mpi_send_displs_[static_cast(b)]); for (int t = 0; t < s_; ++t) { - for (int u = 0; u < s_; ++u) { - pack_off_[block_idx_(b, t, u)] = cur; - cur += static_cast(slots_[static_cast(u)].send_counts[b * s_ + t]); - } + const size_t g = static_cast(b) * static_cast(s_) + static_cast(t); + base_send_[g] = cur; + cur += static_cast(col_sum_[g]); } } + // Pass B: the exclusive prefix over source partitions; col_sum_ is free to be reused for it here. + std::fill(col_sum_.begin(), col_sum_.end(), 0LL); + for (int u = 0; u < s_; ++u) { + const int *row = counts_row_(u); + size_t *off = pack_off_.data() + pack_idx_(u, 0); + for (size_t g = 0; g < p; ++g) { + off[g] = base_send_[g] + static_cast(col_sum_[g]); + col_sum_[g] += row[g]; + } + } + // Grow-only, no zero-fill: pack_send_'s blocks tile [0, total_send) exactly. + grow_(stage_send_, total_send * elem); + } + + // recv_col_[a*S + t] = what partition t receives from rank a: the rows published in Phase P0. + auto fill_recv_col_from_rows_() -> void { for (int a = 0; a < r_; ++a) { - size_t cur = static_cast(mpi_recv_displs_[static_cast(a)]); for (int t = 0; t < s_; ++t) { + recv_col_[static_cast(a) * static_cast(s_) + static_cast(t)] = row_recv_(t)[a]; + } + } + } + + auto fill_recv_col_from_counts_recv_() -> void { + for (int a = 0; a < r_; ++a) { + for (int t = 0; t < s_; ++t) { + const int *blk = counts_recv_.data() + counts_idx_(a, t, 0); + long long sum = 0; for (int su = 0; su < s_; ++su) { - scatter_off_[block_idx_(a, t, su)] = cur; - cur += static_cast(recv_count(t, a, su)); + sum += blk[su]; } + recv_col_[static_cast(a) * static_cast(s_) + static_cast(t)] = sum; } } } + // Partition 0's recv-side staging sizing, from recv_col_. Only the per-(rank, partition) base; the + // post-B4 scatter re-derives the per-source offsets as it walks (a, su). + auto size_staging_recv_(size_t elem) -> void { + for (int a = 0; a < r_; ++a) { + long long recv_sum = 0; + for (int t = 0; t < s_; ++t) { + recv_sum += recv_col_[static_cast(a) * static_cast(s_) + static_cast(t)]; + } + mpi_recv_counts_[static_cast(a)] = checked_mpi_count(recv_sum, "Per-rank recv count"); + } + long long recv_running = 0; + for (int a = 0; a < r_; ++a) { + mpi_recv_displs_[static_cast(a)] = checked_mpi_count(recv_running, "Recv displacement"); + recv_running += mpi_recv_counts_[static_cast(a)]; + } + const size_t total_recv = static_cast(checked_mpi_count(recv_running, "Total recv count")); + for (int a = 0; a < r_; ++a) { + size_t cur = static_cast(mpi_recv_displs_[static_cast(a)]); + for (int t = 0; t < s_; ++t) { + const size_t g = static_cast(a) * static_cast(s_) + static_cast(t); + base_recv_[g] = cur; + cur += static_cast(recv_col_[g]); + } + } + grow_(stage_recv_, total_recv * elem); + } + auto pack_send_(int local_partition, size_t elem) -> void { const int u = local_partition; // Own slot only — no peer's published send buffer is read here, which is what lets every // partition pack concurrently in the B2→B3 window. const std::byte *src = slots_[static_cast(u)].ptr; - const int *my_send_counts = slots_[static_cast(u)].send_counts; + const int *my_send_counts = counts_row_(u); const int *my_send_displs = slots_[static_cast(u)].send_displs; - for (int b = 0; b < r_; ++b) { - for (int t = 0; t < s_; ++t) { - const int cnt = my_send_counts[b * s_ + t]; - if (cnt != 0) { - std::memcpy(stage_send_.data() + pack_off_[block_idx_(b, t, u)] * elem, - src + static_cast(my_send_displs[b * s_ + t]) * elem, - static_cast(cnt) * elem); - } + const size_t *off = pack_off_.data() + pack_idx_(u, 0); + const int p = r_ * s_; + for (int g = 0; g < p; ++g) { + const int cnt = my_send_counts[g]; + if (cnt != 0) { + std::memcpy(stage_send_.data() + off[g] * elem, + src + static_cast(my_send_displs[g]) * elem, + static_cast(cnt) * elem); } } } @@ -501,7 +594,7 @@ class HybridComm { int mpi_rank_ = 0; std::vector slots_; - // Partition-0-managed shared state (written by partition 0, read by all between barriers). + // Partition-0-managed, except the two Phase P0 tables at the end, whose rows each partition owns. // S*S per rank, the counts alltoall. std::vector counts_send_; std::vector counts_recv_; @@ -510,9 +603,24 @@ class HybridComm { std::vector mpi_send_displs_; std::vector mpi_recv_counts_; std::vector mpi_recv_displs_; - // [R*S*S] block starts (elements) in the staging buffers. + // [S x P] payload block starts in stage_send_, row u per global destination g. No scatter-side + // table: see size_staging_recv_. std::vector pack_off_; - std::vector scatter_off_; + // [P] staging bases: base_send_[b*S+t] starts partition (b,t)'s region of the message to rank b, + // base_recv_[a*S+t] starts local partition t's region of the message from rank a. + std::vector base_send_; + std::vector base_recv_; + std::vector col_sum_; + std::vector recv_col_; + // [S x counts_stride_] send-count matrix, row u owned by partition u; padded so no two rows share a + // 64-byte line. Lifetime rule: publish_counts_row_. + std::vector counts_matrix_store_; + int *counts_matrix_ = nullptr; + size_t counts_stride_ = 0; + // [S x rows_stride_] per-partition per-rank recv totals, same ownership and padding rules. + std::vector rows_store_; + long long *rows_ = nullptr; + size_t rows_stride_ = 0; // Aggregated MPI payload staging, HWM-sized. std::vector stage_send_; std::vector stage_recv_; diff --git a/cpp/monoprop/detail/mpi/MPICompat.cpp b/cpp/monoprop/detail/mpi/MPICompat.cpp index 30141d4b..133a9fb7 100644 --- a/cpp/monoprop/detail/mpi/MPICompat.cpp +++ b/cpp/monoprop/detail/mpi/MPICompat.cpp @@ -119,12 +119,10 @@ auto alltoall_counts(const int *send_counts, int *recv_counts, int n, Comm comm) #endif } -auto resolve_recv(std::span send_counts, const Comm &comm, RecvLayoutCache &cache) -> const RecvLayout & { +auto check_exchange_layout_width(std::span send_counts, const Comm &comm) -> void { const auto n = static_cast(send_counts.size()); const int comm_size = mpi::size(comm); - // alltoall_counts moves comm_size ints each way regardless of `n`, so a send vector that is not - // exactly one entry per rank reads and writes out of bounds — reachable because layouts outlive - // propagator copies and pare rebuilds. + // MPI_Alltoallv reads comm_size counts and displacements whatever the span holds. if (n != comm_size) { throw CollectiveArgumentError( std::format("Exchange layout has {} send counts but the communicator has {} ranks — a graph built for one " @@ -132,22 +130,6 @@ auto resolve_recv(std::span send_counts, const Comm &comm, RecvLayout n, comm_size)); } - if (cache.comm_size == comm_size && static_cast(cache.layout.counts.size()) == n) { - return cache.layout; - } - - RecvLayout &out = cache.layout; - out.counts.resize(static_cast(n)); - alltoall_counts(send_counts.data(), out.counts.data(), n, comm); - out.displs.resize(static_cast(n)); - long long total = 0; - for (int i = 0; i < n; ++i) { - out.displs[static_cast(i)] = checked_mpi_count(total); - total += out.counts[static_cast(i)]; - } - out.total = checked_mpi_count(total); - cache.comm_size = comm_size; - return out; } } // namespace monoprop::mpi diff --git a/cpp/monoprop/detail/mpi/RecvLayout.h b/cpp/monoprop/detail/mpi/RecvLayout.h deleted file mode 100644 index f48f9d4a..00000000 --- a/cpp/monoprop/detail/mpi/RecvLayout.h +++ /dev/null @@ -1,35 +0,0 @@ -// Copyright 2026 Algorithmiq -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -#pragma once - -#include - -// Kept MPI-free and dependency-light so graph-encoding types (LayerExchangeLayout) can embed the cache -// without pulling in or the exchange machinery (see Exchange.h). - -namespace monoprop::mpi { - -struct RecvLayout { - std::vector counts; - std::vector displs; - int total = 0; -}; - -struct RecvLayoutCache { - RecvLayout layout; - int comm_size = -1; -}; - -} // namespace monoprop::mpi diff --git a/cpp/monoprop/detail/pare/PareGraph.cpp b/cpp/monoprop/detail/pare/PareGraph.cpp index 4d6dccec..b490a0a9 100644 --- a/cpp/monoprop/detail/pare/PareGraph.cpp +++ b/cpp/monoprop/detail/pare/PareGraph.cpp @@ -26,15 +26,6 @@ namespace monoprop { namespace { -template -auto for_each_remote_rank(const LayerTraversal &layer, size_t my_rank, Func &&func) -> void { - for (size_t rank = 0; rank < layer.cross_rank_rank_count(); ++rank) { - if (rank != my_rank) { - func(rank); - } - } -} - // Bits of `present` whose absolute index is in the keep set. inline auto keep_mask_for_block(const std::vector &keep, size_t base, uint64_t present) -> uint64_t { uint64_t mask = 0; @@ -84,17 +75,13 @@ auto mark_cross_rank_endpoints_kept(const LayerTraversal &layer, size_t my_rank, nodes_to_keep[idx] = 1; } }; - for (size_t rank = 0; rank < layer.cross_rank_rank_count(); ++rank) { - layer.for_each_cross_rank_sin_recv_range(rank, - 0, - layer.cross_rank_sin_recv_size(rank), - [&mark](size_t, size_t tgt_idx, int) { mark(tgt_idx); }); - } - for_each_remote_rank(layer, my_rank, [&layer, &mark](size_t rank) { - layer.for_each_cross_rank_sin_send_range(rank, - 0, - layer.cross_rank_sin_send_size(rank), - [&mark](size_t, size_t src_idx) { mark(src_idx); }); + layer.for_each_occupied_slot([&mark, my_rank](size_t rank, const detail::CrossRankSlotView &slot) { + for (size_t k = 0; k < slot.sin_send_count; ++k) { + mark(detail::slot_sin_recv_index(slot, k)); + if (rank != my_rank) { + mark(detail::slot_sin_send_index(slot, k)); + } + } }); } diff --git a/cpp/tests/ExchangeLayoutOracle.h b/cpp/tests/ExchangeLayoutOracle.h new file mode 100644 index 00000000..b34cf078 --- /dev/null +++ b/cpp/tests/ExchangeLayoutOracle.h @@ -0,0 +1,50 @@ +// Copyright 2026 Algorithmiq +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#pragma once + +// The independent reference for a layer's exchange layout, which derive_exchange_layout must match +// elementwise: counts[r] = scale * send_counts[r], displacements their prefix sum. +// Outside the library on purpose: an oracle refactored beside its subject asserts nothing. + +#include +#include +#include +#include + +#include "monoprop/detail/graph_encoding/MPGraphEncodingTypes.h" + +namespace exchange_layout_oracle { + +inline auto build_layer_exchange_layout(const std::vector &send_counts, + int scale, + const char *what = "Layer exchange") -> monoprop::LayerExchangeLayout { + const std::string count_label = std::format("{} count", what); + const std::string displacement_label = std::format("{} displacement", what); + + monoprop::LayerExchangeLayout layout; + layout.counts.resize(send_counts.size()); + layout.displs.resize(send_counts.size()); + size_t total = 0; + for (size_t r = 0; r < send_counts.size(); ++r) { + const size_t count = static_cast(scale) * send_counts[r]; + layout.counts[r] = monoprop::detail::checked_mpi_int(count, count_label.c_str()); + layout.displs[r] = monoprop::detail::checked_mpi_int(total, displacement_label.c_str()); + total += count; + } + layout.total_count = total; + return layout; +} + +} // namespace exchange_layout_oracle diff --git a/cpp/tests/README.md b/cpp/tests/README.md index f167fe1f..6fc34490 100644 --- a/cpp/tests/README.md +++ b/cpp/tests/README.md @@ -70,6 +70,9 @@ name and cannot address suite-nested cases, tests use flat - **`GraphBuildHarness.h`**: direct Layer/MPGraph construction helpers (`core_with_gate`, `layer_with_gate`, `graph_with_gates`) for white-box MPGraph transform tests. +- **`ExchangeLayoutOracle.h`**: `build_layer_exchange_layout` — the independent + reference for a layer's exchange counts and displacements, which + `derive_exchange_layout` in the library is checked against. - **`TestData.{h,cpp}`**: the `CaseData` struct and msgpack fixture loader. - **`boost-test.cmake` / `boostAddTests.cmake`**: CMake test discovery. @@ -106,6 +109,8 @@ name and cannot address suite-nested cases, tests use flat - **Simulator / operator lifecycle**: `simulator_copy_tests.cpp`, `update_initial_operator.cpp`, `ctor_validation_tests.cpp` (constructor guard rails + MPGraph bounds). +- **Exchange layout preconditions**: `exchange_layout_precondition_tests.cpp` + (the layout width check, driven through ShmComm). New `*.cpp` files are auto-discovered on the next configure — no CMake edit needed. diff --git a/cpp/tests/exchange_layout_precondition_tests.cpp b/cpp/tests/exchange_layout_precondition_tests.cpp new file mode 100644 index 00000000..1015db8f --- /dev/null +++ b/cpp/tests/exchange_layout_precondition_tests.cpp @@ -0,0 +1,52 @@ +// Copyright 2026 Algorithmiq +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +// Preconditions on a layer's exchange layout. Driven through ShmComm, so it needs no ranks. + +#include + +#include + +#include "monoprop/detail/mpi/Comm.h" +#include "monoprop/detail/mpi/Exchange.h" +#include "monoprop/detail/mpi/MPICompat.h" +#include "monoprop/detail/mpi/ShmComm.h" + +using monoprop::mpi::CollectiveArgumentError; +using monoprop::mpi::Comm; +using monoprop::mpi::ShmComm; + +BOOST_AUTO_TEST_CASE(exchange_layout_width_is_checked) { + // Sized 4 but driven by one thread: the throw must precede any collective, or this hangs. + ShmComm world(4); + const Comm comm = Comm::make_shm(&world, /*rank=*/0); + BOOST_REQUIRE_EQUAL(monoprop::mpi::size(comm), 4); + + const std::vector too_short(3, 0); + BOOST_CHECK_THROW(monoprop::mpi::check_exchange_layout_width(too_short, comm), CollectiveArgumentError); + + const std::vector too_long(5, 0); + BOOST_CHECK_THROW(monoprop::mpi::check_exchange_layout_width(too_long, comm), CollectiveArgumentError); + + const std::vector empty; + BOOST_CHECK_THROW(monoprop::mpi::check_exchange_layout_width(empty, comm), CollectiveArgumentError); +} + +BOOST_AUTO_TEST_CASE(exchange_layout_of_the_right_width_is_accepted) { + ShmComm world(1); + const Comm comm = Comm::make_shm(&world, /*rank=*/0); + + const std::vector exact(1, 0); + BOOST_CHECK_NO_THROW(monoprop::mpi::check_exchange_layout_width(exact, comm)); +} diff --git a/cpp/tests/graph_encoding_tests.cpp b/cpp/tests/graph_encoding_tests.cpp index 26bfe2cf..6ff1c1a0 100644 --- a/cpp/tests/graph_encoding_tests.cpp +++ b/cpp/tests/graph_encoding_tests.cpp @@ -19,12 +19,16 @@ #include #include +#include #include +#include "ExchangeLayoutOracle.h" #include "monoprop/detail/graph_encoding/MPGraphEncodingStorage.h" using namespace monoprop; +using exchange_layout_oracle::build_layer_exchange_layout; + BOOST_AUTO_TEST_CASE(graph_encoding_word_builder_push_index_coalesces_within_word) { CosineWordBuilder b; b.push_index(0); @@ -119,47 +123,93 @@ BOOST_AUTO_TEST_CASE(graph_encoding_packed_phase_at_reads_int8_values) { BOOST_AUTO_TEST_CASE(graph_encoding_exchange_layout_scale_and_displacements) { const std::vector send_counts = {3, 0, 5}; - const auto s1 = detail::build_layer_exchange_layout(send_counts, /*scale=*/1); + const auto s1 = build_layer_exchange_layout(send_counts, /*scale=*/1); BOOST_CHECK((s1.counts == std::vector{3, 0, 5})); BOOST_CHECK((s1.displs == std::vector{0, 3, 3})); // prefix sum: 0, 0+3, 3+0 BOOST_CHECK_EQUAL(s1.total_count, 8U); - const auto s2 = detail::build_layer_exchange_layout(send_counts, /*scale=*/2); + const auto s2 = build_layer_exchange_layout(send_counts, /*scale=*/2); BOOST_CHECK((s2.counts == std::vector{6, 0, 10})); BOOST_CHECK((s2.displs == std::vector{0, 6, 6})); BOOST_CHECK_EQUAL(s2.total_count, 16U); +} + +// The derived layout must equal the ExchangeLayoutOracle.h reference elementwise, not literals. +namespace { - BOOST_CHECK_GT(detail::layer_exchange_layout_storage_bytes(s1), 0U); +auto slot_partners(const std::vector &sin_send_counts) -> std::vector { + std::vector data(sin_send_counts.size()); + for (size_t r = 0; r < sin_send_counts.size(); ++r) { + for (size_t k = 0; k < sin_send_counts[r]; ++k) { + data[r].sin_send_indices.push_back(k); + data[r].sin_recv_entries.push_back({k, 1}); + } + } + return data; } -// Production only builds scale=1; the 2x layout reaches MPI through this accessor, which is -// unreachable at comm size 1, so the default non-MPI suite would otherwise never touch it. +} // namespace + +BOOST_AUTO_TEST_CASE(graph_encoding_derived_layout_matches_the_layout_it_replaces) { + const std::vector counts{3, 0, 5, 2}; + const auto storage = detail::build_packed_cross_rank_storage(slot_partners(counts)); + + for (size_t my_rank = 0; my_rank < counts.size(); ++my_rank) { + std::vector expected_counts = counts; + expected_counts[my_rank] = 0; -BOOST_AUTO_TEST_CASE(graph_encoding_derivative_exchange_layout_is_twice_the_evolution_layout) { - LayerCore core; - core.evolution_exchange_layout = detail::build_layer_exchange_layout({3, 0, 5}, /*scale=*/1); + for (const int scale : {1, 2}) { + const auto reference = build_layer_exchange_layout(expected_counts, scale); + LayerExchangeLayout derived; + detail::derive_exchange_layout(storage, my_rank, scale, derived); - const auto &derivative = core.derivative_exchange_layout(); - BOOST_CHECK((derivative.counts == std::vector{6, 0, 10})); - BOOST_CHECK((derivative.displs == std::vector{0, 6, 6})); - BOOST_CHECK_EQUAL(derivative.total_count, 16U); + BOOST_CHECK(derived.counts == reference.counts); + BOOST_CHECK(derived.displs == reference.displs); + BOOST_CHECK_EQUAL(derived.total_count, reference.total_count); + } + } +} + +BOOST_AUTO_TEST_CASE(graph_encoding_derived_layout_reuses_its_scratch) { + // Reused across layers: a stale tail reads as a real count for an unused slot. + const auto wide = detail::build_packed_cross_rank_storage(slot_partners({1, 2, 3, 4})); + const auto narrow = detail::build_packed_cross_rank_storage(slot_partners({7, 7})); - // Cached: the second read returns the same object, so eval-time MPI holds a stable pointer. - BOOST_CHECK_EQUAL(&core.derivative_exchange_layout(), &derivative); + LayerExchangeLayout scratch; + detail::derive_exchange_layout(wide, /*my_rank=*/0, 1, scratch); + BOOST_CHECK_EQUAL(scratch.counts.size(), 4U); + detail::derive_exchange_layout(narrow, /*my_rank=*/0, 1, scratch); + BOOST_CHECK_EQUAL(scratch.counts.size(), 2U); + BOOST_CHECK((scratch.counts == std::vector{0, 7})); + BOOST_CHECK_EQUAL(scratch.total_count, 7U); +} + +BOOST_AUTO_TEST_CASE(graph_encoding_a_zero_traffic_slot_still_gets_a_valid_displacement) { + // Count 0, but the displacement must still advance. + const auto storage = detail::build_packed_cross_rank_storage(slot_partners({0, 4, 0, 0, 6})); + LayerExchangeLayout derived; + detail::derive_exchange_layout(storage, /*my_rank=*/3, 1, derived); - // Reset drops the cache (relabel copies cores and must not inherit eval-time state). - core.reset_derivative_exchange_layout(); - BOOST_CHECK_EQUAL(core.derivative_exchange_layout().total_count, 16U); + BOOST_CHECK((derived.counts == std::vector{0, 4, 0, 0, 6})); + BOOST_CHECK((derived.displs == std::vector{0, 0, 4, 4, 4})); + BOOST_CHECK_EQUAL(derived.total_count, 10U); + for (size_t r = 1; r < derived.displs.size(); ++r) { + BOOST_CHECK_GE(derived.displs[r], derived.displs[r - 1]); + } } BOOST_AUTO_TEST_CASE(graph_encoding_derivative_exchange_layout_overflow_throws) { - // A count that fits int at 1x but not at 2x. build_layer_storage_unified runs this derivation - // eagerly, so the throw lands in build_graph and not inside the gradient collective window. + // Fits int at 1x but not 2x. The 2x layout is derived eagerly, so the throw lands in build_graph + // rather than inside a collective window where the peers are already committed. + // Declared, not materialised: 2^30 real endpoints exhaust a 16 GB runner. const size_t just_over_half = static_cast(std::numeric_limits::max()) / 2 + 1; + PackedCrossRankStorage storage; + storage.world_size = 2; + storage.occupied.push_back(CrossRankOccupiedSlot{.slot = 0, .sin_send_count = just_over_half, .in_count = 0}); - LayerCore core; - core.evolution_exchange_layout = detail::build_layer_exchange_layout({just_over_half}, 1); - BOOST_CHECK_THROW(detail::build_derivative_exchange_layout(core.evolution_exchange_layout), std::overflow_error); + LayerExchangeLayout derived; + BOOST_CHECK_NO_THROW(detail::derive_exchange_layout(storage, /*my_rank=*/1, 1, derived)); + BOOST_CHECK_THROW(detail::derive_exchange_layout(storage, /*my_rank=*/1, 2, derived), std::overflow_error); } BOOST_AUTO_TEST_CASE(graph_encoding_d_from_b_derivation_both_arms) { @@ -187,3 +237,136 @@ BOOST_AUTO_TEST_CASE(graph_encoding_d_from_b_derivation_both_arms) { BOOST_CHECK_EQUAL(detail::cross_rank_sin_send_index(storage, 0, 0), 10U); BOOST_CHECK_EQUAL(detail::cross_rank_sin_send_index(storage, 0, 4), 22U); } + +// The accounting split behind graph_memory_breakdown(): slot-record cost tracks the world, not traffic. + +BOOST_AUTO_TEST_CASE(graph_encoding_occupied_slots_counts_only_slots_carrying_traffic) { + // Zeros at the front, in the interior and at the back -- the three places a scan loses count. + const auto storage = detail::build_packed_cross_rank_storage(slot_partners({0, 3, 0, 0, 7, 0})); + + BOOST_CHECK_EQUAL(storage.rank_count(), 6U); + BOOST_CHECK_EQUAL(detail::cross_rank_occupied_slots(storage), 2U); + BOOST_CHECK_EQUAL(detail::cross_rank_endpoint_count(storage), 10U); + // An occupied slot holds at least one endpoint. + BOOST_CHECK_LE(detail::cross_rank_occupied_slots(storage), detail::cross_rank_endpoint_count(storage)); +} + +BOOST_AUTO_TEST_CASE(graph_encoding_slot_record_bytes_track_the_traffic_not_the_world) { + // Same traffic, four times the world. + const auto narrow = detail::build_packed_cross_rank_storage(slot_partners({5, 0, 0, 0})); + const auto wide = + detail::build_packed_cross_rank_storage(slot_partners({5, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0})); + + BOOST_CHECK_EQUAL(narrow.rank_count(), 4U); + BOOST_CHECK_EQUAL(wide.rank_count(), 16U); + BOOST_CHECK_EQUAL(detail::cross_rank_endpoint_count(narrow), detail::cross_rank_endpoint_count(wide)); + BOOST_CHECK_EQUAL(detail::cross_rank_occupied_slots(narrow), detail::cross_rank_occupied_slots(wide)); + BOOST_CHECK_EQUAL(detail::cross_rank_slot_record_bytes(narrow), detail::cross_rank_slot_record_bytes(wide)); + BOOST_CHECK_EQUAL(detail::cross_rank_slot_record_bytes(wide), 1U * sizeof(CrossRankOccupiedSlot)); + // A slice of cross_rank_bytes, not an addition to it. + BOOST_CHECK_LT(detail::cross_rank_slot_record_bytes(wide), detail::cross_rank_storage_bytes(wide)); +} + +BOOST_AUTO_TEST_CASE(graph_encoding_a_layer_retains_no_exchange_layout) { + // A built layer holds the slot records and nothing else sized by P. + // Neither side of the exchange is retained. + const auto core = detail::build_layer_storage_unified(slot_partners({3, 0, 5}), /*my_rank=*/1); + + // Recoverable from the slot records alone: 3 + 5, my_rank's own slot excluded. + LayerExchangeLayout derived; + detail::derive_exchange_layout(core->cross_rank, /*my_rank=*/1, /*scale=*/1, derived); + BOOST_CHECK_EQUAL(derived.total_count, 8U); + + // Held L x P times across a job, so a new P-sized member costs once per layer per partition. + BOOST_CHECK_LE(sizeof(LayerCore), 256U); +} + +BOOST_AUTO_TEST_CASE(graph_encoding_skewed_endpoint_counts_are_refused) { + // One count and one offset serve both B and D, and the type does not prevent a skew: unchecked it + // mis-derives Q and reads a wrong-but-valid endpoint rather than throwing. + std::vector data(1); + data[0].sin_send_indices.push_back(1); + data[0].sin_send_indices.push_back(2); + data[0].sin_recv_entries.push_back({1, 1}); // one D against two B + data[0].in_count = 1; + + BOOST_CHECK_THROW(detail::build_packed_cross_rank_storage(data), std::logic_error); +} + +BOOST_AUTO_TEST_CASE(graph_encoding_an_in_block_past_the_endpoint_list_is_refused) { + // in_count bounds a block inside B, and the out-block size is an unsigned subtraction, so one past + // the end wraps rather than going negative. + std::vector data(1); + data[0].sin_send_indices.push_back(1); + data[0].sin_send_indices.push_back(2); + data[0].sin_recv_entries.push_back({1, 1}); + data[0].sin_recv_entries.push_back({2, 1}); + data[0].in_count = 3; // three inside two + + BOOST_CHECK_THROW(detail::build_packed_cross_rank_storage(data), std::logic_error); + + // in_count == B.size() is legal: an all-in slot with an empty out-block. + data[0].in_count = 2; + BOOST_CHECK_NO_THROW(detail::build_packed_cross_rank_storage(data)); +} + +BOOST_AUTO_TEST_CASE(graph_encoding_occupied_sweep_matches_the_dense_sweep_it_replaces) { + // Zeros at front, interior and back, plus a slot whose in_count splits the B list. + auto data = slot_partners({0, 3, 0, 0, 7, 0}); + data[1].in_count = 1; + data[4].in_count = 4; + const auto storage = detail::build_packed_cross_rank_storage(data); + + // The offsets prefix over all slots; the empty ones contributed zero. + struct Expected { + size_t slot, offset, count, in_count; + }; + const std::vector expected{{1, 0, 3, 1}, {4, 3, 7, 4}}; + + std::vector seen; + detail::for_each_occupied_slot(storage, [&](size_t slot, const detail::CrossRankSlotView &view) { + seen.push_back({slot, view.phase_offset, view.sin_send_count, view.in_count}); + }); + + BOOST_REQUIRE_EQUAL(seen.size(), expected.size()); + for (size_t k = 0; k < expected.size(); ++k) { + BOOST_CHECK_EQUAL(seen[k].slot, expected[k].slot); + BOOST_CHECK_EQUAL(seen[k].offset, expected[k].offset); + BOOST_CHECK_EQUAL(seen[k].count, expected[k].count); + BOOST_CHECK_EQUAL(seen[k].in_count, expected[k].in_count); + } + + // The single-slot resolver must agree with the sweep, including on an absent slot. + for (const auto &e : expected) { + const auto view = detail::cross_rank_slot(storage, e.slot); + BOOST_CHECK_EQUAL(view.phase_offset, e.offset); + BOOST_CHECK_EQUAL(view.sin_send_count, e.count); + } + BOOST_CHECK_EQUAL(detail::cross_rank_slot(storage, 0).sin_send_count, 0U); + BOOST_CHECK_EQUAL(detail::cross_rank_slot(storage, 5).sin_send_count, 0U); + BOOST_CHECK_EQUAL(storage.sin_send_size(3), 0U); + BOOST_CHECK_EQUAL(storage.sin_send_size(4), 7U); +} + +BOOST_AUTO_TEST_CASE(graph_encoding_self_slot_is_resolved_without_a_search) { + auto storage = detail::build_packed_cross_rank_storage(slot_partners({0, 3, 0, 7, 0})); + + detail::resolve_self_slot(storage, 3); + BOOST_CHECK_EQUAL(storage.self_offset, 3U); // slot 1's three endpoints precede it + const auto self = detail::cross_rank_self_slot(storage); + BOOST_CHECK_EQUAL(self.sin_send_count, 7U); + BOOST_CHECK_EQUAL(self.phase_offset, 3U); + + // A rank whose own slot carries nothing resolves to an empty view, not a neighbour's. + detail::resolve_self_slot(storage, 2); + BOOST_CHECK_EQUAL(storage.self_pos, kNoSelfSlot); + BOOST_CHECK_EQUAL(detail::cross_rank_self_slot(storage).sin_send_count, 0U); +} + +BOOST_AUTO_TEST_CASE(graph_encoding_a_layer_with_no_cross_rank_traffic_stores_no_slots) { + const auto storage = detail::build_packed_cross_rank_storage(slot_partners({0, 0, 0, 0, 0, 0, 0, 0})); + + BOOST_CHECK_EQUAL(storage.rank_count(), 8U); // the world is still eight wide + BOOST_CHECK_EQUAL(detail::cross_rank_occupied_slots(storage), 0U); + BOOST_CHECK_EQUAL(detail::cross_rank_slot_record_bytes(storage), 0U); +} diff --git a/cpp/tests/hybrid_comm_tests.cpp b/cpp/tests/hybrid_comm_tests.cpp index 6c8db293..bb6ca426 100644 --- a/cpp/tests/hybrid_comm_tests.cpp +++ b/cpp/tests/hybrid_comm_tests.cpp @@ -14,6 +14,9 @@ // HybridComm transport equivalence: R MPI ranks x S in-process partitions must behave as one flat P=R*S // SPMD world, with only partition 0 touching MPI, exactly as PartitionGroup drives it. +// +// The staging layout never reaches a caller, so no case here distinguishes a consistently transposed +// tiling from the shipped one; only the bit-identity comment in HybridComm.h pins that choice. #include @@ -104,7 +107,7 @@ BOOST_AUTO_TEST_CASE(hybrid_comm_allreduce_sum_global) { } // begin_alltoallv must deliver each source's block contiguously in ascending global source order with -// tags intact (Resolve.h's positional pairing). +// tags intact (Resolve.h's positional pairing). With no known recv counts this drives the fused verb. BOOST_AUTO_TEST_CASE(hybrid_comm_alltoallv_source_order_and_tags) { if (world_size() < 2) { return; @@ -265,6 +268,160 @@ BOOST_AUTO_TEST_CASE(hybrid_comm_alltoallv_resolve_fused) { BOOST_CHECK_EQUAL(failures.load(), 0); } +namespace { + +// Counts depending on both ends of the leg; every 5th round is a high-water round, so the smaller +// rounds after it run over stale staged bytes. +auto pair_count(int g, int d, int round) -> int { + if (round % 5 == 4) { + return (g * 3 + d * 7) % 11 + 1; + } + return (g * 2 + d * 3 + round) % 4; +} + +// Unique per (source, destination, index), so a block from the wrong source cannot match. +auto pair_tag(int g, int d, int j) -> int { + return (g * 100 + d) * 1000 + j; +} + +} // namespace + +// The only direct coverage of HybridComm::alltoallv, and the only count matrix varying along both +// indices, which is what pins the indexing of col_sum_ and recv_col_. +BOOST_AUTO_TEST_CASE(hybrid_comm_alltoallv_pairwise_counts) { + if (world_size() < 2) { + return; + } + const int R = world_size(); + const int rounds = 12; + for (const int S : {1, 2, 3}) { + const int P = R * S; + std::atomic failures{0}; + std::atomic payload_checks{0}; + auto errs = run_hybrid(S, [&](HybridComm &hyb, int u) { + Comm c = Comm::make_hybrid(&hyb, u); + const int g = monoprop::mpi::rank(c); + std::vector sc(static_cast(P)), sd(static_cast(P)); + std::vector rc(static_cast(P)), rd(static_cast(P)); + std::vector send; + std::vector recv; // reassigned per round; stage_recv_'s HWM is what carries stale bytes + for (int round = 0; round < rounds; ++round) { + send.clear(); + int so = 0; + int ro = 0; + for (int d = 0; d < P; ++d) { + const int n = pair_count(g, d, round); + sc[static_cast(d)] = n; + sd[static_cast(d)] = so; + so += n; + for (int j = 0; j < n; ++j) { + send.push_back(pair_tag(g, d, j)); + } + const int m = pair_count(d, g, round); // the transpose: what d sends me + rc[static_cast(d)] = m; + rd[static_cast(d)] = ro; + ro += m; + } + recv.assign(static_cast(ro), -1); + const auto args = monoprop::mpi::FlatAlltoallvArgs{.send = send.data(), + .send_counts = sc.data(), + .send_displs = sd.data(), + .recv = recv.data(), + .recv_counts = rc.data(), + .recv_displs = rd.data()} + .bytes(); + hyb.alltoallv(u, args, monoprop::mpi::datatype::get()); + for (int src = 0; src < P; ++src) { + const int m = rc[static_cast(src)]; + for (int j = 0; j < m; ++j) { + payload_checks.fetch_add(1); + if (recv[static_cast(rd[static_cast(src)] + j)] != pair_tag(src, g, j)) { + failures.fetch_add(1); + } + } + } + } + }); + for (const auto &e : errs) { + BOOST_CHECK(e == nullptr); + } + // Assertions inside count-dependent loops can all be skipped, so count the payload comparisons. + BOOST_CHECK_GT(payload_checks.load(), 0); + BOOST_CHECK_EQUAL(failures.load(), 0); + } +} + +// The same pairwise counts through the fused verb, whose recv staging is sized from the count matrix +// rather than from published rows. rc/rd are outputs here, and are checked. +BOOST_AUTO_TEST_CASE(hybrid_comm_alltoallv_resolve_pairwise_counts) { + if (world_size() < 2) { + return; + } + const int R = world_size(); + const int rounds = 12; + for (const int S : {1, 2, 3}) { + const int P = R * S; + std::atomic failures{0}; + // Two counters: the layout checks run unconditionally, so one combined counter would pass on them. + std::atomic layout_checks{0}; + std::atomic payload_checks{0}; + auto errs = run_hybrid(S, [&](HybridComm &hyb, int u) { + Comm c = Comm::make_hybrid(&hyb, u); + const int g = monoprop::mpi::rank(c); + std::vector sc(static_cast(P)), sd(static_cast(P)); + std::vector rc(static_cast(P)), rd(static_cast(P)); + std::vector send; + std::vector recv; // reused across rounds (staging + recv HWM) + for (int round = 0; round < rounds; ++round) { + send.clear(); + int so = 0; + for (int d = 0; d < P; ++d) { + const int n = pair_count(g, d, round); + sc[static_cast(d)] = n; + sd[static_cast(d)] = so; + so += n; + for (int j = 0; j < n; ++j) { + send.push_back(pair_tag(g, d, j)); + } + } + hyb.alltoallv_resolve(u, + {.send = send.data(), + .send_counts = sc.data(), + .send_displs = sd.data(), + .recv = recv, + .recv_counts = rc.data(), + .recv_displs = rd.data()}, + monoprop::mpi::datatype::get()); + int total = 0; + for (int src = 0; src < P; ++src) { + const int m = pair_count(src, g, round); + layout_checks.fetch_add(1); + if (rc[static_cast(src)] != m || rd[static_cast(src)] != total) { + failures.fetch_add(1); + } + for (int j = 0; j < m; ++j) { + payload_checks.fetch_add(1); + if (recv[static_cast(total + j)] != pair_tag(src, g, j)) { + failures.fetch_add(1); + } + } + total += m; + } + layout_checks.fetch_add(1); + if (static_cast(recv.size()) != total) { + failures.fetch_add(1); + } + } + }); + for (const auto &e : errs) { + BOOST_CHECK(e == nullptr); + } + BOOST_CHECK_GT(layout_checks.load(), 0); + BOOST_CHECK_GT(payload_checks.load(), 0); + BOOST_CHECK_EQUAL(failures.load(), 0); + } +} + BOOST_AUTO_TEST_CASE(hybrid_comm_allreduce_sum_inplace_global) { if (world_size() < 2) { return; diff --git a/cpp/tests/large_cosine_storage_tests.cpp b/cpp/tests/large_cosine_storage_tests.cpp index 2110b355..e40f967f 100644 --- a/cpp/tests/large_cosine_storage_tests.cpp +++ b/cpp/tests/large_cosine_storage_tests.cpp @@ -14,6 +14,8 @@ #include +#include +#include #include #include @@ -45,7 +47,7 @@ BOOST_AUTO_TEST_CASE(pruned_layer_supports_cos_counts_above_u32) { } p.in_count = 12; // in-block size P (D indices are derived from B via in_count) - auto storage = detail::build_layer_storage_unified(std::move(cross_rank), /*my_rank=*/0); + auto storage = detail::build_layer_storage_unified(cross_rank, /*my_rank=*/0); // An engaged pruned_cos stores its filtered cosine list explicitly, so num_cos_inds() reports its // total_count. @@ -64,7 +66,7 @@ BOOST_AUTO_TEST_CASE(pruned_layer_supports_cos_counts_above_u32) { lt.for_each_cross_rank_sin_send_range(1, 0, 1, [&](size_t, size_t i) { b_idx = i; }); BOOST_CHECK_EQUAL(b_idx, 200UL); - // D[0] is derived from B: Q = sin_recv_count - in_count = 20 - 12 = 8, so D[0] = out-block[0] = 100, + // D[0] derives from B: Q = 20 - 12 = 8, so D[0] = out-block[0] = 100, // stored phase = -(out_phases[0]) = -(+1) = -1. size_t d_idx = static_cast(-1); int d_phi = 0; @@ -79,10 +81,14 @@ BOOST_AUTO_TEST_CASE(pruned_layer_supports_cos_counts_above_u32) { // The per-rank cross-rank counts index into one layer's term set, so under the wide build they must // be TermIndex-wide; uint32_t would silently cap a single partition/layer at ~2^32 terms. BOOST_AUTO_TEST_CASE(cross_rank_partner_range_counts_track_term_index_width) { - CrossRankPartnerRange r{}; + CrossRankOccupiedSlot r{}; BOOST_CHECK_EQUAL(sizeof(r.sin_send_count), sizeof(TermIndex)); - BOOST_CHECK_EQUAL(sizeof(r.sin_recv_count), sizeof(TermIndex)); BOOST_CHECK_EQUAL(sizeof(r.in_count), sizeof(TermIndex)); + // One record per occupied slot, so its width scales traffic, not P squared. + // Written out rather than imported: a check that borrows the value it checks cannot fail. + constexpr size_t slot_field = std::max(sizeof(uint32_t), alignof(TermIndex)); + BOOST_CHECK_EQUAL(alignof(CrossRankOccupiedSlot), alignof(TermIndex)); + BOOST_CHECK_EQUAL(sizeof(CrossRankOccupiedSlot), slot_field + 2 * sizeof(TermIndex)); } #if defined(monoprop_WIDE_TERM_INDEX) @@ -105,7 +111,7 @@ BOOST_AUTO_TEST_CASE(cross_rank_sin_send_index_round_trips_above_u32) { BOOST_CHECK_EQUAL(detail::cross_rank_sin_send_index(storage, 1, 0), big_in); BOOST_CHECK_EQUAL(detail::cross_rank_sin_send_index(storage, 1, 1), big_out); - // D[0] is derived from B: Q = sin_recv_count - in_count = 1, so D[0] = out-block[0] = big_out. + // D[0] derives from B: Q = 1, so D[0] = out-block[0] = big_out. BOOST_CHECK_EQUAL(detail::cross_rank_sin_recv_index(storage, 1, 0), big_out); } #endif diff --git a/sonar-project.properties b/sonar-project.properties index 6172475a..43274285 100644 --- a/sonar-project.properties +++ b/sonar-project.properties @@ -21,7 +21,7 @@ sonar.coverageReportPaths=cpp-coverage-sonar.xml,cpp-coverage-through-python-bin # Deliberate design decisions, not defects. Each entry is scoped to one file and one rule so # unrelated findings in the same file still surface. Do not widen these to project-level # patterns. The rationale lives next to the code in every case. -sonar.issue.ignore.multicriteria=barrier,hybridtypes,commvalue1,commvalue2,commvalue3,typeerasure1,typeerasure2,commimplicit,allocrebind,includeorder,layercore +sonar.issue.ignore.multicriteria=barrier,hybridtypes,commvalue1,commvalue2,commvalue3,typeerasure1,typeerasure2,commimplicit,allocrebind,includeorder # S8417 wants seq_cst. PartitionBarrier is a sense-reversing generation barrier whose # acquire/release/relaxed orderings are deliberate and documented in-file; seq_cst would put @@ -71,9 +71,3 @@ sonar.issue.ignore.multicriteria.allocrebind.resourceKey=cpp/monoprop/TypeAliase # defined first. sonar.issue.ignore.multicriteria.includeorder.ruleKey=cpp:S954 sonar.issue.ignore.multicriteria.includeorder.resourceKey=cpp/monoprop/TypeAliases.h - -# S5414 flags LayerCore's one private cache member among public fields. LayerCore is a -# public-field aggregate written from several sites; privatising it would cost ~8 accessors -# for no invariant that reset_derivative_exchange_layout() does not already enforce. -sonar.issue.ignore.multicriteria.layercore.ruleKey=cpp:S5414 -sonar.issue.ignore.multicriteria.layercore.resourceKey=cpp/monoprop/detail/graph_encoding/MPGraphEncodingTypes.h diff --git a/src/monoprop/bindings/binder.h b/src/monoprop/bindings/binder.h index bd2b4796..4076c7ff 100644 --- a/src/monoprop/bindings/binder.h +++ b/src/monoprop/bindings/binder.h @@ -265,7 +265,7 @@ auto bind_monomial_propagator(nb::module_ &mod) -> void { {"inverted_index_bytes", b.inverted_index_bytes}, {"matched_scratch_bytes", b.matched_scratch_bytes}, {"total_bytes", b.total_bytes()}, - // Diagnostics (not part of total_bytes; see the struct). + // Diagnostics, outside total_bytes(). {"d_invidx_dense_bytes", b.inverted_index_dense_bytes}, {"d_invidx_sparse_bytes", b.inverted_index_sparse_bytes}, {"d_invidx_dense_columns", b.inverted_index_dense_columns}, @@ -273,5 +273,23 @@ auto bind_monomial_propagator(nb::module_ &mod) -> void { {"d_state_coeffs_nonzero", b.state_coeffs_nonzero}, {"d_init_operator_entries", b.init_operator_entries}}; }); + + // The graph does not partition: its arrays are indexed by the flat world, so these grow with a P + // the MPI rank count never reveals. + cls.def("graph_memory_breakdown", [](const MonomialPropagator &self) { + const auto b = self.graph_memory_usage(); + return std::map{{"layer_descriptor_bytes", b.layer_descriptor_bytes}, + {"layer_storage_object_bytes", b.layer_storage_object_bytes}, + {"cos_data_bytes", b.cos_data_bytes}, + {"cross_rank_bytes", b.cross_rank_bytes}, + {"exchange_layout_bytes", b.exchange_layout_bytes}, + {"total_bytes", b.total_bytes()}, + // Diagnostics, outside total_bytes(). + {"d_slot_record_bytes", b.slot_record_bytes}, + {"d_layer_cores", b.layer_cores}, + {"d_slot_records", b.slot_records}, + {"d_occupied_slots", b.occupied_slots}, + {"d_cross_rank_endpoints", b.cross_rank_endpoints}}; + }); } } // namespace monoprop::bindings::detail diff --git a/tests/test_monoprop_smoke.py b/tests/test_monoprop_smoke.py index e4addfd5..82160252 100644 --- a/tests/test_monoprop_smoke.py +++ b/tests/test_monoprop_smoke.py @@ -155,3 +155,62 @@ def test_bound_graph_methods_accept_declared_arguments( assert core.graph_layers() is not None core.update_initial_operator(op_dict=problem.operator.terms) + + +# The ledger keys are a measurement contract: an A/B reads this dictionary from two builds and +# subtracts, so a key that appears, disappears or crosses ``total_bytes`` compares two different +# quantities. Nothing on the C++ side notices. +_GRAPH_LEDGER_KEYS = ( + "layer_descriptor_bytes", + "layer_storage_object_bytes", + "cos_data_bytes", + "cross_rank_bytes", + "exchange_layout_bytes", +) +_GRAPH_DIAGNOSTIC_KEYS = ( + "d_slot_record_bytes", + "d_layer_cores", + "d_slot_records", + "d_occupied_slots", + "d_cross_rank_endpoints", +) + + +@parametrize_with_cases( + "problem", cases=CasesFermionicProblem, has_tag="has_commutator_data" +) +def test_graph_memory_breakdown_keys_and_totals(problem, serial_comm): + """``graph_memory_breakdown()`` had no Python-side coverage at all before this.""" + monomial_circuit = problem.monomial_circuit + core = _make_bound_core(problem, serial_comm, schrodinger=False) + _evolve_bound_core(core, monomial_circuit) + + breakdown = core.graph_memory_breakdown() + + # Exact, not a superset: both a dropped key and an undecided new one are what this catches. + assert set(breakdown) == { + *_GRAPH_LEDGER_KEYS, + "total_bytes", + *_GRAPH_DIAGNOSTIC_KEYS, + } + + # total_bytes() sums the ledger keys only; the d_ keys are counts or subsets, so folding one in + # would double-count. + assert breakdown["total_bytes"] == sum(breakdown[k] for k in _GRAPH_LEDGER_KEYS) + assert breakdown["total_bytes"] == core.graph_memory_bytes() + assert breakdown["total_bytes"] > 0 + + # Zero because the per-layer counts/displs are no longer retained. The key is kept rather than + # deleted so an A/B against a build that retained them shows the drop instead of losing the row. + assert breakdown["exchange_layout_bytes"] == 0 + + # The slot array is indexed by the flat world, so slot_records / layer_cores recovers P. Stated as + # divisibility, not `== 1`, because the MPI gate also runs this at 16 partitions per rank. + assert breakdown["d_layer_cores"] > 0 + assert breakdown["d_slot_records"] % breakdown["d_layer_cores"] == 0 + assert breakdown["d_slot_records"] // breakdown["d_layer_cores"] >= 1 + + # An occupied slot holds at least one endpoint, so the endpoint count is the P-independent ceiling. + assert breakdown["d_occupied_slots"] <= breakdown["d_slot_records"] + assert breakdown["d_occupied_slots"] <= breakdown["d_cross_rank_endpoints"] + assert breakdown["d_slot_record_bytes"] <= breakdown["cross_rank_bytes"]