Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
Loading...
Searching...
No Matches
tpl_ca_parallel_engine.H
Go to the documentation of this file.
1/*
2 Aleph_w
3
4 Data structures & Algorithms
5 version 2.0.0b
6 https://github.com/lrleon/Aleph-w
7
8 This file is part of Aleph-w library
9
10 Copyright (c) 2002-2026 Leandro Rabindranath Leon
11
12 Permission is hereby granted, free of charge, to any person obtaining a copy
13 of this software and associated documentation files (the "Software"), to deal
14 in the Software without restriction, including without limitation the rights
15 to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
16 copies of the Software, and to permit persons to whom the Software is
17 furnished to do so, subject to the following conditions:
18
19 The above copyright notice and this permission notice shall be included in all
20 copies or substantial portions of the Software.
21
22 THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
23 IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
24 FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
25 AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
26 LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
27 OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
28 SOFTWARE.
29*/
30
80#ifndef TPL_CA_PARALLEL_ENGINE_H
81#define TPL_CA_PARALLEL_ENGINE_H
82
83#include <array>
84#include <concepts>
85#include <cstddef>
86#include <functional>
87#include <future>
88#include <span>
89#include <thread>
90#include <type_traits>
91#include <utility>
92#include <vector>
93
94#include <thread_pool.H>
95
96#include <ca-tiling.H>
97#include <ca-traits.H>
98#include <tpl_ca_bit_storage.H>
99#include <tpl_ca_concepts.H>
100#include <tpl_ca_neighborhood.H>
101
102namespace Aleph {
103namespace CA {
104
105namespace ca_parallel_detail {
106
108template <typename O>
109struct is_tile : std::false_type
110{
111};
112
113template <std::size_t W, std::size_t H>
114struct is_tile<Tile<W, H>> : std::true_type
115{
116};
117
118template <typename O>
119inline constexpr bool is_tile_v = is_tile<O>::value;
120
121template <typename Storage>
122struct is_bit_cell_storage : std::false_type
123{
124};
125
126template <std::size_t N>
127struct is_bit_cell_storage<Bit_Cell_Storage<N>> : std::true_type
128{
129};
130
131template <typename Lattice>
132inline constexpr bool uses_bit_cell_storage_v = requires {
133 typename Lattice::storage_type;
134} and is_bit_cell_storage<typename Lattice::storage_type>::value;
135
136} // namespace ca_parallel_detail
137
145{
147 ThreadPool *pool = nullptr;
148
153 std::size_t num_partitions = 0;
154
158 std::size_t min_parallel_cells = 4096;
159};
160
185template <typename Lattice, typename Rule, typename Neighborhood, typename Order = RowMajor>
187{
188 static_assert(LatticeLike<Lattice>, "Parallel_Synchronous_Engine requires a LatticeLike lattice");
189 static_assert(NeighborhoodLike<Neighborhood>,
190 "Parallel_Synchronous_Engine requires a NeighborhoodLike neighborhood");
191 static_assert(RuleLike<Rule, Lattice>, "Parallel_Synchronous_Engine requires a RuleLike rule");
192 static_assert(Neighborhood::rank_v == Lattice::rank, "Neighborhood rank must match lattice rank");
193
194public:
196 using rule_type = Rule;
198 using order_type = Order;
199
203
205 static constexpr std::size_t rank = Lattice::rank;
207 static constexpr std::size_t neighbour_count = Neighborhood::size_v;
208
210 using hook_type = std::function<void(std::size_t, const Lattice &)>;
211
212private:
215 Rule rule_;
219 std::size_t step_count_ = 0;
221
222 // -------------------------------------------------------------------
223 // Per-cell update used by both the sequential fallback and the
224 // worker tasks. Captured by value-of-reference; safe to invoke from
225 // multiple workers as long as their `c` arguments are disjoint.
226 // -------------------------------------------------------------------
227 void compute_cell(const coord_type &c) const = delete; // not used
228
236 void update_cell(const coord_type &c, std::span<state_type> nbuf)
237 {
238 const std::size_t neigh_sz = std::min(nh_.size(), nbuf.size());
239 gather_neighbors(nh_, cur_buf_, c, std::span<state_type>(nbuf.data(), neigh_sz));
240 const Cell_Context<rank> ctx{step_count_, c};
243 }
244
245 // -------------------------------------------------------------------
246 // Iteration kernels for each (rank, order) combination. They take a
247 // half-open row range `[r_begin, r_end)` so the parallel scheduler
248 // can hand each worker a disjoint slab.
249 // -------------------------------------------------------------------
250
252 {
253 std::array<state_type, neighbour_count> nbuf{};
254 coord_type c{};
255 for (ca_size_t i = r_begin; i < r_end; ++i)
256 {
257 c[0] = static_cast<ca_index_t>(i);
258 update_cell(c, nbuf);
259 }
260 }
261
263 {
264 std::array<state_type, neighbour_count> nbuf{};
265 const ca_size_t cols = cur_buf_.size(1);
266 coord_type c{};
267 for (ca_size_t i = r_begin; i < r_end; ++i)
268 for (ca_size_t j = 0; j < cols; ++j)
269 {
270 c[0] = static_cast<ca_index_t>(i);
271 c[1] = static_cast<ca_index_t>(j);
272 update_cell(c, nbuf);
273 }
274 }
275
276 template <std::size_t TileW, std::size_t TileH>
278 {
279 std::array<state_type, neighbour_count> nbuf{};
280 const ca_size_t cols = cur_buf_.size(1);
281 coord_type c{};
282 // Walk tiles strictly inside the assigned row slab; the schedule
283 // therefore covers each cell exactly once across all workers.
284 for (ca_size_t ti = r_begin; ti < r_end; ti += TileH)
285 for (ca_size_t tj = 0; tj < cols; tj += TileW)
286 {
287 const ca_size_t i_end = ti + TileH < r_end ? ti + TileH : r_end;
288 const ca_size_t j_end = tj + TileW < cols ? tj + TileW : cols;
289 for (ca_size_t i = ti; i < i_end; ++i)
290 for (ca_size_t j = tj; j < j_end; ++j)
291 {
292 c[0] = static_cast<ca_index_t>(i);
293 c[1] = static_cast<ca_index_t>(j);
294 update_cell(c, nbuf);
295 }
296 }
297 }
298
300 {
301 std::array<state_type, neighbour_count> nbuf{};
302 const ca_size_t n1 = cur_buf_.size(1);
303 const ca_size_t n2 = cur_buf_.size(2);
304 coord_type c{};
305 for (ca_size_t i = r_begin; i < r_end; ++i)
306 for (ca_size_t j = 0; j < n1; ++j)
307 for (ca_size_t k = 0; k < n2; ++k)
308 {
309 c[0] = static_cast<ca_index_t>(i);
310 c[1] = static_cast<ca_index_t>(j);
311 c[2] = static_cast<ca_index_t>(k);
312 update_cell(c, nbuf);
313 }
314 }
315
319 {
320 if constexpr (rank == 1)
321 {
323 }
324 else if constexpr (rank == 2)
325 {
326 if constexpr (ca_parallel_detail::is_tile_v<Order>)
328 else
330 }
331 else if constexpr (rank == 3)
332 {
334 }
335 else
336 {
337 static_assert(rank <= 3, "Parallel_Synchronous_Engine currently supports rank <= 3");
338 }
339 }
340
344 [[nodiscard]] std::size_t effective_partitions(ThreadPool *pool) const noexcept
345 {
346 if (pool == nullptr)
347 return 1;
348 const std::size_t pool_workers = pool->num_threads();
350 if (parts == 0)
351 parts = 1;
352
353 // Cannot split into more pieces than there are rows along axis 0:
354 // an empty range would still cost a task scheduling roundtrip.
355 const ca_size_t outer = cur_buf_.size(0);
356 if (parts > outer)
357 parts = outer == 0 ? 1 : outer;
358
360 return 1;
361 return parts;
362 }
363
364public:
381
391 requires std::default_initializable<Neighborhood>
393 {}
394
399 {
400 return cur_buf_;
401 }
402
407 {
408 return step_count_;
409 }
410
413 {
414 return cur_buf_.extents();
415 }
416
419 {
420 return cfg_;
421 }
422
425 {
426 cfg_ = cfg;
427 }
428
431 template <typename F>
432 void on_pre_step(F &&f)
433 {
434 pre_hook_ = std::forward<F>(f);
435 }
436
438 template <typename F>
439 void on_post_step(F &&f)
440 {
441 post_hook_ = std::forward<F>(f);
442 }
443
459 void step()
460 {
461 if constexpr (HasRefreshHalo<Lattice>)
462 cur_buf_.refresh_halo();
463
464 if (pre_hook_)
466
467 ThreadPool *pool = cfg_.pool != nullptr ? cfg_.pool : &Aleph::default_pool();
468 const std::size_t parts = effective_partitions(pool);
469 const ca_size_t outer = cur_buf_.size(0);
470
471 if (parts <= 1)
472 {
473 // Sequential fast path. Identical iteration to the synchronous
474 // engine and to the parallel branch with parts == 1.
476 }
477 else
478 {
479 // Parallel branch: enqueue one task per partition and wait for
480 // every future. We collect futures up-front so the destructor
481 // semantics of `std::future` cannot block us in unexpected
482 // orders if the rule throws.
483 std::vector<std::future<void>> futures;
484 futures.reserve(parts);
485 const std::array<ca_size_t, rank> ext = cur_buf_.extents();
486 for (std::size_t p = 0; p < parts; ++p)
487 {
489 if (rng.empty())
490 continue;
491 futures.emplace_back(pool->enqueue([this, b = rng.begin, e = rng.end]
492 {
493 this->process_slab(b, e);
494 }));
495 }
496 // Phase 1: wait for every worker to finish without re-throwing,
497 // so no worker is still writing to nxt_buf_ when we propagate.
498 for (auto &f : futures)
499 f.wait();
500 // Phase 2: collect exceptions; re-throw the first one if any.
501 for (auto &f : futures)
502 f.get();
503 }
504
506 ++step_count_;
507
508 if (post_hook_)
510 }
511
521 void run(const std::size_t steps)
522 {
523 for (std::size_t i = 0; i < steps; ++i)
524 step();
525 }
526};
527
528} // namespace CA
529} // namespace Aleph
530
531#endif // TPL_CA_PARALLEL_ENGINE_H
size_t steps
Definition ca-c-api.h:126
size_t cols
Definition ca-c-api.h:105
Domain-decomposition utilities for the parallel CA engine.
Common typedefs and tag types for the Cellular Automata module.
Bit-packed row-major storage for N-dimensional binary CAs.
Lattice that adds boundary-aware access on top of a storage.
void set(const coord_type &c, const state_type &v)
Strict write: throws if c is out of range.
typename Storage::state_type state_type
typename Storage::extents_type extents_type
const extents_type & extents() const noexcept
void swap(Lattice &other) noexcept(noexcept(store_.swap(other.store_)))
O(1) swap.
typename Storage::coord_type coord_type
static constexpr std::size_t rank
ca_size_t size() const noexcept
state_type at(const coord_type &c) const
Strict access: throws if c is out of range.
Parallel synchronous double-buffered engine.
std::function< void(std::size_t, const Lattice &)> hook_type
Hook signature: (step_index, frame).
const Parallel_Engine_Config & config() const noexcept
typename Lattice::extents_type extents_type
static constexpr std::size_t neighbour_count
Number of neighbours.
void update_cell(const coord_type &c, std::span< state_type > nbuf)
Update a single cell c: gather its neighbours, evaluate the rule and write the result into nxt_buf_.
void on_post_step(F &&f)
Register a hook fired after every step().
Parallel_Synchronous_Engine(Lattice initial, Rule r, Parallel_Engine_Config cfg={})
Build an engine with a default-constructed neighbourhood.
void update_slab_1d(const ca_size_t r_begin, const ca_size_t r_end)
const extents_type & extents() const noexcept
std::size_t steps_run() const noexcept
Return the number of completed steps.
void set_config(Parallel_Engine_Config cfg) noexcept
Replace the configuration. Takes effect on the next step().
void run(const std::size_t steps)
Run several synchronous steps.
static constexpr std::size_t rank
Lattice dimension.
void update_slab_2d_tile(const ca_size_t r_begin, const ca_size_t r_end)
void process_slab(const ca_size_t r_begin, const ca_size_t r_end)
Drive the iteration over a single row slab according to Order.
void update_slab_2d_row_major(const ca_size_t r_begin, const ca_size_t r_end)
const Lattice & frame() const noexcept
Return the current frame.
void compute_cell(const coord_type &c) const =delete
std::size_t effective_partitions(ThreadPool *pool) const noexcept
Resolve the effective number of partitions for the next step().
void step()
Apply the rule to every cell once and swap buffers.
void on_pre_step(F &&f)
Register a hook fired before every step().
Parallel_Synchronous_Engine(Lattice initial, Rule r, Neighborhood n, Parallel_Engine_Config cfg={})
Build an engine on top of an existing initial lattice.
void update_slab_3d_row_major(const ca_size_t r_begin, const ca_size_t r_end)
A reusable thread pool for efficient parallel task execution.
auto enqueue(F &&f, Args &&... args) -> std::future< std::invoke_result_t< F, Args... > >
Submit a task for execution and get a future for the result.
Shape (per-axis sizes) of an mdspan, mixing compile-time and run-time extents.
Definition ah-mdspan.H:303
Detect whether a lattice exposes a refresh_halo() method (the hallmark of a Ghost_Lattice or any othe...
Storage + topology that carries the cell values.
Connectivity pattern around a coordinate.
Local transition function (any supported signature).
static mt19937 rng
#define N
Definition fib.C:294
size_t blossom_maximum_cardinality_matching(const GT &g, DynDlist< typename GT::Arc * > &matching, SA sa=SA())
Alias of compute_maximum_cardinality_general_matching().
Definition Blossom.H:466
std::span< const T > Neighbor_View
Read-only view over a contiguous range of neighbour values.
Definition ca-traits.H:90
std::ptrdiff_t ca_index_t
Signed coordinate component used by lattices and neighborhoods.
Definition ca-traits.H:60
void gather_neighbors(const Nbh &nh, const L &lat, const typename L::coord_type &center, std::span< T > out)
Populate out[0..nh.size()) with neighbour values of center.
constexpr bool should_run_sequential(const ca_size_t cells, const ca_size_t num_partitions, const ca_size_t min_cells) noexcept
Decide whether a workload should run sequentially.
Definition ca-tiling.H:163
std::size_t ca_size_t
Unsigned size component used for extents and counts.
Definition ca-traits.H:63
State apply_rule(const Rule &r, const State &s, Neighbor_View< State > v, const Cell_Context< Rank > &ctx)
Invoke a rule, optionally forwarding the per-cell context.
Main namespace for Aleph-w library functions.
Definition ah-arena.H:89
ThreadPool & default_pool()
Return the default shared thread pool instance.
and
Check uniqueness with explicit hash + equality functors.
STL namespace.
Per-cell context handed to rules that need to know "where" and "when" they are firing.
Definition ca-traits.H:106
Configuration for Parallel_Synchronous_Engine.
std::size_t num_partitions
Number of partitions per step.
ThreadPool * pool
Thread pool to schedule on. nullptr means Aleph::default_pool().
std::size_t min_parallel_cells
Workloads with fewer cells than this are executed sequentially without touching the pool.
Half-open integer interval [begin, end).
Definition ca-tiling.H:90
static constexpr Range1D slab(const std::array< ca_size_t, Rank > &extents, const ca_size_t parts, const ca_size_t idx) noexcept
Number of partitions along axis 0 with a balanced split.
Definition ca-tiling.H:190
Iterate a 2D lattice in W x H tiles.
Definition ca-traits.H:200
Detect whether O is a Tile<W, H> instantiation.
static int * k
gsl_rng * r
A modern, efficient thread pool for parallel task execution.
Bit-packed dense storage for boolean cellular automata.
C++20 concepts for the Cellular Automata module.
Neighborhoods catalogue for Aleph::CA.