diff --git a/benchmark/CACHE_MIGRATION.md b/benchmark/CACHE_MIGRATION.md new file mode 100644 index 0000000000..b7304d5c30 --- /dev/null +++ b/benchmark/CACHE_MIGRATION.md @@ -0,0 +1,120 @@ +# TTL cache migration experiments + +The production prototype uses ttlcache v3.4.1 without `Start`, callbacks, or +loaders. Reads promote LRU recency without renewing TTL. Writes call +`DeleteExpired` before insertion, so expired recent entries do not displace live +entries. `Has` rejects expired entries immediately and does not promote recency; +`Len` reclaims expiration and returns live entry count. Both differ from +Hashicorp's retained-entry `Contains`/`Len` behavior. Nonpositive TTL means no +expiration, replacing Hashicorp's ten-year sentinel; nonpositive capacity means +unlimited. Expired references remain until writes, length inspection, clearing, +or disposal. Unlimited caches have no general memory bound. + +## Evaluation outcome: blocked (2026-10-05) + +The checked first-write experiment failed the selected 0% latency budget after +ten paired repetitions, each with 100 idle cycles, capacity 50,000, four CPUs, +100 ms TTL, integer keys, and representative file-hash payloads. Filling was +validated before expiration in every trial. After inactivity, Hashicorp retained +zero entries and the adapter retained all 50,000 expired entries. + +| Median of per-run percentiles | Hashicorp | Adapter | +| --- | ---: | ---: | +| Write p95 | 0.0259 ms | 35.8209 ms | +| Write p99 | 0.0468 ms | 42.0160 ms | +| Reader call p99 | 0.0303 ms | 41.9898 ms | + +All six primary latency comparisons failed their Bonferroni-adjusted paired +bootstrap gate. For write p99, the relative-change interval was +**+56,338% to +170,551%**, entirely above the 0% budget. A separate ten-cycle +mutex/block profile attributed **99.76% of mutex contention delay** to +`ttlcache.DeleteExpired` called by the adapter's `Set`. This identifies the +synchronous purge of the entire expired batch as the cause of reader stalls. + +The adapter, process-tree prototype, regression tests, and experiment harness +are retained as evaluation work. Other caches remain on Hashicorp. +The system harness changes, remaining migration, complete benchmark matrix, +repository-wide privileged checks, and system benchmark were not executed: +they remain blocked by this failed prototype gate. Publication of the evaluation +branch was explicitly authorized after the failed gate; this is not an accepted +production migration or a resolution of #1014. No PR was created, and Docker +was not started. + +Passed: adapter/process-tree race and leak tests; ten repetitions of the core +unit tests; process-tree subtree race tests; focused builds and vet; Python +latency-gate regression tests; and whitespace checks. Allocation/heap benchmark +support exists but no memory-plateau acceptance result is claimed. + +Local evidence (not committed): + +- `/home/linux/.cache/node-agent-1014-idle-checked-50000-cpu4/`: raw samples, + source snapshots, binary, manifest, formal `report.json`, and `benchstat.txt`. +- `/home/linux/.cache/node-agent-1014-profiles/`: adapter and Hashicorp + mutex/block profiles and raw profiling experiment output. +- The earlier four-implementation run at + `/home/linux/.cache/node-agent-1014-idle-50000-cpu4/` was interrupted to add + full-batch validation and is marked diagnostic-only. It is not gate evidence. + +## First migration gate + +Run from the repository root with Go 1.27. The output directory must not exist: + +```sh +python3 benchmark/cache-migration.py \ + --scenario idle --capacity 50000 --cpus 4 --trials 100 \ + --output /absolute/path/to/fresh-output +``` + +Each process fills the cache with precomputed file-hash values, waits two TTLs, +then releases concurrent readers and performs the first write. Write latency, +reader call latency, and reader response latency (including scheduling delay +since release) are recorded separately, including individual samples and maxima. +The current experiment also checks that filling completed before the oldest +entry expired, and records physical retention after inactivity. Mutex/block +profiles are needed to confirm reader overlap with cleanup; barrier release alone +does not establish lock contention. + +Every implementation has identical capacity, TTL, values and idle duration. +Active ttlcache must physically expire a readiness sentinel before testing, then +stop and join its sweeper at exit. Hashicorp has exactly one cache per process, +including Go benchmark calibration. Each repetition runs in a fresh process; +the implementation order reverses on alternate repetitions. + +`observations.json` contains samples. `manifest.json` and `sources/` capture the +baseline SHA, toolchain, dimensions, tracked diff and exact experiment sources, +including untracked prototype files. The linked test binary is retained for +replay. No Docker daemon, Kubernetes cluster, or Git remote is changed. + +The gate compares adapter/Hashicorp p95 and p99 for each latency metric. It uses +paired bootstrap intervals of median relative changes, a fixed seed, 100,000 +resamples, and Bonferroni correction over 270 predeclared possible comparisons +for 95% family-wide coverage. An upper bound <=0% passes; a lower bound >0% +fails; otherwise the result is inconclusive. Inconclusive experiments extend from +ten to thirty pairs. Missing, non-finite, or insufficient evidence cannot pass. +Fail/inconclusive exits nonzero and blocks migration. Publishing an evaluation +branch requires explicit authorization and must retain the failed-gate findings. + +## Steady-state and allocation diagnostics + +Use `--scenario parallel` for sampled `RunParallel` p95/p99 latency, and +`--scenario allocations` for separate, uninstrumented ns/op, B/op and allocs/op. +Use `--pattern hits|misses|churn|expiration`, `--payload int|hash|process`, and +capacities 1000/10000/50000 with CPU settings 1/4/16. Allocation diagnostics do +not establish latency equivalence. The latency sample buffer is bounded per +worker, samples every 64 operations, and retains the most recent samples. + +`TestCacheRetainedMemory` is an opt-in isolated experiment enabled by +`CACHE_BENCH_MEMORY=1` and `CACHE_BENCH_IMPL`. It inserts twenty capacity-sized +batches of newly allocated process-tree values, checks stored entry count, and +records post-GC heap bytes after each batch and clearing. Read its raw samples +and heap profiles; a successful cardinality assertion alone does not prove a +stable memory plateau. + +## Scope of an unsuccessful prototype + +If the first-write gate fails, retain the prototype and its evidence for review. +Do not migrate other subsystems, modify the system benchmark harness, start +Docker, or create a PR. Publish the evaluation only on explicit authorization; +do not describe it as a completed migration. The system benchmark safety/telemetry +changes and repository-wide privileged validation remain gated on a successful +prototype. Baseline package tests are not evidence of migration completion. diff --git a/benchmark/cache-migration.py b/benchmark/cache-migration.py new file mode 100644 index 0000000000..fb557d01ae --- /dev/null +++ b/benchmark/cache-migration.py @@ -0,0 +1,177 @@ +#!/usr/bin/env python3 +"""Isolated cache experiments. A failed/inconclusive latency gate exits nonzero. + +Each child runs exactly one scenario, one implementation, one CPU setting and +one repetition. Hashicorp's cache is reused across benchmark calibration. +The default experiment targets the first-write-after-inactivity risk first. +""" + +import argparse +import hashlib +import json +import math +import os +from pathlib import Path +import random +import re +import statistics +import subprocess +import sys + + +# Predeclared upper bound on primary comparisons: steady state has 3 payloads +# * 4 patterns * 3 capacities * 3 CPU settings * 2 percentiles = 216; +# idle has 3 capacities * 3 CPUs * 3 latency metrics * 2 percentiles = 54. +COMPARISONS = 270 + + +def classify(deltas, comparisons=COMPARISONS, resamples=100_000): + """Bonferroni-adjusted paired bootstrap interval of median relative change.""" + if len(deltas) < 10 or any(not math.isfinite(x) for x in deltas): + return {"status": "inconclusive", "reason": "need >=10 finite paired observations"} + rng = random.Random(1014) + count = len(deltas) + bootstrap = sorted( + statistics.median(deltas[rng.randrange(count)] for _ in range(count)) + for _ in range(resamples) + ) + tail = 0.025 / comparisons + lower = bootstrap[int(tail * resamples)] + upper = bootstrap[min(resamples - 1, math.ceil((1 - tail) * resamples) - 1)] + status = "pass" if upper <= 0 else "fail" if lower > 0 else "inconclusive" + return {"status": status, "median_change_percent": statistics.median(deltas), + "lower_percent": lower, "upper_percent": upper, "pairs": count, + "family_comparisons": comparisons, "family_confidence": 0.95} + + +def gate(records, scenario): + rows = {"hashicorp": {}, "adapter": {}} + for row in records: + if row["implementation"] in rows: + rows[row["implementation"]][row["repetition"]] = row + repetitions = sorted(set(rows["hashicorp"]) & set(rows["adapter"])) + metrics = ([f"{name}_p{p}_ns" for name in ("write", "reader_call", "reader_response") + for p in (95, 99)] if scenario == "idle" else ["p95-ns", "p99-ns"]) + decisions = {} + for metric in metrics: + deltas = [] + for repetition in repetitions: + baseline = rows["hashicorp"][repetition].get(metric) + candidate = rows["adapter"][repetition].get(metric) + if (not isinstance(baseline, (int, float)) or not isinstance(candidate, (int, float)) + or not math.isfinite(baseline) or not math.isfinite(candidate) + or baseline <= 0 or candidate < 0): + break + deltas.append(100 * (candidate / baseline - 1)) + decisions[metric] = classify(deltas) if len(deltas) == len(repetitions) else { + "status": "inconclusive", "reason": "missing or invalid latency observations"} + status = ("fail" if any(x["status"] == "fail" for x in decisions.values()) else + "pass" if all(x["status"] == "pass" for x in decisions.values()) else "inconclusive") + return {"status": status, "latency_budget_percent": 0, "metrics": decisions} + + +def run_child(binary, args, implementation, repetition, output): + env = dict(os.environ, CACHE_BENCH_IMPL=implementation, + CACHE_BENCH_CAPACITY=str(args.capacity), CACHE_BENCH_TRIALS=str(args.trials), + CACHE_BENCH_TTL_MS=str(args.ttl_ms), CACHE_BENCH_PATTERN=args.pattern, + CACHE_BENCH_PAYLOAD=args.payload, + CACHE_BENCH_LATENCY="0" if args.scenario == "allocations" else "1", + GOMAXPROCS=str(args.cpus)) + command = [str(binary), "-test.count=1", f"-test.cpu={args.cpus}"] + if args.scenario == "idle": + command += ["-test.run=^TestCacheIdleLatency$", "-test.v"] + else: + command += ["-test.run=^$", "-test.bench=^BenchmarkCacheParallel$", + "-test.benchtime=2s", "-test.benchmem"] + child = subprocess.run(command, env=env, text=True, capture_output=True) + raw_path = output / f"{repetition:02d}-{implementation}.txt" + raw_path.write_text(child.stdout + child.stderr) + if child.returncode: + raise RuntimeError(f"{implementation} child failed; see {raw_path}") + if args.scenario == "idle": + result = next((json.loads(line.split("=", 1)[1]) for line in child.stdout.splitlines() + if line.startswith("CACHE_IDLE_RESULT=")), None) + if result is None: + raise RuntimeError(f"missing idle result in {raw_path}") + else: + line = next((line for line in child.stdout.splitlines() + if line.startswith("BenchmarkCacheParallel")), "") + result = {unit: float(value) for value, unit in + re.findall(r"([\d.eE+-]+)\s+(ns/op|B/op|allocs/op|p95-ns|p99-ns|samples)", line)} + if args.scenario != "allocations" and result.get("samples", 0) < 10_000: + raise RuntimeError(f"insufficient latency samples in {raw_path}") + if not math.isfinite(result.get("ns/op", float("nan"))) or result["ns/op"] <= 0: + raise RuntimeError(f"invalid operation timing in {raw_path}") + result["implementation"] = implementation + result["repetition"] = repetition + return result + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--output", type=Path, required=True) + parser.add_argument("--scenario", choices=("idle", "parallel", "allocations"), default="idle") + parser.add_argument("--capacity", type=int, default=50_000) + parser.add_argument("--cpus", type=int, default=4) + parser.add_argument("--trials", type=int, default=100) + parser.add_argument("--ttl-ms", type=int, default=100) + parser.add_argument("--pattern", choices=("hits", "misses", "churn", "expiration"), default="hits") + parser.add_argument("--payload", choices=("int", "hash", "process"), default="hash") + parser.add_argument("--implementations", nargs="+", choices=("hashicorp", "adapter", "passive", "active"), + default=["hashicorp", "adapter", "passive", "active"]) + parser.add_argument("--repetitions", type=int, default=10) + args = parser.parse_args() + if min(args.capacity, args.cpus, args.trials, args.ttl_ms) <= 0 or args.repetitions < 10: + parser.error("positive dimensions and at least ten repetitions required") + if not {"hashicorp", "adapter"}.issubset(args.implementations): + parser.error("both hashicorp and adapter required for the gate") + if args.scenario == "idle" and args.payload != "hash": + parser.error("the idle scenario uses representative file-hash payloads; select --payload hash") + output = args.output.resolve() + output.mkdir(parents=True, exist_ok=False) + binary = output / "cache.test" + root = Path(__file__).resolve().parent.parent + build_env = dict(os.environ, GOTOOLCHAIN="go1.27.0") + subprocess.run(["go", "test", "-mod=readonly", "-c", "-o", str(binary), "./internal/ttlcache"], + cwd=root, env=build_env, check=True) + manifest = vars(args).copy() + manifest["output"] = str(output) + manifest["baseline_sha"] = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=root, text=True).strip() + manifest["go_version"] = subprocess.check_output(["go", "version"], env=build_env, text=True).strip() + manifest["working_diff"] = subprocess.check_output(["git", "diff"], cwd=root, text=True) + manifest["source_snapshot"] = {} + paths = list((root / "internal/ttlcache").glob("*.go")) + [ + Path(__file__).resolve(), root / "benchmark/cache_migration_test.py", root / "go.mod", root / "go.sum", + root / "pkg/processtree/process_tree_manager.go", root / "pkg/processtree/process_tree_manager_test.go"] + for path in paths: + relative = str(path.relative_to(root)) + target = output / "sources" / relative + target.parent.mkdir(parents=True, exist_ok=True) + source = path.read_bytes() + target.write_bytes(source) + manifest["source_snapshot"][relative] = hashlib.sha256(source).hexdigest() + (output / "manifest.json").write_text(json.dumps(manifest, indent=2)) + records = [] + total = args.repetitions + repetition = 0 + while repetition < total: + order = args.implementations if repetition % 2 == 0 else list(reversed(args.implementations)) + for implementation in order: + print(f"repetition {repetition + 1}/{total}: {implementation}", flush=True) + records.append(run_child(binary, args, implementation, repetition, output)) + (output / "observations.json").write_text(json.dumps(records, indent=2)) + repetition += 1 + if repetition == total: + report = ({"status": "diagnostic", "scenario": "allocations", + "note": "uninstrumented B/op and allocs/op; this is not a latency gate"} + if args.scenario == "allocations" else gate(records, args.scenario)) + if report["status"] == "inconclusive" and total < 30: + total = 30 + print("inconclusive: extending paired experiment to 30 repetitions", flush=True) + (output / "report.json").write_text(json.dumps(report, indent=2)) + print(json.dumps(report, indent=2), flush=True) + return 0 if report["status"] in ("pass", "diagnostic") else 1 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/benchmark/cache_migration_test.py b/benchmark/cache_migration_test.py new file mode 100644 index 0000000000..c8feeb3330 --- /dev/null +++ b/benchmark/cache_migration_test.py @@ -0,0 +1,36 @@ +import importlib.util +from pathlib import Path +import unittest + +spec = importlib.util.spec_from_file_location("cache_migration", Path(__file__).with_name("cache-migration.py")) +module = importlib.util.module_from_spec(spec) +spec.loader.exec_module(module) + + +class LatencyGateTests(unittest.TestCase): + def test_budget_and_uncertainty(self): + for deltas, expected in [([-1] * 10, "pass"), ([1] * 10, "fail"), + ([-1, 1] * 5, "inconclusive"), ([0] * 10, "pass")]: + with self.subTest(expected=expected): + self.assertEqual(module.classify(deltas, resamples=10_000)["status"], expected) + + def test_invalid_evidence(self): + for deltas in [[], [0] * 9, [float("nan")] * 10, [float("inf")] * 10]: + self.assertEqual(module.classify(deltas)["status"], "inconclusive") + + def test_missing_observations_do_not_pass(self): + self.assertEqual(module.gate([], "idle")["status"], "inconclusive") + + def test_nonfinite_and_zero_latencies_do_not_pass(self): + for invalid in [None, "bad", float("inf"), float("nan"), 0, -1]: + records = [] + for repetition in range(10): + for implementation in ("hashicorp", "adapter"): + records.append({"implementation": implementation, "repetition": repetition, + "p95-ns": invalid if implementation == "hashicorp" else 1, + "p99-ns": 1}) + self.assertEqual(module.gate(records, "parallel")["status"], "inconclusive") + + +if __name__ == "__main__": + unittest.main() diff --git a/go.mod b/go.mod index 9d69a3c308..75b6e7b483 100644 --- a/go.mod +++ b/go.mod @@ -30,6 +30,7 @@ require ( github.com/hashicorp/golang-lru/v2 v2.0.7 github.com/iceber/iouring-go v0.0.0-20230403020409-002cfd2e2a90 github.com/inspektor-gadget/inspektor-gadget v0.45.1-0.20251020222545-c91c23581ebf + github.com/jellydator/ttlcache/v3 v3.4.1 github.com/joncrlsn/dque v0.0.0-20241024143830-7723fd131a64 github.com/kubescape/backend v0.0.39 github.com/kubescape/go-logger v0.0.32 @@ -60,6 +61,7 @@ require ( go.opentelemetry.io/otel/sdk v1.43.0 go.opentelemetry.io/otel/sdk/metric v1.43.0 go.opentelemetry.io/otel/trace v1.43.0 + go.uber.org/goleak v1.3.0 go.uber.org/multierr v1.11.0 golang.org/x/net v0.56.0 golang.org/x/sync v0.22.0 diff --git a/go.sum b/go.sum index 981d2408f0..77b1f49255 100644 --- a/go.sum +++ b/go.sum @@ -1104,6 +1104,8 @@ github.com/jcmturner/rpc/v2 v2.0.3 h1:7FXXj8Ti1IaVFpSAziCZWNzbNuZmnvw/i6CqLNdWfZ github.com/jcmturner/rpc/v2 v2.0.3/go.mod h1:VUJYCIDm3PVOEHw8sgt091/20OJjskO/YJki3ELg/Hc= github.com/jedib0t/go-pretty/v6 v6.6.8/go.mod h1:YwC5CE4fJ1HFUDeivSV1r//AmANFHyqczZk+U6BDALU= github.com/jellevandenhooff/dkim v0.0.0-20150330215556-f50fe3d243e1/go.mod h1:E0B/fFc00Y+Rasa88328GlI/XbtyysCtTHZS8h7IrBU= +github.com/jellydator/ttlcache/v3 v3.4.1 h1:bOdXmXiycyK6E6Qjyuj5vl+/vU3SCOoDs8a86NbHjAQ= +github.com/jellydator/ttlcache/v3 v3.4.1/go.mod h1:j7LO12PNghFg5+0v9budMAT4rDK4JY969jb9vOdOBBk= github.com/jeremywohl/flatten v1.0.1/go.mod h1:4AmD/VxjWcI5SRB0n6szE2A6s2fsNHDLO0nAlMHgfLQ= github.com/jessevdk/go-flags v1.5.0/go.mod h1:Fw0T6WPc1dYxT4mKEZRfG5kJhaTDP9pj1c2EWnYs/m4= github.com/jinzhu/copier v0.4.0 h1:w3ciUoD19shMCRargcpm0cm91ytaBhDvuRpz1ODO/U8= diff --git a/internal/ttlcache/benchmark_test.go b/internal/ttlcache/benchmark_test.go new file mode 100644 index 0000000000..40f737aad3 --- /dev/null +++ b/internal/ttlcache/benchmark_test.go @@ -0,0 +1,343 @@ +package ttlcache + +import ( + "encoding/json" + "fmt" + "math" + "os" + "runtime" + "sort" + "strconv" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/armosec/armoapi-go/armotypes" + "github.com/hashicorp/golang-lru/v2/expirable" + jellycache "github.com/jellydator/ttlcache/v3" + "go.uber.org/goleak" +) + +// Invoke exactly one benchmark/scenario per process, with -count=1 and one +// -cpu value. Reuse the fixture across Go's calibration invocations so the +// Hashicorp baseline creates exactly one unstoppable sweeper per process. +var benchmarkFixture any +var benchmarkStop func() + +func TestMain(m *testing.M) { + code := m.Run() + if benchmarkStop != nil { + benchmarkStop() + } + os.Exit(code) +} + +func TestActiveComparisonLifecycle(t *testing.T) { + defer goleak.VerifyNone(t, goleak.IgnoreCurrent()) + for range 10 { + c := newBenchmarkCache[int]("active", 10, time.Second) + c.set(1, 1) + benchmarkStop() + } +} + +type benchmarkCache[V any] struct { + get func(int) (V, bool) + set func(int, V) + clear func() + stored func() int +} + +func newBenchmarkCache[V any](implementation string, capacity int, ttl time.Duration) benchmarkCache[V] { + switch implementation { + case "hashicorp": + c := expirable.NewLRU[int, V](capacity, nil, ttl) + return benchmarkCache[V]{get: c.Get, set: func(k int, v V) { c.Add(k, v) }, clear: c.Purge, stored: c.Len} + case "adapter": + c := New[int, V](capacity, ttl) + return benchmarkCache[V]{get: c.Get, set: c.Set, clear: c.DeleteAll, stored: func() int { + metrics := c.cache.Metrics() + return int(metrics.Insertions - metrics.Evictions) + }} + case "passive", "active": + c := jellycache.New[int, V](jellycache.WithCapacity[int, V](uint64(capacity)), + jellycache.WithTTL[int, V](ttl), jellycache.WithDisableTouchOnHit[int, V]()) + if implementation == "active" { + // Has alone cannot prove the sweeper is running: expiration is checked + // even in passive mode. Wait for physical eviction of a sentinel. + var zero V + c.Set(-1, zero, time.Millisecond) + done := make(chan struct{}) + go func() { c.Start(); close(done) }() + deadline := time.Now().Add(time.Second) + for c.Metrics().Evictions == 0 { + if time.Now().After(deadline) { + panic("active sweeper did not evict its readiness sentinel") + } + time.Sleep(time.Millisecond) + } + benchmarkStop = func() { c.Stop(); <-done } + } + return benchmarkCache[V]{ + get: func(k int) (V, bool) { + if item := c.Get(k); item != nil { + return item.Value(), true + } + var zero V + return zero, false + }, + set: func(k int, v V) { c.Set(k, v, jellycache.DefaultTTL) }, clear: c.DeleteAll, + stored: func() int { + metrics := c.Metrics() + return int(metrics.Insertions - metrics.Evictions) + }, + } + default: + panic("CACHE_BENCH_IMPL must be hashicorp, adapter, passive, or active") + } +} + +func envInt(name string, fallback int) int { + if s := os.Getenv(name); s != "" { + v, err := strconv.Atoi(s) + if err != nil || v <= 0 { + panic(name + " must be a positive integer") + } + return v + } + return fallback +} + +func percentile(samples []time.Duration, quantile float64) float64 { + sort.Slice(samples, func(i, j int) bool { return samples[i] < samples[j] }) + return float64(samples[int(math.Ceil(quantile*float64(len(samples))))-1]) +} + +type benchmarkHashes struct{ SHA1Hash, MD5Hash string } + +func benchmarkProcess(i int) armotypes.Process { + return armotypes.Process{PID: uint32(i + 1), Comm: "worker", Children: []armotypes.Process{ + {PID: uint32(i + 2), Comm: "child"}, + }} +} + +func BenchmarkCacheParallel(b *testing.B) { + if os.Getenv("CACHE_BENCH_IMPL") == "" { + b.Skip("use benchmark/cache-migration.py to run isolated comparisons") + } + switch os.Getenv("CACHE_BENCH_PAYLOAD") { + case "int", "": + runParallel(b, func(i int) int { return i }) + case "hash": + runParallel(b, func(i int) *benchmarkHashes { + return &benchmarkHashes{fmt.Sprintf("%040x", i), fmt.Sprintf("%032x", i)} + }) + case "process": + runParallel(b, benchmarkProcess) + default: + b.Fatal("unknown CACHE_BENCH_PAYLOAD") + } +} + +func runParallel[V any](b *testing.B, makeValue func(int) V) { + capacity := envInt("CACHE_BENCH_CAPACITY", 1000) + pattern := os.Getenv("CACHE_BENCH_PATTERN") + ttl := time.Minute + if pattern == "expiration" { + ttl = time.Millisecond + } + if benchmarkFixture == nil { + benchmarkFixture = newBenchmarkCache[V](os.Getenv("CACHE_BENCH_IMPL"), capacity, ttl) + } + c := benchmarkFixture.(benchmarkCache[V]) + c.clear() + values := make([]V, capacity) + for i := range values { + values[i] = makeValue(i) + c.set(i, values[i]) + } + if pattern != "hits" && pattern != "misses" && pattern != "churn" && pattern != "expiration" { + b.Fatal("CACHE_BENCH_PATTERN must be hits, misses, churn, or expiration") + } + latency := os.Getenv("CACHE_BENCH_LATENCY") == "1" + var worker atomic.Uint64 + var mu sync.Mutex + var samples []time.Duration + b.ReportAllocs() + b.ResetTimer() + b.RunParallel(func(pb *testing.PB) { + seed := worker.Add(1) + // Bounded per-worker samples, separate from allocation measurements. + var local []time.Duration + if latency { + local = make([]time.Duration, 0, 8192) + } + iteration := uint64(0) + for pb.Next() { + seed ^= seed << 13 + seed ^= seed >> 7 + seed ^= seed << 17 + key := int(seed % uint64(capacity)) + measure := latency && iteration%64 == 0 + var start time.Time + if measure { + start = time.Now() + } + switch pattern { + case "hits", "expiration": + if seed%10 < 9 { + _, _ = c.get(key) + } else { + c.set(key, values[key]) + } + case "misses": + if seed%10 < 9 { + _, _ = c.get(key + capacity) + } else { + c.set(key, values[key]) + } + case "churn": + c.set(int(seed), values[key]) + } + if measure { + d := time.Since(start) + if len(local) < cap(local) { + local = append(local, d) + } else { + local[(iteration/64)%uint64(len(local))] = d + } + } + iteration++ + } + if latency { + mu.Lock() + samples = append(samples, local...) + mu.Unlock() + } + }) + b.StopTimer() + if len(samples) != 0 { + b.ReportMetric(percentile(samples, .95), "p95-ns") + b.ReportMetric(percentile(samples, .99), "p99-ns") + b.ReportMetric(float64(len(samples)), "samples") + } +} + +// TestCacheIdleLatency is deliberately separate from the steady-state +// benchmarks: a rare first-write spike must not disappear in a mixed percentile. +func TestCacheIdleLatency(t *testing.T) { + if os.Getenv("CACHE_BENCH_IMPL") == "" { + t.Skip("explicit isolated performance experiment") + } + capacity := envInt("CACHE_BENCH_CAPACITY", 50000) + trials := envInt("CACHE_BENCH_TRIALS", 100) + readers := runtime.GOMAXPROCS(0) + ttl := time.Duration(envInt("CACHE_BENCH_TTL_MS", 100)) * time.Millisecond + c := newBenchmarkCache[*benchmarkHashes](os.Getenv("CACHE_BENCH_IMPL"), capacity, ttl) + values := make([]*benchmarkHashes, capacity) + for i := range values { + values[i] = &benchmarkHashes{fmt.Sprintf("%040x", i), fmt.Sprintf("%032x", i)} + } + writes := make([]time.Duration, 0, trials) + fills := make([]time.Duration, 0, trials) + retained := make([]int, 0, trials) + readCalls := make([]time.Duration, 0, trials*readers) + readResponses := make([]time.Duration, 0, trials*readers) + for range trials { + c.clear() + fill := time.Now() + for i, value := range values { + c.set(i, value) + } + fills = append(fills, time.Since(fill)) + if c.stored() != capacity { + t.Fatalf("filled %d entries, expected %d; increase CACHE_BENCH_TTL_MS", c.stored(), capacity) + } + if _, found := c.get(0); !found { + t.Fatalf("first entry expired during fill (%s); increase CACHE_BENCH_TTL_MS", time.Since(fill)) + } + // Let Hashicorp's asynchronous bucket cleanup finish too. The TTL and + // idle interval are identical for every implementation. + time.Sleep(2 * ttl) + retained = append(retained, c.stored()) + start := make(chan struct{}) + var ready, finished sync.WaitGroup + var release time.Time + calls := make([]time.Duration, readers) + responses := make([]time.Duration, readers) + for reader := range readers { + ready.Add(1) + finished.Add(1) + go func() { + defer finished.Done() + ready.Done() + <-start + call := time.Now() + _, _ = c.get(reader % capacity) + calls[reader] = time.Since(call) + responses[reader] = time.Since(release) + }() + } + ready.Wait() + release = time.Now() + close(start) + write := time.Now() + c.set(capacity, values[0]) + writes = append(writes, time.Since(write)) + finished.Wait() + readCalls = append(readCalls, calls...) + readResponses = append(readResponses, responses...) + } + metrics := map[string]any{"implementation": os.Getenv("CACHE_BENCH_IMPL"), "capacity": capacity, + "cpus": readers, "trials": trials, "ttl_ms": ttl.Milliseconds(), + "fill_samples_ns": fills, "retained_after_idle": retained} + for name, samples := range map[string][]time.Duration{"write": writes, "reader_call": readCalls, "reader_response": readResponses} { + metrics[name+"_p95_ns"] = percentile(samples, .95) + metrics[name+"_p99_ns"] = percentile(samples, .99) + metrics[name+"_max_ns"] = float64(samples[len(samples)-1]) + metrics[name+"_samples_ns"] = samples + } + data, err := json.Marshal(metrics) + if err != nil { + t.Fatal(err) + } + fmt.Printf("CACHE_IDLE_RESULT=%s\n", data) +} + +func TestCacheRetainedMemory(t *testing.T) { + if os.Getenv("CACHE_BENCH_MEMORY") != "1" { + t.Skip("explicit isolated heap experiment") + } + capacity := envInt("CACHE_BENCH_CAPACITY", 50000) + c := newBenchmarkCache[armotypes.Process](os.Getenv("CACHE_BENCH_IMPL"), capacity, time.Minute) + snapshot := func() uint64 { + runtime.GC() + var stats runtime.MemStats + runtime.ReadMemStats(&stats) + runtime.KeepAlive(c) + return stats.HeapAlloc + } + baseline := snapshot() + observations := make([]uint64, 0, 20) + for cycle := range 20 { + for i := range capacity { + key := cycle*capacity + i + c.set(key, benchmarkProcess(key)) + } + if c.stored() != capacity { + t.Fatalf("stored entry count %d differs from capacity %d", c.stored(), capacity) + } + observations = append(observations, snapshot()) + } + populated := snapshot() + c.clear() + cleared := snapshot() + data, err := json.Marshal(map[string]any{"implementation": os.Getenv("CACHE_BENCH_IMPL"), + "capacity": capacity, "payload": "process", "baseline_heap_bytes": baseline, + "heap_samples_bytes": observations, "populated_heap_bytes": populated, "cleared_heap_bytes": cleared}) + if err != nil { + t.Fatal(err) + } + fmt.Printf("CACHE_MEMORY_RESULT=%s\n", data) +} diff --git a/internal/ttlcache/cache.go b/internal/ttlcache/cache.go new file mode 100644 index 0000000000..a3fc337904 --- /dev/null +++ b/internal/ttlcache/cache.go @@ -0,0 +1,63 @@ +// Package ttlcache provides passive, fixed-TTL LRU caches. It starts no +// goroutines and registers no asynchronous callbacks. +// +// Reads update LRU order without extending expiration. Expired entries are +// immediately unreadable, but their references remain until a write, Len, +// DeleteAll, or disposal. Finite capacities bound stored entry count; +// unlimited caches do not have a memory bound. +package ttlcache + +import ( + "time" + + jellycache "github.com/jellydator/ttlcache/v3" +) + +type Cache[K comparable, V any] struct { + cache *jellycache.Cache[K, V] +} + +// New creates a cache. Nonpositive capacity means unlimited; nonpositive TTL +// means no expiration, replacing expirable's ten-year sentinel. +func New[K comparable, V any](capacity int, ttl time.Duration) *Cache[K, V] { + if capacity < 0 { + capacity = 0 + } + if ttl <= 0 { + ttl = jellycache.NoTTL + } + return &Cache[K, V]{cache: jellycache.New[K, V]( + jellycache.WithCapacity[K, V](uint64(capacity)), + jellycache.WithTTL[K, V](ttl), + jellycache.WithDisableTouchOnHit[K, V](), + )} +} + +func (c *Cache[K, V]) Get(key K) (V, bool) { + if item := c.cache.Get(key); item != nil { + return item.Value(), true + } + var zero V + return zero, false +} + +// Set renews expiration and reclaims expired entries before capacity eviction. +func (c *Cache[K, V]) Set(key K, value V) { + c.cache.DeleteExpired() + c.cache.Set(key, value, jellycache.DefaultTTL) +} + +// Has rejects expired entries without updating LRU order or expiration. +// Unlike expirable.Contains, it does not recognize retained expired entries. +func (c *Cache[K, V]) Has(key K) bool { return c.cache.Has(key) } + +func (c *Cache[K, V]) Delete(key K) { c.cache.Delete(key) } + +func (c *Cache[K, V]) DeleteAll() { c.cache.DeleteAll() } + +// Len reclaims expired entries and reports live entries, rather than counting +// retained expired entries as expirable.Len does. +func (c *Cache[K, V]) Len() int { + c.cache.DeleteExpired() + return c.cache.Len() +} diff --git a/internal/ttlcache/cache_test.go b/internal/ttlcache/cache_test.go new file mode 100644 index 0000000000..4687ea4b9e --- /dev/null +++ b/internal/ttlcache/cache_test.go @@ -0,0 +1,146 @@ +package ttlcache + +import ( + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + "go.uber.org/goleak" +) + +func waitUntilExpired(deadline time.Time) { + time.Sleep(time.Until(deadline) + time.Millisecond) +} + +func TestExpirationAndRenewal(t *testing.T) { + c := New[string, int](3, 80*time.Millisecond) + c.Set("a", 1) + deadline := c.cache.Get("a").ExpiresAt() + require.True(t, c.Has("a")) + value, found := c.Get("a") + require.True(t, found) + require.Equal(t, 1, value) + require.Equal(t, deadline, c.cache.Get("a").ExpiresAt(), "Get and Has must not renew TTL") + waitUntilExpired(deadline) + require.False(t, c.Has("a")) + _, found = c.Get("a") + require.False(t, found) + require.Zero(t, c.Len()) + + c.Set("a", 2) + deadline = c.cache.Get("a").ExpiresAt() + time.Sleep(5 * time.Millisecond) + c.Set("a", 3) + renewed := c.cache.Get("a").ExpiresAt() + require.True(t, renewed.After(deadline)) + value, found = c.Get("a") + require.True(t, found) + require.Equal(t, 3, value) +} + +func TestRecency(t *testing.T) { + t.Run("Get promotes", func(t *testing.T) { + c := New[string, int](2, time.Minute) + c.Set("a", 1) + c.Set("b", 2) + _, _ = c.Get("a") + c.Set("c", 3) + require.True(t, c.Has("a")) + require.False(t, c.Has("b")) + require.Equal(t, 2, c.Len()) + }) + t.Run("Has does not promote", func(t *testing.T) { + c := New[string, int](2, time.Minute) + c.Set("a", 1) + c.Set("b", 2) + require.True(t, c.Has("a")) + c.Set("c", 3) + require.False(t, c.Has("a")) + require.True(t, c.Has("b")) + }) +} + +func TestExpiredRecentEntryIsReclaimedBeforeLiveEviction(t *testing.T) { + c := New[string, int](2, time.Minute) + c.Set("live", 1) + // Make the expired entry more recent than the live one. A plain Set without + // DeleteExpired would evict "live" when inserting "new". + c.cache.Set("expired", 2, time.Millisecond) + waitUntilExpired(c.cache.Get("expired").ExpiresAt()) + c.Set("new", 3) + require.True(t, c.Has("live")) + require.True(t, c.Has("new")) + require.False(t, c.Has("expired")) + require.Equal(t, 2, c.Len()) +} + +func TestZeroNilDeleteAndClear(t *testing.T) { + c := New[string, *int](2, time.Minute) + c.Set("nil", nil) + value, found := c.Get("nil") + require.True(t, found) + require.Nil(t, value) + c.Delete("missing") + c.Delete("nil") + require.False(t, c.Has("nil")) + c.Set("nil", nil) + c.DeleteAll() + require.Zero(t, c.Len()) + zero := New[string, int](1, time.Minute) + zero.Set("zero", 0) + v, found := zero.Get("zero") + require.True(t, found) + require.Zero(t, v) +} + +func TestNonpositiveCapacityAndTTL(t *testing.T) { + for _, capacity := range []int{0, -1} { + for _, ttl := range []time.Duration{0, -time.Second} { + c := New[int, int](capacity, ttl) + for i := range 100 { + c.Set(i, i) + } + require.Equal(t, 100, c.Len()) + require.True(t, c.cache.Get(0).ExpiresAt().IsZero()) + } + } +} + +func TestConcurrentOperations(t *testing.T) { + c := New[int, int](100, time.Millisecond) + var wg sync.WaitGroup + for worker := range 8 { + wg.Add(1) + go func() { + defer wg.Done() + for i := range 2000 { + key := (i + worker) % 200 + switch i % 6 { + case 0: + c.Set(key, i) + case 1: + _, _ = c.Get(key) + case 2: + _ = c.Has(key) + case 3: + c.Delete(key) + case 4: + _ = c.Len() + case 5: + c.DeleteAll() + } + } + }() + } + wg.Wait() + require.LessOrEqual(t, c.Len(), 100) +} + +func TestConstructionAndDiscardStartsNoGoroutines(t *testing.T) { + defer goleak.VerifyNone(t, goleak.IgnoreCurrent()) + for range 1000 { + c := New[int, int](10, time.Millisecond) + c.Set(1, 1) + } +} diff --git a/pkg/processtree/process_tree_manager.go b/pkg/processtree/process_tree_manager.go index 0d5373d1bf..8163f942d3 100644 --- a/pkg/processtree/process_tree_manager.go +++ b/pkg/processtree/process_tree_manager.go @@ -6,7 +6,7 @@ import ( "time" "github.com/armosec/armoapi-go/armotypes" - "github.com/hashicorp/golang-lru/v2/expirable" + "github.com/kubescape/node-agent/internal/ttlcache" "github.com/kubescape/node-agent/pkg/config" containerprocesstree "github.com/kubescape/node-agent/pkg/processtree/container" "github.com/kubescape/node-agent/pkg/processtree/conversion" @@ -23,7 +23,7 @@ type treeCacheKey struct { type ProcessTreeManagerImpl struct { creator processtreecreator.ProcessTreeCreator containerTree containerprocesstree.ContainerProcessTree - containerProcessTreeCache *expirable.LRU[treeCacheKey, armotypes.Process] // (containerID, pid) -> cached result + containerProcessTreeCache *ttlcache.Cache[treeCacheKey, armotypes.Process] // (containerID, pid) -> cached result mutex sync.RWMutex config config.Config } @@ -35,7 +35,7 @@ func NewProcessTreeManager( config config.Config, ) ProcessTreeManager { - containerProcessTreeCache := expirable.NewLRU[treeCacheKey, armotypes.Process](10000, nil, 1*time.Minute) + containerProcessTreeCache := ttlcache.New[treeCacheKey, armotypes.Process](10000, 1*time.Minute) ptm := &ProcessTreeManagerImpl{ creator: creator, @@ -95,7 +95,7 @@ func (ptm *ProcessTreeManagerImpl) GetContainerProcessTree(containerID string, p } // Cache the result - ptm.containerProcessTreeCache.Add(cacheKey, containerSubtree) + ptm.containerProcessTreeCache.Set(cacheKey, containerSubtree) return containerSubtree, nil } diff --git a/pkg/processtree/process_tree_manager_test.go b/pkg/processtree/process_tree_manager_test.go index ca4d039197..9319dac423 100644 --- a/pkg/processtree/process_tree_manager_test.go +++ b/pkg/processtree/process_tree_manager_test.go @@ -5,8 +5,10 @@ import ( "testing" "time" + "github.com/armosec/armoapi-go/armotypes" containercollection "github.com/inspektor-gadget/inspektor-gadget/pkg/container-collection" eventtypes "github.com/inspektor-gadget/inspektor-gadget/pkg/types" + "github.com/kubescape/node-agent/internal/ttlcache" "github.com/kubescape/node-agent/pkg/config" "github.com/kubescape/node-agent/pkg/ebpf/events" containerprocesstree "github.com/kubescape/node-agent/pkg/processtree/container" @@ -14,6 +16,7 @@ import ( "github.com/kubescape/node-agent/pkg/utils" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "go.uber.org/goleak" ) // TestManager_GetProcessBootTimeNs exercises the accessor that the network-stream @@ -85,4 +88,34 @@ func TestManager_GetContainerProcessTree_CarriesStartTime(t *testing.T) { // The boot-relative identity for the same process, from the side map. assert.Equal(t, uint64(2_220_000_000), mgr.GetProcessBootTimeNs(self)) + + impl := mgr.(*ProcessTreeManagerImpl) + key := treeCacheKey{containerID: containerID, pid: self} + impl.containerProcessTreeCache.Set(key, armotypes.Process{Comm: "cached"}) + cached, err := mgr.GetContainerProcessTree(containerID, self, true) + require.NoError(t, err) + require.Equal(t, "cached", cached.Comm) + bypassed, err := mgr.GetContainerProcessTree(containerID, self, false) + require.NoError(t, err) + require.Equal(t, "nginx", bypassed.Comm) + + impl.containerProcessTreeCache = ttlcache.New[treeCacheKey, armotypes.Process](10, time.Millisecond) + impl.containerProcessTreeCache.Set(key, armotypes.Process{Comm: "expired"}) + time.Sleep(5 * time.Millisecond) + expired, err := mgr.GetContainerProcessTree(containerID, self, true) + require.NoError(t, err) + require.Equal(t, "nginx", expired.Comm) +} + +func TestManagerCacheLifecycleDoesNotLeak(t *testing.T) { + defer goleak.VerifyNone(t, goleak.IgnoreCurrent()) + for range 10 { + tree := containerprocesstree.NewContainerProcessTree() + cfg := config.Config{} + cfg.ExitCleanup.CleanupInterval = time.Minute + creator := processtreecreator.NewProcessTreeCreator(tree, cfg) + mgr := NewProcessTreeManager(creator, tree, cfg) + mgr.Start() + mgr.Stop() + } }