Skip to content

Commit cf8eaf9

Browse files
committed
Harden cgroups snapshot cache concurrency
1 parent 96f5f29 commit cf8eaf9

49 files changed

Lines changed: 2846 additions & 809 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.agents/sow/done/SOW-0025-20260627-cgroups-cache-concurrency-contract.md‎

Lines changed: 534 additions & 0 deletions
Large diffs are not rendered by default.

‎bench/drivers/c/bench_posix.c‎

Lines changed: 27 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1041,12 +1041,22 @@ static int run_lookup_bench(int duration_sec)
10411041
}
10421042

10431043
nipc_cgroups_cache_t cache;
1044-
memset(&cache, 0, sizeof(cache));
1045-
cache.items = items;
1046-
cache.item_count = 16;
1047-
cache.populated = true;
1048-
cache.systemd_enabled = 1;
1049-
cache.generation = 1;
1044+
nipc_client_config_t ccfg = {
1045+
.supported_profiles = NIPC_PROFILE_BASELINE,
1046+
.max_request_batch_items = 1,
1047+
.max_response_payload_bytes = RESPONSE_BUF_SIZE,
1048+
.auth_token = AUTH_TOKEN,
1049+
};
1050+
nipc_cgroups_cache_init(&cache, "/tmp", "bench_lookup", &ccfg);
1051+
if (!nipc_cgroups_cache_seed_for_tests(&cache, items, 16, 1, 1)) {
1052+
fprintf(stderr, "lookup: failed to seed cache\n");
1053+
for (int i = 0; i < 16; i++) {
1054+
free(items[i].name);
1055+
free(items[i].path);
1056+
}
1057+
nipc_cgroups_cache_close(&cache);
1058+
return 1;
1059+
}
10501060

10511061
uint64_t lookups = 0;
10521062
uint64_t hits = 0;
@@ -1057,12 +1067,16 @@ static int run_lookup_bench(int duration_sec)
10571067

10581068
while (now_ns() < wall_end) {
10591069
/* Cycle through all 16 items */
1060-
for (int i = 0; i < 16; i++) {
1061-
const nipc_cgroups_cache_item_t *found =
1062-
nipc_cgroups_cache_lookup(&cache, items[i].hash, items[i].name);
1063-
if (found)
1064-
hits++;
1065-
lookups++;
1070+
nipc_cgroups_cache_read_guard_t guard;
1071+
if (nipc_cgroups_cache_read_lock(&cache, &guard)) {
1072+
for (int i = 0; i < 16; i++) {
1073+
const nipc_cgroups_cache_item_view_t *found =
1074+
nipc_cgroups_cache_get(&guard, items[i].hash, items[i].name);
1075+
if (found)
1076+
hits++;
1077+
lookups++;
1078+
}
1079+
nipc_cgroups_cache_read_unlock(&guard);
10661080
}
10671081
}
10681082

@@ -1088,6 +1102,7 @@ static int run_lookup_bench(int duration_sec)
10881102
free(items[i].name);
10891103
free(items[i].path);
10901104
}
1105+
nipc_cgroups_cache_close(&cache);
10911106

10921107
return (hits == lookups) ? 0 : 1;
10931108
}

‎bench/drivers/c/bench_windows.c‎

Lines changed: 27 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1453,12 +1453,22 @@ static int run_lookup_bench(int duration_sec)
14531453
}
14541454

14551455
nipc_cgroups_cache_t cache;
1456-
memset(&cache, 0, sizeof(cache));
1457-
cache.items = items;
1458-
cache.item_count = 16;
1459-
cache.populated = true;
1460-
cache.systemd_enabled = 1;
1461-
cache.generation = 1;
1456+
nipc_client_config_t ccfg = {
1457+
.supported_profiles = NIPC_PROFILE_BASELINE,
1458+
.max_request_batch_items = 1,
1459+
.max_response_payload_bytes = RESPONSE_BUF_SIZE,
1460+
.auth_token = AUTH_TOKEN,
1461+
};
1462+
nipc_cgroups_cache_init(&cache, ".", "bench_lookup", &ccfg);
1463+
if (!nipc_cgroups_cache_seed_for_tests(&cache, items, 16, 1, 1)) {
1464+
fprintf(stderr, "lookup: failed to seed cache\n");
1465+
for (int i = 0; i < 16; i++) {
1466+
free(items[i].name);
1467+
free(items[i].path);
1468+
}
1469+
nipc_cgroups_cache_close(&cache);
1470+
return 1;
1471+
}
14621472

