Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 21 additions & 3 deletions src/core/platform.c
Original file line number Diff line number Diff line change
Expand Up @@ -369,7 +369,7 @@ static uint64_t cache_sysfs_llc_bytes(void) {
}
#endif

uint64_t ray_cache_llc_bytes(void) {
static uint64_t cache_llc_probe(void) {
static uint64_t cached = UINT64_MAX;
if (cached != UINT64_MAX) return cached;
uint64_t bytes = 0;
Expand Down Expand Up @@ -642,7 +642,7 @@ uint32_t ray_physical_core_count(void) {
/* Sum of every level-3 cache instance reported by the processor topology
* (each SYSTEM_LOGICAL_PROCESSOR_INFORMATION cache record is one instance).
* 0 when the query fails. */
uint64_t ray_cache_llc_bytes(void) {
static uint64_t cache_llc_probe(void) {
static uint64_t cached = UINT64_MAX;
if (cached != UINT64_MAX) return cached;
uint64_t bytes = 0;
Expand Down Expand Up @@ -800,7 +800,7 @@ ray_err_t ray_thread_join(ray_thread_t t) {
}

uint32_t ray_thread_count(void) { return 1; }
uint64_t ray_cache_llc_bytes(void) { return 0; }
static uint64_t cache_llc_probe(void) { return 0; }

/* Semaphore — counter-only. Single-threaded so wait never blocks (the
* counter must already be positive when wait fires). */
Expand All @@ -818,3 +818,21 @@ void ray_sem_wait(ray_sem_t* s) {
void ray_sem_signal(ray_sem_t* s) { (*s)++; }

#endif /* RAY_OS_WASM */

#ifdef DEBUG
/* Test pin for the probed LLC size: routing that bounds replicated state by
* the cache (group dense slabs) otherwise picks a different strategy on
* every CI runner. 0 restores the platform probe. */
static uint64_t g_llc_for_test = 0;

void ray_cache_llc_set_for_test(uint64_t bytes) {
g_llc_for_test = bytes;
}
#endif

uint64_t ray_cache_llc_bytes(void) {
#ifdef DEBUG
if (g_llc_for_test) return g_llc_for_test;
#endif
return cache_llc_probe();
}
4 changes: 4 additions & 0 deletions src/core/platform.h
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,10 @@ uint32_t ray_physical_core_count(void);
* 0 when the platform cannot report it. Bounds replicated per-task state
* whose random-access working set must stay cache-resident to scale. */
uint64_t ray_cache_llc_bytes(void);
#ifdef DEBUG
/* Pin ray_cache_llc_bytes to `bytes` (0 = probe again). */
void ray_cache_llc_set_for_test(uint64_t bytes);
#endif

void ray_parallel_begin(void);
void ray_parallel_end(void);
Expand Down
17 changes: 13 additions & 4 deletions test/test_agg_contract.c
Original file line number Diff line number Diff line change
Expand Up @@ -1734,7 +1734,13 @@ static test_result_t test_cancelled_group(void) {
* 20-worker pool and a 100k-slot slab the raw replication (20 slabs) leaves
* most caches, and the run must use at most floor(0.75 * LLC / slab) task
* slabs (never fewer than the pool when everything fits). The result is
* identical either way. */
* identical either way.
*
* The LLC is pinned to 32 MB (the size the route assumes when none is
* reported) for the routed query: ~10 slabs fit, so the bound bites but the
* run is not cache-starved. Unpinned, a small-cache runner (a macOS CI VM
* reports a few MB of L2) fits fewer than three slabs and the route rightly
* switches to partition ownership, failing the strategy assertion. */
static test_result_t test_dense_cache_bound(void) {
ray_pool_destroy();
TEST_ASSERT_EQ_I(ray_pool_init_total(20), RAY_OK);
Expand All @@ -1743,21 +1749,24 @@ static test_result_t test_dense_cache_bound(void) {
"(set cb_t (table [k v] (list (as 'I32 (% (* cb_i 7919) 100000)) (% cb_i 13))))");
TEST_ASSERT_NOT_NULL(setup); TEST_ASSERT_FALSE(RAY_IS_ERR(setup)); ray_release(setup);
agg_route_reset();
ray_cache_llc_set_for_test(32ull << 20);
ray_t* r = ray_eval_str("(select {from:cb_t by:k s:(sum v)})");
TEST_ASSERT_NOT_NULL(r); TEST_ASSERT_FALSE(RAY_IS_ERR(r));
agg_route_stats_t stats = agg_route_stats();
uint64_t llc = ray_cache_llc_bytes();
ray_cache_llc_set_for_test(0); /* before any assert can return */
TEST_ASSERT_NOT_NULL(r); TEST_ASSERT_FALSE(RAY_IS_ERR(r));
TEST_ASSERT_EQ_I(stats.routes[AGG_ROUTE_V2_DENSE], 1);
TEST_ASSERT_EQ_I(stats.dense_strategy, AGG_DENSE_TASK_LOCAL);
TEST_ASSERT_TRUE(stats.dense_tasks >= 2 && stats.dense_tasks <= 20);
TEST_ASSERT_EQ_I(ray_table_nrows(r), 100000);
uint64_t llc = ray_cache_llc_bytes();
if (llc > 0) {
{
size_t block = agg_resolve(OP_SUM, RAY_I64)->state_size;
double slots = (double)stats.dense_local_slots / stats.dense_tasks;
double slab = slots * (block + sizeof(int64_t) + 1);
double budget = (double)llc * 0.75;
uint32_t cap = slab * 20 > budget ? (uint32_t)(budget / slab) : 20;
if (cap < 2) cap = 2;
TEST_ASSERT_TRUE(cap >= 3 && cap < 20); /* the pin makes the bound bite */
TEST_ASSERT_EQ_I(stats.dense_tasks, cap);
}
/* The bounded run computes the same sums as the serial engine. */
Expand Down
Loading