diff --git a/bench/warmup_test.go b/bench/warmup_test.go new file mode 100644 index 0000000..c7ff9d0 --- /dev/null +++ b/bench/warmup_test.go @@ -0,0 +1,182 @@ +package bench_test + +import ( + "fmt" + "strings" + "testing" + + "github.com/stretchr/testify/require" + + ascache "github.com/sshaplygin/as-cache" + "github.com/sshaplygin/as-cache/bench" +) + +// switchOnceTo asks for a switch to target on its first selection and for no +// change on every later one, so the switch lands at a known request. Bandit +// calls are serialised by the cache, so done needs no lock. +type switchOnceTo struct { + target ascache.PolicyType + done bool +} + +func (b *switchOnceTo) RecordStats(ascache.ShadowStats) {} + +func (b *switchOnceTo) SelectPolicy() ascache.PolicyType { + if b.done { + return ascache.Undefined + } + b.done = true + + return b.target +} + +// warmupWindows are the reported request ranges, as offsets from the switch. +var warmupWindows = [][2]int{{0, 1000}, {1000, 5000}, {5000, 20000}, {20000, 100000}} + +type warmupRow struct { + name string + hits []int +} + +// TestSwitchWarmupCost measures what a switch costs in the requests right +// after it, per migration strategy. +// +// A switch is only worth making if the policy switched to earns back what the +// switch loses. The loss is concentrated in the requests immediately after it, +// where a cold policy starts empty, a warm one starts with the outgoing +// policy's contents, and a gradual one fills as keys are asked for. An average +// over a whole run hides that shape entirely, so hits are reported in windows +// measured from the switch. +// +// The switch is scripted rather than chosen, so every strategy switches at the +// same request, from LRU to LFU, on the evidence suite's zipf workload and +// cache size. Two references bracket the result: LFU replayed on its own over +// the whole trace, which is an LFU that never had to warm up, and LRU replayed +// on its own, which is not switching at all. +func TestSwitchWarmupCost(t *testing.T) { + if testing.Short() { + t.Skip("evidence run; use make evidence") + } + + const ( + capacity = 500 + switchAt = 100000 + ) + w := bench.Zipf(2*switchAt, 20000, 1.1, 1) + + lruBuilder, lfuBuilder := fixedPolicy(t, "LRU"), fixedPolicy(t, "LFU") + + rows := make([]warmupRow, 0, 8) + for _, ref := range []struct { + name string + builder bench.PolicyBuilder + }{ + {"LFU all along (no warm-up)", lfuBuilder}, + {"LRU, never switched", lruBuilder}, + } { + policy, err := ref.builder.Build(capacity) + require.NoError(t, err) + rows = append(rows, warmupRow{ref.name, windowHits(policy, w, switchAt)}) + } + + for _, cfg := range []struct { + name string + strategy ascache.MigrationStrategy + maxRequests int64 + }{ + {"cold", ascache.MigrationCold, 0}, + {"warm", ascache.MigrationWarm, 0}, + {"gradual", ascache.MigrationGradual, 0}, + {"gradual, capped at 10 Gets", ascache.MigrationGradual, 10}, + {"gradual, capped at 100 Gets", ascache.MigrationGradual, 100}, + {"gradual, capped at 1000 Gets", ascache.MigrationGradual, 1000}, + } { + lru, err := lruBuilder.Build(capacity) + require.NoError(t, err) + lfu, err := lfuBuilder.Build(capacity) + require.NoError(t, err) + + cache, err := ascache.NewAdaptiveCache( + []ascache.Policy[string, int]{lru, lfu}, + &switchOnceTo{target: ascache.LFU}, + &ascache.Settings{ + // The Get completing request switchAt ends the first epoch, so + // request index switchAt is the first one the new policy serves. + EpochRequests: switchAt, + EvictPartialCapacityFilling: true, + MigrationStrategy: cfg.strategy, + MigrationMaxRequests: cfg.maxRequests, + }, + ) + require.NoError(t, err) + + hits := windowHits(cache, w, switchAt) + require.Equal(t, ascache.LFU, cache.ActivePolicy(), "%s: the scripted switch must have happened", cfg.name) + require.NoError(t, cache.Close()) + + rows = append(rows, warmupRow{cfg.name, hits}) + } + + t.Logf("\nhit rate after a switch from LRU to LFU at request %d (%s, %d requests, cache %d)\n%s", + switchAt, w.Name, len(w.Keys), capacity, warmupTable(rows)) +} + +func fixedPolicy(t *testing.T, name string) bench.PolicyBuilder { + t.Helper() + + for _, builder := range bench.FixedPolicies() { + if builder.Name == name { + return builder + } + } + require.FailNow(t, "no fixed policy named "+name) + + return bench.PolicyBuilder{} +} + +// windowHits replays w through c as a read-through cache, filling on every miss +// as bench.Replay does, and counts the hits falling in each of warmupWindows. +func windowHits(c bench.Cache, w bench.Workload, switchAt int) []int { + hits := make([]int, len(warmupWindows)) + + for i, key := range w.Keys { + _, ok := c.Get(key) + if !ok { + c.Add(key, i) + + continue + } + if i < switchAt { + continue + } + + offset := i - switchAt + for j, window := range warmupWindows { + if offset >= window[0] && offset < window[1] { + hits[j]++ + } + } + } + + return hits +} + +func warmupTable(rows []warmupRow) string { + var b strings.Builder + + b.WriteString("| Configuration |") + for _, window := range warmupWindows { + fmt.Fprintf(&b, " %d-%d |", window[0], window[1]) + } + b.WriteString("\n| --- |" + strings.Repeat(" --- |", len(warmupWindows)) + "\n") + + for _, row := range rows { + fmt.Fprintf(&b, "| %s |", row.name) + for j, window := range warmupWindows { + fmt.Fprintf(&b, " %.2f%% |", 100*float64(row.hits[j])/float64(window[1]-window[0])) + } + b.WriteString("\n") + } + + return b.String() +} diff --git a/cache.go b/cache.go index 3d6f498..300af38 100644 --- a/cache.go +++ b/cache.go @@ -54,6 +54,10 @@ type AdaptiveCache[K comparable, V any] struct { migrateFrom PolicyType migrationKeys []K migrationRealKeys map[K]struct{} + // migrationRequests counts Gets served while the current gradual window + // is open, against Settings.MigrationMaxRequests. It needs no atomic: + // every Get in a window already holds the write lock. + migrationRequests int64 // --- Control Plane --- @@ -110,8 +114,9 @@ type AdaptiveCache[K comparable, V any] struct { // epochTicker is nil when the cache ends its epochs on request count // alone, since time.NewTicker rejects a non-positive duration. epochTicker *time.Ticker - // epochRequests counts Get calls since the last request-driven epoch. It - // is mutated on the read path, so it must be atomic. + // epochRequests counts every Get since construction; an epoch runs on each + // multiple of Settings.EpochRequests (see countRequest). It is mutated on + // the read path, so it must be atomic. epochRequests atomic.Int64 settings *Settings @@ -124,6 +129,13 @@ type AdaptiveCache[K comparable, V any] struct { // recordActiveSample counts the active policy's result for a key that is part // of the measured sample. Unsampled keys are served normally but not counted, // so the active arm's evidence covers the same substream as every shadow's. +// +// On the read-lock path it runs after the lock is released, so it is not +// atomic with the Get it records. An epoch collecting in that gap reports the +// sample one epoch late; a switch landing in it credits the sample to the +// policy just made active. switchLocked clears the counters, but it cannot +// reach a Get already past its lookup, so what remains is bounded by the Gets +// in flight at that instant, not by how long the bandit took to decide. func (c *AdaptiveCache[K, V]) recordActiveSample(sampled, hit bool) { if !sampled { return @@ -185,6 +197,16 @@ func (c *AdaptiveCache[K, V]) get(key K) (V, bool) { val, found := c.policies[c.activePolicy].Get(key) c.recordActiveSample(sampled, found) + // A window still open after this Get counts it against the cap. Closing + // demotes the source, which is safe here: the value above came from the + // active policy, and the source is not the active policy. + if c.migrating { + c.migrationRequests++ + if limit := c.settings.MigrationMaxRequests; limit > 0 && c.migrationRequests >= limit { + c.closeMigrationLocked() + } + } + return val, found } diff --git a/docs/configuration.md b/docs/configuration.md index 218b8a9..db4146f 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -21,6 +21,10 @@ type Settings struct { // Default: MigrationCold. MigrationStrategy MigrationStrategy + // MigrationMaxRequests caps a MigrationGradual window at N Get calls. + // Zero sets no cap. See "Migration Strategies". + MigrationMaxRequests int64 + // ObserveOnly measures every arm without ever switching. // See docs/advisor-mode.md. ObserveOnly bool @@ -44,7 +48,21 @@ type Settings struct { | --- | --- | --- | | `MigrationCold` (default) | New active policy starts empty | Simple; causes a temporary miss spike | | `MigrationWarm` | All key/value pairs copied at switch time | No miss spike; O(n) work at switch | -| `MigrationGradual` | Keys promoted on Get; one key drained per Add | Spreads migration cost; window closes at the next epoch at the latest | +| `MigrationGradual` | Keys promoted on Get; one key drained per Add | Spreads migration cost; every `Get` takes the write lock while the window is open, which closes at the next epoch at the latest | + +A gradual window serialises reads for as long as it is open. On a long epoch +that can be most of the epoch, so `MigrationMaxRequests` caps it at a number of +`Get` calls: when the cap is reached the window closes and the old policy is +demoted, and any key not promoted by then is gone — a later `Get` for it is a +miss, the same as it would have been under `MigrationCold`. Only `Get` counts, +because `Get` is what takes the write lock; `Add` drains a key per call and +shortens the window anyway. Zero, the default, sets no cap. + +Measured on zipf with a switch from LRU to LFU, a cap of 100 cost 6.5 points of +hit rate over the first 1,000 requests after the switch, a cap of 10 made gradual +behave like cold, and a cap of 1,000 was never reached because read-through +traffic had already drained the window. [Evidence](evidence.md#what-does-a-switch-cost-right-after-it) +has the table, including what each strategy costs. ## Reducing shadow overhead diff --git a/docs/evidence.md b/docs/evidence.md index ab7731c..da01a87 100644 --- a/docs/evidence.md +++ b/docs/evidence.md @@ -407,3 +407,43 @@ Higher rates cost more and buy no better ranking here, so 0.05 is a reasonable default. Raise it if your keyspace is small enough that 5% of it is only a handful of keys -- `MinShadowCapacity` guards the degenerate end by raising the effective rate rather than letting a miniature shrink into noise. + +## What does a switch cost right after it? + +Every figure above is a whole-run average, which hides where a switch's cost +falls: in the requests immediately after it. `TestSwitchWarmupCost` scripts one +switch, from LRU to LFU at request 100,000 of the zipf workload above (200,000 +requests, cache 500), and reports hit rate in windows measured from the switch. +The switch is forced rather than chosen, so every strategy switches at the same +request. The run is deterministic: two runs produce identical tables. + +| Configuration | 0-1000 | 1000-5000 | 5000-20000 | 20000-100000 | +| --- | --- | --- | --- | --- | +| LFU all along (no warm-up) | 77.00% | 74.28% | 74.59% | 73.90% | +| LRU, never switched | 70.50% | 67.90% | 67.58% | 66.82% | +| cold | 59.50% | 69.45% | 72.93% | 73.52% | +| warm | 70.50% | 70.03% | 72.91% | 73.55% | +| gradual | 70.40% | 70.12% | 72.88% | 73.52% | +| gradual, capped at 10 Gets | 59.90% | 69.47% | 72.93% | 73.52% | +| gradual, capped at 100 Gets | 63.90% | 69.58% | 72.99% | 73.53% | +| gradual, capped at 1000 Gets | 70.40% | 70.12% | 72.88% | 73.52% | + +- **Cold pays for the switch up front.** Its first 1,000 requests serve 59.50%, + 11.0 points below not switching at all and 17.5 below an LFU that never had to + warm up. By the next window it is already ahead of not switching. +- **Warm and gradual show no dip.** Both serve the first 1,000 requests within + 0.1 points of LRU's own rate, because the entries LRU held are there to be hit. +- **None of them becomes the LFU that was there all along.** From 20,000 to + 100,000 requests after the switch every strategy sits at 73.52-73.55%, against + 73.90%. Moving the contents does not move the access history an LFU running + from the start would have built, and the gap is what that history was worth + here. +- **The window cap costs only when it bites.** At 10 Gets the gradual window + closes almost immediately and the first 1,000 requests look like cold + (59.90%); at 100 they lose 6.5 points against uncapped gradual. At 1,000 the + row is identical to uncapped: under read-through traffic every miss is an + `Add`, every `Add` drains a pending key, and the window emptied on its own + somewhere between 100 and 1,000 Gets. + +Reproduce with `cd bench && go test -run TestSwitchWarmupCost -v .`, or +`make evidence`. diff --git a/docs/policies.md b/docs/policies.md index 6f9d830..7377492 100644 --- a/docs/policies.md +++ b/docs/policies.md @@ -29,7 +29,7 @@ cache, err := ascache.NewAdaptiveCache( | LFU | `policies.NewLFU` | this repository's O(1) LFU; strong on stationary popularity, weak when it shifts | | 2Q | `policies.NewTwoQueue` | scan-resistant; a scan cannot flush the working set | | Random | `policies.NewRandomPolicy` | no bookkeeping; the control arm worth beating | -| TTL | `policies.NewTTL` | expiry as well as recency | +| TTL | `policies.NewTTL` | expiry as well as recency; expiry runs on the wall clock, so its hit rate depends on how fast traffic arrives, and a replay is reproducible only while the TTL is far longer than the run | | ARC | `policies/arc.NewPolicy` | separate module — see below | | W-TinyLFU | `policies/tinylfu.NewPolicy` | separate module; the strongest baseline | | S3-FIFO | `policies/fifo.NewS3FIFOPolicy` | separate module; three FIFO queues, and deterministic | diff --git a/epoch.go b/epoch.go index a49e60f..dda79eb 100644 --- a/epoch.go +++ b/epoch.go @@ -28,20 +28,25 @@ func (c *AdaptiveCache[K, V]) runAdaptiveSelect() { // the call that completes it. Caller must hold no lock: runEpoch takes the // write lock. // -// Exactly one caller per epoch sees the count equal the limit, so exactly one -// epoch runs however many goroutines are in Get. The limit is subtracted -// rather than the counter reset, so requests arriving mid-crossing still -// count towards the next epoch. +// The counter only ever increases, and an epoch runs on the call whose +// increment returns a multiple of the limit. Add hands every caller a distinct +// value, so each multiple is seen by exactly one caller: one epoch per limit +// requests however many goroutines are in Get, and a request arriving +// mid-crossing counts toward the next epoch. +// +// Do not reintroduce "compare with the limit, then subtract it". Callers that +// increment between one caller's comparison and its subtraction push the count +// past the limit unobserved, nothing subtracts again, and request-driven +// epochs stop for good. func (c *AdaptiveCache[K, V]) countRequest() { limit := c.settings.EpochRequests if limit <= 0 { return } - if c.epochRequests.Add(1) != limit { + if c.epochRequests.Add(1)%limit != 0 { return } - c.epochRequests.Add(-limit) c.runEpoch() } diff --git a/epoch_attribution_test.go b/epoch_attribution_test.go new file mode 100644 index 0000000..865646e --- /dev/null +++ b/epoch_attribution_test.go @@ -0,0 +1,131 @@ +package ascache + +import ( + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// parkedSwitchBandit parks inside its first SelectPolicy until released and +// then asks for a switch to next; every later selection asks for no change. +// It keeps every report it is given. +type parkedSwitchBandit struct { + entered chan struct{} + release chan struct{} + next PolicyType + + mu sync.Mutex + calls int + reports []ShadowStats +} + +func (b *parkedSwitchBandit) RecordStats(s ShadowStats) { + b.mu.Lock() + defer b.mu.Unlock() + + b.reports = append(b.reports, s) +} + +func (b *parkedSwitchBandit) SelectPolicy() PolicyType { + b.mu.Lock() + b.calls++ + first := b.calls == 1 + b.mu.Unlock() + + if !first { + return Undefined + } + close(b.entered) + <-b.release + + return b.next +} + +func (b *parkedSwitchBandit) takeReports() []ShadowStats { + b.mu.Lock() + defer b.mu.Unlock() + + out := b.reports + b.reports = nil + + return out +} + +// TestEpoch_OutgoingSamplesAreNotCreditedToIncomingPolicy pins who a sample +// belongs to across a switch. +// +// The active arm's evidence is counted on the cache, not on the policy, and +// read-and-reset when an epoch collects. The epoch then releases the lock while +// the bandit decides, and Gets keep arriving, served by the policy that is +// still active and counted into those same cache-level counters. If the bandit +// then switches, the next epoch reports whatever accumulated as the evidence +// of the policy that has just become active -- requests it never served. +// The window is as long as the bandit takes to decide. +func TestEpoch_OutgoingSamplesAreNotCreditedToIncomingPolicy(t *testing.T) { + lru := newMockPolicy[string, int](LRU, 10) + lfu := newMockPolicy[string, int](LFU, 10) + bandit := &parkedSwitchBandit{ + entered: make(chan struct{}), + release: make(chan struct{}), + next: LFU, + } + + ac, err := NewAdaptiveCache([]Policy[string, int]{lru, lfu}, bandit, &Settings{ + EpochDuration: 24 * time.Hour, + EvictPartialCapacityFilling: true, + }) + require.NoError(t, err) + t.Cleanup(func() { _ = ac.Close() }) + + ac.Add("a", 1) + + epochOne := make(chan struct{}) + go func() { + ac.runEpoch() + close(epochOne) + }() + + select { + case <-bandit.entered: + case <-time.After(5 * time.Second): + require.FailNow(t, "epoch one never reached the bandit") + } + + // Epoch one has collected and released the lock. LRU is still active and + // serves these; each is recorded as a sample of the active arm. + const servedByLRU = 5 + for range servedByLRU { + _, ok := ac.Get("a") + require.True(t, ok) + } + + close(bandit.release) + select { + case <-epochOne: + case <-time.After(5 * time.Second): + require.FailNow(t, "epoch one never finished") + } + require.Equal(t, LFU, ac.ActivePolicy(), "epoch one's selection switches to LFU") + + bandit.takeReports() // epoch one's, delivered before the bandit parked + + // Epoch two. Not one request has reached LFU since it became active. + ac.runEpoch() + + reports := bandit.takeReports() + var incoming *ShadowStats + for i := range reports { + if reports[i].Policy == LFU { + incoming = &reports[i] + } + } + require.NotNil(t, incoming, "epoch two must report LFU") + + assert.Zero(t, incoming.Hits+incoming.Misses, + "LFU served nothing since becoming active, yet epoch two reported %d hits and %d misses "+ + "for it: requests LRU served while the bandit was deciding were credited to LFU", + incoming.Hits, incoming.Misses) +} diff --git a/epoch_requests_liveness_test.go b/epoch_requests_liveness_test.go new file mode 100644 index 0000000..af8c4e1 --- /dev/null +++ b/epoch_requests_liveness_test.go @@ -0,0 +1,78 @@ +package ascache + +import ( + "sync" + "sync/atomic" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// tickCountingBandit counts selections, one per reporting epoch, so a test +// can tell whether the epoch clock is still advancing. +type tickCountingBandit struct { + selections atomic.Int64 +} + +func (b *tickCountingBandit) RecordStats(ShadowStats) {} + +func (b *tickCountingBandit) SelectPolicy() PolicyType { + b.selections.Add(1) + + return LRU +} + +// TestEpochRequests_ClockSurvivesContention checks that the request-driven +// epoch clock is still running after a contended burst. +// +// countRequest triggers an epoch when its counter equals the limit, then +// subtracts the limit. If enough other callers increment between one caller's +// increment and its subtraction, the counter passes the limit without anyone +// observing it equal, and nothing ever subtracts again: every later Get sees +// a count above the limit and returns. The clock does not skip an epoch; it +// stops. With EpochRequests at 1 the window needs only one concurrent Get. +func TestEpochRequests_ClockSurvivesContention(t *testing.T) { + t.Parallel() + + bandit := &tickCountingBandit{} + cache, err := NewAdaptiveCache[string, int]( + []Policy[string, int]{ + newEvictingPolicy[string, int](LRU, 4), + newEvictingPolicy[string, int](LFU, 4), + }, + bandit, + &Settings{EpochRequests: 1, EvictPartialCapacityFilling: true}, + ) + require.NoError(t, err) + defer cache.Close() + + const ( + goroutines = 16 + perRoutine = 2000 + ) + + var wg sync.WaitGroup + for range goroutines { + wg.Add(1) + go func() { + defer wg.Done() + for range perRoutine { + cache.Get("k") + } + }() + } + wg.Wait() + + // The burst is over and nothing else is running. With EpochRequests at 1, + // each of these Gets must end exactly one epoch. + const serial = 10 + before := bandit.selections.Load() + for range serial { + cache.Get("k") + } + + assert.Equal(t, before+serial, bandit.selections.Load(), + "the request-driven epoch clock stopped after contention: %d selections before, "+ + "%d after %d uncontended Gets", before, bandit.selections.Load(), serial) +} diff --git a/errors.go b/errors.go index 3cfa1f3..386c85f 100644 --- a/errors.go +++ b/errors.go @@ -30,3 +30,7 @@ var ErrInvalidEpochDuration = errors.New( // ErrInvalidEpochRequests is returned by NewAdaptiveCache when // Settings.EpochRequests is negative. var ErrInvalidEpochRequests = errors.New("epoch requests must not be negative") + +// ErrInvalidMigrationMaxRequests is returned by NewAdaptiveCache when +// Settings.MigrationMaxRequests is negative. +var ErrInvalidMigrationMaxRequests = errors.New("migration max requests must not be negative") diff --git a/evicting_test.go b/evicting_test.go index a603122..01da420 100644 --- a/evicting_test.go +++ b/evicting_test.go @@ -370,6 +370,33 @@ func TestGradualMigration_NeverServesAZeroFromTheSource(t *testing.T) { cache.Get("fresh" + strconv.Itoa(i)) } + // Peek and Contains must not surface a placeholder either. They are checked + // here, while the window is still open: the Gets below promote the real + // values, after which nothing about the window is observable any more. + // + // The fresh keys carry the discriminating half. Nobody stored them, so the + // only copies anywhere are zero placeholders written by the fan-out, and a + // Peek or Contains that ever reached past the active policy into one would + // report a key present that was never written. + cache.mu.RLock() + stillOpen := cache.migrating + cache.mu.RUnlock() + require.True(t, stillOpen, "the window must still be open for Peek and Contains to be tested against it") + + for i := 1; i <= capacity; i++ { + key := "real" + strconv.Itoa(i) + if value, found := cache.Peek(key); found { + assert.NotZero(t, value, + "Peek(%q) returned a zero as a hit: a shadow placeholder surfaced without a Get", key) + } + } + for i := range capacity * 3 { + key := "fresh" + strconv.Itoa(i) + value, found := cache.Peek(key) + assert.False(t, found, "Peek(%q) reported a key nobody stored (value %d)", key, value) + assert.False(t, cache.Contains(key), "Contains(%q) reported a key nobody stored", key) + } + for i := 1; i <= capacity; i++ { key := "real" + strconv.Itoa(i) value, found := cache.Get(key) diff --git a/migration.go b/migration.go index 9458255..521b8e7 100644 --- a/migration.go +++ b/migration.go @@ -70,6 +70,7 @@ func (c *AdaptiveCache[K, V]) clearMigrationState() { c.migrateFrom = Undefined c.migrationKeys = nil c.migrationRealKeys = nil + c.migrationRequests = 0 } // drainOneKey migrates one pending key from the migration source policy into diff --git a/migration_limit_test.go b/migration_limit_test.go new file mode 100644 index 0000000..922b0c9 --- /dev/null +++ b/migration_limit_test.go @@ -0,0 +1,118 @@ +package ascache + +import ( + "errors" + "strconv" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// makeLimitedGradualCache builds a two-policy cache under MigrationGradual with +// the given window cap, holding pending real values k1..k5 (stored as 1..5, so +// a zero read back is unambiguously a shadow placeholder), and opens a window +// by switching from LRU to LFU. +func makeLimitedGradualCache(t *testing.T, maxRequests int64) *AdaptiveCache[string, int] { + t.Helper() + + ac, err := NewAdaptiveCache( + []Policy[string, int]{ + newMockPolicy[string, int](LRU, 10), + newMockPolicy[string, int](LFU, 10), + }, + &mockBandit{next: LRU}, + &Settings{ + EpochDuration: 24 * time.Hour, + EvictPartialCapacityFilling: true, + MigrationStrategy: MigrationGradual, + MigrationMaxRequests: maxRequests, + }, + ) + require.NoError(t, err) + t.Cleanup(func() { _ = ac.Close() }) + + for i := 1; i <= 5; i++ { + ac.Add("k"+strconv.Itoa(i), i) + } + triggerSwitch(ac, LFU) + require.True(t, windowOpen(ac), "switching away from a non-empty policy must open a window") + + return ac +} + +func windowOpen(ac *AdaptiveCache[string, int]) bool { + ac.mu.RLock() + defer ac.mu.RUnlock() + + return ac.migrating +} + +// TestMigrationGradual_MaxRequestsClosesWindow pins the cap. Gets for keys +// nobody stored promote nothing, so without a cap the window would stay open +// until the next epoch and every one of those Gets would take the write lock. +// With a cap of 3 the third Get closes it, and Gets go back to the read-lock +// path. +func TestMigrationGradual_MaxRequestsClosesWindow(t *testing.T) { + ac := makeLimitedGradualCache(t, 3) + + ac.Get("absent-1") + ac.Get("absent-2") + require.True(t, windowOpen(ac), "the window must stay open below the cap") + + ac.Get("absent-3") + assert.False(t, windowOpen(ac), "the Get reaching the cap must close the window") + + ac.mu.RLock() + source := ac.migrateFrom + ac.mu.RUnlock() + assert.Equal(t, Undefined, source, "a closed window leaves no migration source") +} + +// TestMigrationGradual_MaxRequestsAbandonsPendingKeys pins what the cap costs. +// Closing the window demotes the source, rewriting it to zero values, so a key +// that was never promoted is gone. The one thing that must not happen is the +// central invariant breaking on the way out: a Get for an abandoned key is a +// miss, never a zero returned as a hit. +func TestMigrationGradual_MaxRequestsAbandonsPendingKeys(t *testing.T) { + ac := makeLimitedGradualCache(t, 2) + + ac.Get("absent-1") + ac.Get("absent-2") + require.False(t, windowOpen(ac), "the cap of 2 is reached") + + for i := 1; i <= 5; i++ { + key := "k" + strconv.Itoa(i) + value, found := ac.Get(key) + if found { + assert.NotZero(t, value, "Get(%q) returned a shadow zero as a hit", key) + } + assert.False(t, found, "Get(%q): a key abandoned by the capped window must be a miss", key) + } +} + +// TestMigrationGradual_ZeroMaxRequestsSetsNoCap pins backward compatibility: +// the zero value behaves exactly as the cache did before the setting existed, +// leaving the window open however many Gets arrive, until something else +// closes it. +func TestMigrationGradual_ZeroMaxRequestsSetsNoCap(t *testing.T) { + ac := makeLimitedGradualCache(t, 0) + + for i := range 1000 { + ac.Get("absent-" + strconv.Itoa(i)) + } + + assert.True(t, windowOpen(ac), "with no cap the window stays open until the epoch ends it") +} + +func TestNewAdaptiveCache_RejectsNegativeMigrationMaxRequests(t *testing.T) { + _, err := NewAdaptiveCache( + []Policy[string, int]{newMockPolicy[string, int](LRU, 10)}, + &mockBandit{next: LRU}, + &Settings{EpochDuration: time.Hour, MigrationMaxRequests: -1}, + ) + + assert.True(t, errors.Is(err, ErrInvalidMigrationMaxRequests), + "a negative cap has no meaning and must be rejected, got %v", err) +} diff --git a/policies/adapt.go b/policies/adapt.go index e33a38f..109a41f 100644 --- a/policies/adapt.go +++ b/policies/adapt.go @@ -87,6 +87,12 @@ func (c *AdaptedCache[K, V]) Add(key K, value V) bool { // still in use at its original capacity while the configured size is smaller, // and without this the cache would hold far more than the caller asked for // while Cap reported the smaller number. +// +// It removes keys[0], which is some entry rather than the oldest one: for 2Q +// the oldest frequent entry, for ARC the oldest recent one (see Keys). That is +// acceptable for the reason Resize gives -- which entries survive a shrink is +// not meaningful -- and the TestKeysOrder_ canaries make an upstream change to +// either order fail a test rather than silently change what this removes. func (c *AdaptedCache[K, V]) enforceCapacityLocked() { for c.cache.Len() > c.size { keys := c.cache.Keys() @@ -140,7 +146,11 @@ func (c *AdaptedCache[K, V]) Purge() { c.cache.Purge() } -// Keys returns the cached keys, oldest first. +// Keys returns the cached keys in the wrapped cache's own order. That is not a +// recency order for the caches this adapter serves: 2Q returns its frequent +// list then its recent one, ARC its recent list then its frequent one, each +// oldest first. TestKeysOrder_TwoQueueIsFrequentThenRecent and policies/arc's +// TestKeysOrder_ARCIsRecentThenFrequent pin both. func (c *AdaptedCache[K, V]) Keys() []K { c.mu.RLock() defer c.mu.RUnlock() diff --git a/policies/arc/keys_order_test.go b/policies/arc/keys_order_test.go new file mode 100644 index 0000000..ccf9c2e --- /dev/null +++ b/policies/arc/keys_order_test.go @@ -0,0 +1,31 @@ +package arc_test + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/sshaplygin/as-cache/policies/arc" +) + +// TestKeysOrder_ARCIsRecentThenFrequent pins ARC's grouping: its recent list +// (T1) first, then its frequent list (T2), each oldest first. That is the +// opposite grouping to 2Q's, which docs/policies.md relies on to explain why +// an adapted cache's Resize replays every entry instead of keeping a head or a +// tail. The 2Q and LRU canaries are in the policies module, keys_order_test.go. +func TestKeysOrder_ARCIsRecentThenFrequent(t *testing.T) { + p, err := arc.NewPolicy[string, int](8) + require.NoError(t, err) + + p.Add("a", 1) + p.Add("b", 2) + p.Add("c", 3) + require.Equal(t, []string{"a", "b", "c"}, p.Keys(), "all three are in the recent list") + + _, ok := p.Get("b") + require.True(t, ok) + + assert.Equal(t, []string{"a", "c", "b"}, p.Keys(), + "a second access moves b to the frequent list, which ARC returns last") +} diff --git a/policies/keys_order_test.go b/policies/keys_order_test.go new file mode 100644 index 0000000..10524ac --- /dev/null +++ b/policies/keys_order_test.go @@ -0,0 +1,62 @@ +package policies_test + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/sshaplygin/as-cache/policies" +) + +// These pin the order in which upstream hashicorp/golang-lru returns Keys(). +// Nothing in golang-lru promises it, and code here depends on it: +// +// - demoteLocked rewrites a demoted policy to zero values by walking Keys(). +// For an LRU that walk re-establishes the same recency order only because +// Keys() runs oldest to newest; were it newest first, every demotion would +// invert the order the policy had learned. +// - docs/policies.md states that 2Q and ARC return opposite groupings, which +// is why AdaptedCache.Resize replays every entry rather than keeping a head +// or a tail. If upstream changes either, that page is wrong. +// +// Each test asserts an exact sequence after one access that moves an entry, so +// a changed order fails rather than passing by coincidence on insertion order. +// ARC lives in its own module; its canary is policies/arc's +// TestKeysOrder_ARCIsRecentThenFrequent. + +// TestKeysOrder_LRUIsOldestToNewest pins the order demotion relies on. +func TestKeysOrder_LRUIsOldestToNewest(t *testing.T) { + p, err := policies.NewLRU[string, int](8) + require.NoError(t, err) + + p.Add("a", 1) + p.Add("b", 2) + p.Add("c", 3) + require.Equal(t, []string{"a", "b", "c"}, p.Keys(), "insertion order, oldest first") + + _, ok := p.Get("a") + require.True(t, ok) + + assert.Equal(t, []string{"b", "c", "a"}, p.Keys(), + "an access moves the key to the newest end; demoteLocked's recency rewrite depends on this") +} + +// TestKeysOrder_TwoQueueIsFrequentThenRecent pins 2Q's grouping: its frequent +// list first, then its recent list, each oldest first. It is not a recency +// order -- the key accessed last comes out first. +func TestKeysOrder_TwoQueueIsFrequentThenRecent(t *testing.T) { + p, err := policies.NewTwoQueue[string, int](8) + require.NoError(t, err) + + p.Add("a", 1) + p.Add("b", 2) + p.Add("c", 3) + require.Equal(t, []string{"a", "b", "c"}, p.Keys(), "all three are in the recent list") + + _, ok := p.Get("b") + require.True(t, ok) + + assert.Equal(t, []string{"b", "a", "c"}, p.Keys(), + "a second access promotes b to the frequent list, which 2Q returns first") +} diff --git a/settings.go b/settings.go index 227bff8..1208483 100644 --- a/settings.go +++ b/settings.go @@ -29,6 +29,14 @@ type Settings struct { // changes. Defaults to MigrationCold (zero value). MigrationStrategy MigrationStrategy + // MigrationMaxRequests caps a MigrationGradual window at this many Get + // calls. While a window is open every Get takes the write lock, so the cap + // bounds how long reads stay serialised. When it is reached the window + // closes and the source is demoted: keys not yet promoted are abandoned, + // and a later Get for one is a miss. Zero sets no cap, and the window + // closes at the next epoch boundary. The other strategies ignore it. + MigrationMaxRequests int64 + // MinHitRateImprovement is the hit-rate advantage, as an absolute // difference in [0,1], that the bandit's selection must hold over the // active policy in the epoch just measured before the switch is applied. @@ -101,6 +109,9 @@ func NewAdaptiveCache[K comparable, V any]( if settings.EpochRequests < 0 { return nil, fmt.Errorf("%w: got %d", ErrInvalidEpochRequests, settings.EpochRequests) } + if settings.MigrationMaxRequests < 0 { + return nil, fmt.Errorf("%w: got %d", ErrInvalidMigrationMaxRequests, settings.MigrationMaxRequests) + } // An epoch has to be ended by something. Either clock is acceptable and // both together are fine; neither leaves a cache that measures every // policy forever and never acts on any of it. diff --git a/shadow.go b/shadow.go index 0bfbd4d..3edab9b 100644 --- a/shadow.go +++ b/shadow.go @@ -35,8 +35,10 @@ func (c *AdaptiveCache[K, V]) fanOutReadLocked(key K) { // // The rewrite walks Keys() so a recency policy re-establishes the same order // and a frequency policy gains one access on every surviving key, leaving the -// relative order intact. That reasoning does not hold for the FIFO-queue -// policies; docs/policies.md records what demotion costs them. +// relative order intact. For an LRU that depends on Keys() running oldest to +// newest, which golang-lru does not promise; policies' +// TestKeysOrder_LRUIsOldestToNewest pins it. The reasoning does not hold for +// the FIFO-queue policies; docs/policies.md records what demotion costs them. // // Caller must hold the write lock, and must have published the new state // first, so no reader can observe a value being dropped. @@ -106,6 +108,14 @@ func (c *AdaptiveCache[K, V]) switchLocked(from, to PolicyType) { delete(c.tenureStats, from) delete(c.tenureStats, to) + // The active arm's samples are counted on the cache rather than on a + // policy, and everything counted since the last collection was served by + // from -- including every Get that arrived while the bandit was deciding. + // Left in place, the next epoch would report it as to's evidence. It is + // dropped, as demotion drops from's own counters. + c.activeSampledHits.Store(0) + c.activeSampledMisses.Store(0) + if !c.migrating { c.demoteLocked(from) } diff --git a/stability.go b/stability.go index 8adbfc1..8e8e99b 100644 --- a/stability.go +++ b/stability.go @@ -53,9 +53,20 @@ func (c *AdaptiveCache[K, V]) allowSwitchLocked(candidate PolicyType) bool { return false } - if c.settings.MinHitRateImprovement > 0 && - hitRate(cand)-hitRate(active) < c.settings.MinHitRateImprovement { - return false + if c.settings.MinHitRateImprovement > 0 { + // hitRate reports 0 for an arm that saw no requests, which reads as an + // arm serving nothing: against an active policy with no traffic any + // candidate with one hit clears the threshold, so the switch would be + // made on no evidence about the policy being replaced. A candidate with + // no traffic already loses on arithmetic; it is rejected here too so + // that stays true whatever hitRate returns for an empty epoch. + if active.Hits+active.Misses == 0 || cand.Hits+cand.Misses == 0 { + return false + } + + if hitRate(cand)-hitRate(active) < c.settings.MinHitRateImprovement { + return false + } } return true diff --git a/stability_test.go b/stability_test.go new file mode 100644 index 0000000..114e72f --- /dev/null +++ b/stability_test.go @@ -0,0 +1,39 @@ +package ascache + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +// TestSwitchStability_ZeroTrafficOnActiveBlocksSwitch pins the edge the +// improvement gate used to get wrong. hitRate reports 0 for an arm that saw no +// requests, so an active policy with no traffic in the measured epoch looked +// like one serving nothing, and any candidate with a single hit cleared the +// threshold against it. That is a switch made on no evidence about the policy +// being replaced -- exactly what MinHitRateImprovement exists to prevent. +func TestSwitchStability_ZeroTrafficOnActiveBlocksSwitch(t *testing.T) { + ac, _, lfu := makeStabilityCache(t, &Settings{MinHitRateImprovement: 0.02}) + + primeActiveStats(ac, 0, 0) + primeStats(lfu, 5, 0) + ac.runEpoch() + + assert.Equal(t, LRU, ac.ActivePolicy(), + "an active policy with no traffic cannot be out-performed; there is nothing to compare") +} + +// TestSwitchStability_ZeroTrafficOnCandidateBlocksSwitch is the other edge. +// It already held before the fix -- a candidate's empty epoch scores 0 and +// loses the comparison -- and is pinned so the fix for the active edge cannot +// quietly turn a candidate with no evidence into a winner. +func TestSwitchStability_ZeroTrafficOnCandidateBlocksSwitch(t *testing.T) { + ac, _, lfu := makeStabilityCache(t, &Settings{MinHitRateImprovement: 0.02}) + + primeActiveStats(ac, 5, 5) + primeStats(lfu, 0, 0) + ac.runEpoch() + + assert.Equal(t, LRU, ac.ActivePolicy(), + "a candidate with no traffic has shown nothing and must not win") +}