14631473
uint64_t lookups = 0;
14641474
uint64_t hits = 0;
@@ -1468,12 +1478,16 @@ static int run_lookup_bench(int duration_sec)
14681478
ULONGLONG tick_deadline = GetTickCount64() + (ULONGLONG)duration_sec * 1000;
14691479

14701480
while (GetTickCount64() < tick_deadline) {
1471-
for (int i = 0; i < 16; i++) {
1472-
const nipc_cgroups_cache_item_t *found =
1473-
nipc_cgroups_cache_lookup(&cache, items[i].hash, items[i].name);
1474-
if (found)
1475-
hits++;
1476-
lookups++;
1481+
nipc_cgroups_cache_read_guard_t guard;
1482+
if (nipc_cgroups_cache_read_lock(&cache, &guard)) {
1483+
for (int i = 0; i < 16; i++) {
1484+
const nipc_cgroups_cache_item_view_t *found =
1485+
nipc_cgroups_cache_get(&guard, items[i].hash, items[i].name);
1486+
if (found)
1487+
hits++;
1488+
lookups++;
1489+
}
1490+
nipc_cgroups_cache_read_unlock(&guard);
14771491
}
14781492
}
14791493

@@ -1498,6 +1512,7 @@ static int run_lookup_bench(int duration_sec)
14981512
free(items[i].name);
14991513
free(items[i].path);
15001514
}
1515+
nipc_cgroups_cache_close(&cache);
15011516

15021517
return (hits == lookups) ? 0 : 1;
15031518
}

‎bench/drivers/go/main.go‎

Lines changed: 13 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -49,19 +49,6 @@ const (
4949
lookupMethodBufBytesPerItem = 256
5050
)
5151

52-
type cacheBucket struct {
53-
index int
54-
used bool
55-
}
56-
57-
func cacheHashName(name string) uint32 {
58-
h := uint32(5381)
59-
for i := 0; i < len(name); i++ {
60-
h = ((h << 5) + h) + uint32(name[i])
61-
}
62-
return h
63-
}
64-
6552
// ---------------------------------------------------------------------------
6653
// Timing helpers
6754
// ---------------------------------------------------------------------------
@@ -1216,19 +1203,18 @@ func runLookupBench(durationSec int) int {
12161203
}
12171204
}
12181205

1219-
lookupIndex := make([]cacheBucket, 32)
1220-
mask, ok := u32FromNonNegativeInt(len(lookupIndex) - 1)
1221-
if !ok {
1222-
fmt.Fprintln(os.Stderr, "cache lookup index too large")
1206+
cache := rawsvc.NewCache("/tmp", "bench_lookup", posix.ClientConfig{
1207+
SupportedProfiles: protocol.ProfileBaseline,
1208+
PreferredProfiles: 0,
1209+
MaxRequestBatchItems: 1,
1210+
MaxResponsePayloadBytes: responseBufSize,
1211+
MaxResponseBatchItems: 1,
1212+
AuthToken: authToken,
1213+
})
1214+
if !cache.SeedForTests(cacheItems, 1, 1) {
1215+
fmt.Fprintln(os.Stderr, "cache seed failed")
12231216
return 1
12241217
}
1225-
for i := range cacheItems {
1226-
slot := (cacheItems[i].Hash ^ cacheHashName(cacheItems[i].Name)) & mask
1227-
for lookupIndex[slot].used {
1228-
slot = (slot + 1) & mask
1229-
}
1230-
lookupIndex[slot] = cacheBucket{index: i, used: true}
1231-
}
12321218

12331219
var lookups, hits uint64
12341220

@@ -1237,22 +1223,14 @@ func runLookupBench(durationSec int) int {
12371223
deadline := time.Duration(durationSec) * time.Second
12381224

12391225
for time.Since(wallStart) < deadline {
1226+
guard := cache.ReadLock()
12401227
for _, it := range items {
1241-
slot := (it.hash ^ cacheHashName(it.name)) & mask
1242-
found := false
1243-
for lookupIndex[slot].used {
1244-
bucketItem := &cacheItems[lookupIndex[slot].index]
1245-
if bucketItem.Hash == it.hash && bucketItem.Name == it.name {
1246-
found = true
1247-
break
1248-
}
1249-
slot = (slot + 1) & mask
1250-
}
1251-
if found {
1228+
if guard.Get(it.hash, it.name) != nil {
12521229
hits++
12531230
}
12541231
lookups++
12551232
}
1233+
guard.Unlock()
12561234
}
12571235

12581236
wallSec := time.Since(wallStart).Seconds()

‎bench/drivers/go/main_windows.go‎

Lines changed: 14 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -53,19 +53,6 @@ const (
5353
clientShmAttachRetryTimeoutWin = 5 * time.Second
5454
)
5555

56-
type cacheBucketWin struct {
57-
index int
58-
used bool
59-
}
60-
61-
func cacheHashNameWin(name string) uint32 {
62-
h := uint32(5381)
63-
for i := 0; i < len(name); i++ {
64-
h = ((h << 5) + h) + uint32(name[i])
65-
}
66-
return h
67-
}
68-
6956
type xorshift32Win struct {
7057
state uint32
7158
}
@@ -1840,14 +1827,17 @@ func runLookupBenchWin(durationSec int) int {
18401827
}
18411828
}
18421829

1843-
lookupIndex := make([]cacheBucketWin, 32)
1844-
mask := uint32(len(lookupIndex) - 1)
1845-
for i := range cacheItems {
1846-
slot := (cacheItems[i].Hash ^ cacheHashNameWin(cacheItems[i].Name)) & mask
1847-
for lookupIndex[slot].used {
1848-
slot = (slot + 1) & mask
1849-
}
1850-
lookupIndex[slot] = cacheBucketWin{index: i, used: true}
1830+
cache := rawsvc.NewCache(".", "bench_lookup", windows.ClientConfig{
1831+
SupportedProfiles: protocol.ProfileBaseline,
1832+
PreferredProfiles: 0,
1833+
MaxRequestBatchItems: 1,
1834+
MaxResponsePayloadBytes: responseBufSizeWin,
1835+
MaxResponseBatchItems: 1,
1836+
AuthToken: authTokenWin,
1837+
})
1838+
if !cache.SeedForTests(cacheItems, 1, 1) {
1839+
fmt.Fprintln(os.Stderr, "cache seed failed")
1840+
return 1
18511841
}
18521842

18531843
var lookups, hits uint64
@@ -1857,22 +1847,14 @@ func runLookupBenchWin(durationSec int) int {
18571847
tickDeadline := tickMSWin() + uint64(durationSec)*1000
18581848

18591849
for tickMSWin() < tickDeadline {
1850+
guard := cache.ReadLock()
18601851
for _, it := range items {
1861-
slot := (it.hash ^ cacheHashNameWin(it.name)) & mask
1862-
found := false
1863-
for lookupIndex[slot].used {
1864-
bucketItem := cacheItems[lookupIndex[slot].index]
1865-
if bucketItem.Hash == it.hash && bucketItem.Name == it.name {
1866-
found = true
1867-
break
1868-
}
1869-
slot = (slot + 1) & mask
1870-
}
1871-
if found {
1852+
if guard.Get(it.hash, it.name) != nil {
18721853
hits++
18731854
}
18741855
lookups++
18751856
}
1857+
guard.Unlock()
18761858
}
18771859

18781860
wallSec := float64(nowNS()-wallStart) / 1e9

‎bench/drivers/rust/src/bench_windows.rs‎

Lines changed: 21 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,8 @@ use netipc::protocol::{
2727
use netipc::service::cgroups::{CgroupsClient, ClientConfig as TypedClientConfig};
2828
#[cfg(windows)]
2929
use netipc::service::raw::{
30-
increment_dispatch, snapshot_dispatch, CgroupsCacheItem, DispatchHandler, ManagedServer,
30+
increment_dispatch, snapshot_dispatch, CgroupsCache, CgroupsCacheItem, DispatchHandler,
31+
ManagedServer,
3132
};
3233
#[cfg(windows)]
3334
use netipc::transport::win_shm::{WinShmContext, PROFILE_HYBRID as WIN_SHM_PROFILE_HYBRID};
@@ -1758,20 +1759,6 @@ fn run_pipeline_batch_client(
17581759

17591760
#[cfg(windows)]
17601761
fn run_lookup_bench(duration_sec: u32) -> i32 {
1761-
#[derive(Clone, Copy, Default)]
1762-
struct Bucket {
1763-
index: usize,
1764-
used: bool,
1765-
}
1766-
1767-
fn hash_name(name: &str) -> u32 {
1768-
let mut h: u32 = 5381;
1769-
for b in name.as_bytes() {
1770-
h = ((h << 5).wrapping_add(h)).wrapping_add(*b as u32);
1771-
}
1772-
h
1773-
}
1774-
17751762
let items: Vec<CgroupsCacheItem> = (0..16)
17761763
.map(|i| CgroupsCacheItem {
17771764
hash: 1000 + i,
@@ -1782,33 +1769,30 @@ fn run_lookup_bench(duration_sec: u32) -> i32 {
17821769
})
17831770
.collect();
17841771

1785-
let mut lookup_index = vec![Bucket::default(); 32];
1786-
let mask = (lookup_index.len() - 1) as u32;
1787-
for (idx, item) in items.iter().enumerate() {
1788-
let mut slot = (item.hash ^ hash_name(&item.name)) & mask;
1789-
while lookup_index[slot as usize].used {
1790-
slot = (slot + 1) & mask;
1791-
}
1792-
lookup_index[slot as usize] = Bucket {
1793-
index: idx,
1794-
used: true,
1795-
};
1796-
}
1772+
let cache = CgroupsCache::new(
1773+
".",
1774+
"bench_lookup",
1775+
ClientConfig {
1776+
supported_profiles: PROFILE_BASELINE,
1777+
preferred_profiles: 0,
1778+
max_request_payload_bytes: 0,
1779+
max_request_batch_items: 1,
1780+
max_response_payload_bytes: RESPONSE_BUF_SIZE as u32,
1781+
max_response_batch_items: 1,
1782+
auth_token: AUTH_TOKEN,
1783+
packet_size: 0,
1784+
},
1785+
);
1786+
cache.seed_for_tests(items.clone(), 1, 1);
17971787

17981788
let mut lookups: u64 = 0;
17991789
let mut hits: u64 = 0;
18001790

18011791
let warmup_deadline = TickDeadline::new(LOOKUP_WARMUP_SEC);
18021792
while !warmup_deadline.expired() {
1793+
let guard = cache.read_lock();
18031794
for item in &items {
1804-
let mut slot = (item.hash ^ hash_name(&item.name)) & mask;
1805-
while lookup_index[slot as usize].used {
1806-
let bucket_item = &items[lookup_index[slot as usize].index];
1807-
if bucket_item.hash == item.hash && bucket_item.name == item.name {
1808-
break;
1809-
}
1810-
slot = (slot + 1) & mask;
1811-
}
1795+
let _ = guard.get(item.hash, &item.name);
18121796
}
18131797
}
18141798

@@ -1817,18 +1801,9 @@ fn run_lookup_bench(duration_sec: u32) -> i32 {
18171801
let tick_deadline = TickDeadline::new(duration_sec);
18181802

18191803
while !tick_deadline.expired() {
1804+
let guard = cache.read_lock();
18201805
for item in &items {
1821-
let mut slot = (item.hash ^ hash_name(&item.name)) & mask;
1822-
let mut found = false;
1823-
while lookup_index[slot as usize].used {
1824-
let bucket_item = &items[lookup_index[slot as usize].index];
1825-
if bucket_item.hash == item.hash && bucket_item.name == item.name {
1826-
found = true;
1827-
break;
1828-
}
1829-
slot = (slot + 1) & mask;
1830-
}
1831-
if found {
1806+
if guard.get(item.hash, &item.name).is_some() {
18321807
hits += 1;
18331808
}
18341809
lookups += 1;

0 commit comments

Comments
 (0)