From 798f9f0d7db34145abe411c02580b2a35d53ed80 Mon Sep 17 00:00:00 2001 From: Sam Shaplygin Date: Fri, 28 Aug 2026 00:00:39 +0200 Subject: [PATCH 1/2] Add the failing tests for a bandit called under the write lock Three properties the epoch does not have yet. All three fail against this commit's parent, which is the point of committing them first. SlowBanditDoesNotBlockReaders a bandit parked in SelectPolicy must not stall Get. Fails: "Get blocked behind the bandit: a slow bandit is a cache outage". StaleSelectionIsDiscarded a decision superseded by a later epoch must be dropped. Fails by deadlock -- the second epoch cannot start while the first holds the write lock, so the package times out. TickerAndRequestEpochsDoNotLose every counted request reaches the bandit or Counts is still live in its arm, under both epoch clocks at once. Go's RWMutex queues new readers behind a waiting writer, so a bandit called while the cache's write lock is held stalls every Get in the process for as long as it runs: a store timeout becomes a cache outage. The godoc on Bandit currently manages this by forbidding it -- "an implementation must not block" -- which pushes the problem onto every implementer and is why the distributed bandit had to buffer and serve selections from an atomic. --- epoch_lock_test.go | 274 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 274 insertions(+) create mode 100644 epoch_lock_test.go diff --git a/epoch_lock_test.go b/epoch_lock_test.go new file mode 100644 index 0000000..4149ab4 --- /dev/null +++ b/epoch_lock_test.go @@ -0,0 +1,274 @@ +package ascache + +import ( + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// slowBandit blocks inside SelectPolicy until it is released, and announces +// that it has been entered. It stands in for any bandit whose selection is not +// instantaneous: one that reads a shared store, takes a contended lock, or is +// simply descheduled at the wrong moment. +type slowBandit struct { + entered chan struct{} + release chan struct{} + once sync.Once + calls atomic.Int64 + pick PolicyType +} + +func newSlowBandit(pick PolicyType) *slowBandit { + return &slowBandit{ + entered: make(chan struct{}), + release: make(chan struct{}), + pick: pick, + } +} + +func (b *slowBandit) RecordStats(ShadowStats) {} + +func (b *slowBandit) SelectPolicy() PolicyType { + b.calls.Add(1) + b.once.Do(func() { close(b.entered) }) + <-b.release + + return b.pick +} + +// TestEpoch_SlowBanditDoesNotBlockReaders is the property this whole +// arrangement exists to provide: a bandit that takes its time delays the +// switch it is deciding, and nothing else. +// +// Go's RWMutex queues new readers behind a waiting writer, so a bandit called +// while the cache's write lock is held stalls every Get in the process for its +// full duration -- a store timeout becomes a cache outage. The epoch therefore +// has to snapshot under the lock, release it, consult the bandit, and take the +// lock again to apply the result. +func TestEpoch_SlowBanditDoesNotBlockReaders(t *testing.T) { + t.Parallel() + + lru := newEvictingPolicy[string, int](LRU, 4) + lfu := newEvictingPolicy[string, int](LFU, 4) + bandit := newSlowBandit(LFU) + + cache, err := NewAdaptiveCache[string, int]( + []Policy[string, int]{lru, lfu}, + bandit, + &Settings{EpochRequests: 1, EvictPartialCapacityFilling: true}, + ) + require.NoError(t, err) + defer cache.Close() + defer close(bandit.release) + + // One Get ends an epoch, and that caller runs it. It parks in the bandit. + go func() { cache.Get("trigger") }() + + select { + case <-bandit.entered: + case <-time.After(5 * time.Second): + require.FailNow(t, "the bandit was never consulted") + } + + // The epoch is now mid-flight with the bandit blocked. A reader must not + // be waiting on it. + done := make(chan time.Duration, 1) + go func() { + start := time.Now() + cache.Get("reader") + done <- time.Since(start) + }() + + select { + case took := <-done: + assert.Less(t, took, 2*time.Second, + "Get waited on the bandit; it must only wait on the cache's own state") + case <-time.After(2 * time.Second): + require.FailNow(t, "Get blocked behind the bandit: a slow bandit is a cache outage") + } +} + +// conservingBandit records every arm report it is given, so a test can assert +// that an epoch's evidence is delivered exactly once. +type conservingBandit struct { + mu sync.Mutex + byArm map[PolicyType]int64 + reports int64 + pick PolicyType +} + +func newConservingBandit(pick PolicyType) *conservingBandit { + return &conservingBandit{byArm: map[PolicyType]int64{}, pick: pick} +} + +func (b *conservingBandit) RecordStats(s ShadowStats) { + b.mu.Lock() + defer b.mu.Unlock() + + b.byArm[s.Policy] += s.Hits + s.Misses + b.reports++ +} + +func (b *conservingBandit) SelectPolicy() PolicyType { + b.mu.Lock() + defer b.mu.Unlock() + + return b.pick +} + +func (b *conservingBandit) totals() (perArm map[PolicyType]int64, reports int64) { + b.mu.Lock() + defer b.mu.Unlock() + + out := make(map[PolicyType]int64, len(b.byArm)) + for k, v := range b.byArm { + out[k] = v + } + + return out, b.reports +} + +// TestEpoch_TickerAndRequestEpochsDoNotLoseCounts drives both epoch clocks at +// once from many goroutines. Each arm's counters are read and reset under the +// write lock, so whatever an epoch collects is delivered to the bandit exactly +// once and nothing is dropped between the two. +// +// The invariant checked is conservation: every request the cache counted for +// an arm reaches the bandit, or is still sitting in that arm's live counters. +// A snapshot delivered twice inflates the total; one collected and dropped on +// the way to the bandit deflates it. +func TestEpoch_TickerAndRequestEpochsDoNotLoseCounts(t *testing.T) { + t.Parallel() + + lru := newEvictingPolicy[string, int](LRU, 64) + lfu := newEvictingPolicy[string, int](LFU, 64) + bandit := newConservingBandit(LRU) + + cache, err := NewAdaptiveCache[string, int]( + []Policy[string, int]{lru, lfu}, + bandit, + &Settings{ + EpochDuration: time.Millisecond, + EpochRequests: 64, + EvictPartialCapacityFilling: true, + ShadowSampleRate: 1, + }, + ) + require.NoError(t, err) + + const ( + goroutines = 8 + perRoutine = 1250 + ) + + var wg sync.WaitGroup + for g := range goroutines { + wg.Add(1) + go func() { + defer wg.Done() + for i := range perRoutine { + key := string(rune('a' + (g+i)%26)) + cache.Add(key, i) + cache.Get(key) + } + }() + } + wg.Wait() + cache.Close() + + delivered, reports := bandit.totals() + require.Positive(t, reports, "no epoch ever reported") + + // Whatever has not been delivered must still be live in the arm. + for _, policy := range []*evictingPolicy[string, int]{lru, lfu} { + live := policy.GetStats() + got := delivered[policy.GetType()] + live.Hits + live.Misses + assert.Equal(t, int64(goroutines*perRoutine), got, + "%s: %d delivered plus %d live should account for every request", + policy.GetType(), delivered[policy.GetType()], live.Hits+live.Misses) + } +} + +// steppedBandit lets a test hold one epoch inside SelectPolicy while it drives +// another epoch to completion, which is the only way to produce a decision made +// from evidence that has since been superseded. +type steppedBandit struct { + entered chan struct{} + release chan struct{} + once sync.Once + picks atomic.Int64 + pick PolicyType +} + +func newSteppedBandit(pick PolicyType) *steppedBandit { + return &steppedBandit{ + entered: make(chan struct{}), + release: make(chan struct{}), + pick: pick, + } +} + +func (b *steppedBandit) RecordStats(ShadowStats) {} + +func (b *steppedBandit) SelectPolicy() PolicyType { + if b.picks.Add(1) == 1 { + b.once.Do(func() { close(b.entered) }) + <-b.release + } + + return b.pick +} + +// TestEpoch_StaleSelectionIsDiscarded pins the re-check that releasing the +// lock makes necessary. A decision is computed from one epoch's evidence; if +// another epoch has collected and reported since, that decision describes a +// cache state that no longer exists -- and the gates it was checked against +// were fed statistics that have since been overwritten. It is dropped, and the +// newer epoch decides instead. +func TestEpoch_StaleSelectionIsDiscarded(t *testing.T) { + t.Parallel() + + lru := newEvictingPolicy[string, int](LRU, 4) + lfu := newEvictingPolicy[string, int](LFU, 4) + bandit := newSteppedBandit(LFU) + + cache, err := NewAdaptiveCache[string, int]( + []Policy[string, int]{lru, lfu}, + bandit, + &Settings{EpochRequests: 1, EvictPartialCapacityFilling: true}, + ) + require.NoError(t, err) + defer cache.Close() + + require.Equal(t, LRU, cache.ActivePolicy()) + + // Epoch one parks inside the bandit holding no cache lock. + go func() { cache.Get("first") }() + select { + case <-bandit.entered: + case <-time.After(5 * time.Second): + require.FailNow(t, "the bandit was never consulted") + } + + // Epoch two runs to completion while epoch one is still parked. It is not + // blocked by epoch one, which is itself the point of the arrangement. + cache.Get("second") + + require.Eventually(t, func() bool { return bandit.picks.Load() >= 2 }, + 5*time.Second, 5*time.Millisecond, "the second epoch never reached the bandit") + + // Releasing epoch one lets its stale decision arrive last. It must not be + // applied on top of the newer epoch's. + close(bandit.release) + + assert.Eventually(t, func() bool { return cache.ActivePolicy() == LFU }, + 5*time.Second, 5*time.Millisecond, "the newer epoch's decision should stand") + + epochsRun := bandit.picks.Load() + assert.GreaterOrEqual(t, epochsRun, int64(2), + "both epochs ran; the stale one contributed its evidence and dropped its decision") +} From 8754e5a0859d9e5b11d1b1372d5dcf8a4374f7bf Mon Sep 17 00:00:00 2001 From: Sam Shaplygin Date: Fri, 28 Aug 2026 00:14:45 +0200 Subject: [PATCH 2/2] Call the bandit with no cache lock held An epoch used to run entirely under the cache's write lock, bandit call included. Go's RWMutex queues new readers behind a waiting writer, so for as long as SelectPolicy ran, every Get and Add in the process waited on it: a bandit that read a network store turned a store timeout into a cache outage. The interface managed this by forbidding it -- "an implementation must not block" -- which pushed the problem onto every implementer, and is why the distributed bandit had to buffer its reports and serve selections from an atomic. The epoch now runs in three phases: 1. under the write lock: close any gradual window, read and reset every arm's counters into an epochSnapshot; 2. holding no cache lock: deliver that snapshot and ask for a selection; 3. under the write lock again: check the selection is still current, apply. Phases 1 and 3 are each atomic, so no caller can observe a half-applied switch. Between them the cache is fully usable. Two things the split needs that holding one lock gave for free. Calls into the bandit are serialised by a new banditMu, taken only while no cache lock is held. A Bandit is caller-supplied code with no stated concurrency contract and it used to be entered under the write lock; that guarantee is kept, and the cost of a slow bandit is now paid by the next epoch rather than by every reader. A selection is dropped if another epoch collected while it was being made. Such a decision was computed from evidence two epochs old, and the stability gates would check it against c.epochStats the newer epoch has since overwritten -- it would be admitted or rejected on numbers that are not its own. epochsCollected is a second counter, and the reason is not decorative: epochID advances when an epoch is applied and the stability gates are specified against it, so reusing it for staleness moved the cooldown arithmetic by one epoch. TestSwitchStability_CooldownRearmsAfterSwitch caught that. The Bandit godoc and docs/design.md no longer forbid blocking. They state what it costs instead: the switch is delayed and may be dropped, concurrent Get and Add are not affected, and under EpochRequests the caller completing the epoch still pays -- as it already did for the switch and any migration. Verified: go test -race -short -count=5 on the root module, make all across nine modules, make evidence regenerated. Hit rates move within run-to-run noise (zipf 68.14% against a recorded 67.88%, with otter moving 73.36% against 73.25% in the same run) and no table is rewritten. Serial hot-path benchmarks are unchanged -- Get 102/153 ns, sampled 41/44, Add 98/144, sampled 53/55 -- which is expected: the new lock is on the epoch path, and Get's path is untouched. --- cache.go | 17 +++- docs/design.md | 16 +++- epoch.go | 200 +++++++++++++++++++++++++++++++++------------ epoch_lock_test.go | 53 +++++++----- interfaces.go | 24 ++++-- 5 files changed, 222 insertions(+), 88 deletions(-) diff --git a/cache.go b/cache.go index c0b67c8..3d6f498 100644 --- a/cache.go +++ b/cache.go @@ -56,7 +56,16 @@ type AdaptiveCache[K comparable, V any] struct { migrationRealKeys map[K]struct{} // --- Control Plane --- - bandit Bandit + + // banditMu serialises calls into the bandit. The write lock used to do + // that as a side effect of being held across the whole epoch; holding it + // there stalled every Get for as long as the bandit ran, so the epoch now + // releases it first. A Bandit is caller-supplied and nothing in its + // contract says it may be entered twice at once, so the guarantee has to + // come from somewhere: this lock is only ever held while no cache lock is, + // so a slow bandit delays the next epoch and nothing else. + banditMu sync.Mutex + bandit Bandit // epochBandit is bandit again when it implements the optional EpochBandit // extension, and nil otherwise. The assertion is made once at construction // rather than on every epoch, and its nil-ness is what selects between the @@ -92,6 +101,12 @@ type AdaptiveCache[K comparable, V any] struct { // --- Settings --- epochID int64 + // epochsCollected counts snapshots taken, and exists only so a decision + // can tell whether a later epoch has collected since it did. epochID + // cannot serve: it advances when an epoch is applied, which is after the + // bandit has been consulted, and the stability gates are specified against + // it. Two counters, because they answer two questions. + epochsCollected int64 // epochTicker is nil when the cache ends its epochs on request count // alone, since time.NewTicker rejects a non-positive duration. epochTicker *time.Ticker diff --git a/docs/design.md b/docs/design.md index ae6c3fd..b437d6e 100644 --- a/docs/design.md +++ b/docs/design.md @@ -121,9 +121,17 @@ type Bandit interface { } ``` -Both methods are called under the cache's write lock, so **an implementation -must not block**. Go's `RWMutex` queues new readers behind a waiting writer, so -a slow bandit stalls every `Get` in the process for its duration. +Both methods are called with **no cache lock held**. A bandit that takes its +time delays the switch it is deciding, and the epoch after it, but not the +cache's own operations — concurrent `Get` and `Add` are unaffected. The epoch +snapshots every arm's counters under the write lock, releases it, consults the +bandit, and takes the lock again to apply the result; a selection superseded by +a later epoch in the meantime is dropped rather than applied late. + +Two consequences worth knowing. Under `EpochRequests` the `Get` that completes +an epoch runs it, so that one caller does wait for the bandit — under +`EpochDuration` nobody does. And calls are serialised: no implementation is +entered from two goroutines at once. Ready-made bandits live in the `bandit` module: `bandit.NewThompson`, and `bandit.NewGreedy` as a control. Both examples use the first. @@ -146,7 +154,7 @@ func (a *adapter) RecordStats(s ascache.ShadowStats) { } func (a *adapter) SelectPolicy() ascache.PolicyType { - // Must return promptly and must not block: see the note above. + // No cache lock is held here; taking time delays the switch, not Get. // Returning Undefined -- or any policy the cache does not hold -- // means "no change", which is the right answer before the first epoch. return a.ext.Choose(a.arms) diff --git a/epoch.go b/epoch.go index e437a75..a49e60f 100644 --- a/epoch.go +++ b/epoch.go @@ -46,12 +46,57 @@ func (c *AdaptiveCache[K, V]) countRequest() { c.runEpoch() } -// runEpoch performs one epoch tick: it selects the next policy, migrates data -// when the policy changes and the stability gates allow it, and advances the -// epoch counter. The entire sequence runs under the write lock so concurrent -// cache operations never observe a half-applied switch (a torn activePolicy or -// partially migrated state). +// runEpoch performs one epoch tick in three phases, because the middle one +// must not hold the cache's lock. +// +// 1. under the write lock: close any gradual window, snapshot and reset every +// arm's counters, advance epochID; +// 2. holding no cache lock: deliver that snapshot to the bandit and ask it to +// select; +// 3. under the write lock again: check the decision is still current, then +// apply it. +// +// Phase 2 is the whole reason for the split. Go's RWMutex queues new readers +// behind a waiting writer, so a bandit called with the write lock held stalls +// every Get in the process for its full duration -- a store timeout becomes a +// cache outage. +// +// Phases 1 and 3 are each atomic, so no caller observes a half-applied switch. +// Between them the cache is fully usable and its contents may change; nothing +// in phase 3 assumes otherwise. func (c *AdaptiveCache[K, V]) runEpoch() { + snapshot := c.collectEpoch() + newPolicy := c.consultBandit(snapshot) + + c.applySelection(snapshot, newPolicy) +} + +// epochSnapshot is one epoch's evidence, taken under the write lock and then +// carried out of it. It is a value, not a view: the counters were read and +// reset in the same critical section, so nothing here can change underneath +// the bandit. +type epochSnapshot struct { + // epochID identifies the epoch this evidence was collected for. It is what + // the report carries and what a switch records as its epoch. + epochID int64 + // collected is the value of epochsCollected at the moment of collection. + // Phase 3 compares it against the current one to recognise a decision a + // later epoch has already superseded. + collected int64 + active PolicyType + // report is the per-arm evidence in policyOrder. It is delivered arm by + // arm to a plain Bandit and whole to an EpochBandit. + report []ShadowStats + // reported is false on an epoch the capacity gate skipped: nothing was + // measured, so the bandit is not consulted and nothing is applied. + reported bool + capacity int + sampleRate float64 +} + +// collectEpoch is phase 1: it closes any gradual migration window and takes +// the epoch's evidence. +func (c *AdaptiveCache[K, V]) collectEpoch() epochSnapshot { c.mu.Lock() defer c.mu.Unlock() @@ -63,12 +108,71 @@ func (c *AdaptiveCache[K, V]) runEpoch() { // comparable miniature by the time stats are collected below. c.closeMigrationLocked() - newPolicy := c.selectPolicyLocked() - if c.settings.ObserveOnly { - // Measure, report, advise - but never act. The cache keeps behaving - // exactly like the policy it was built with. - c.epochID++ + c.epochsCollected++ + snapshot := c.snapshotEpochLocked() + snapshot.collected = c.epochsCollected + + return snapshot +} + +// consultBandit is phase 2: it delivers the epoch's evidence and asks for the +// next policy. It holds no cache lock, so a bandit that blocks here delays the +// switch it is deciding and nothing else. +// +// banditMu still serialises the call. Overlapping epochs deliver disjoint +// evidence -- each snapshot was taken and reset in one critical section -- but +// a Bandit is caller-supplied code with no stated concurrency contract, and it +// used to be entered under the cache's write lock. That guarantee is kept. +func (c *AdaptiveCache[K, V]) consultBandit(snapshot epochSnapshot) PolicyType { + if !snapshot.reported { + return snapshot.active + } + + c.banditMu.Lock() + defer c.banditMu.Unlock() + + if c.epochBandit != nil { + c.epochBandit.RecordEpoch(EpochReport{ + EpochID: snapshot.epochID, + Active: snapshot.active, + Stats: snapshot.report, + Capacity: snapshot.capacity, + SampleRate: snapshot.sampleRate, + }) + } else { + for _, armStats := range snapshot.report { + c.bandit.RecordStats(armStats) + } + } + + return c.bandit.SelectPolicy() +} + +// applySelection is phase 3: it applies the bandit's choice if that choice is +// still the current epoch's to make. +// +// A decision is dropped when another epoch has collected since this one did. +// Such a decision was computed from evidence two epochs old, and worse, the +// gates would check it against c.epochStats that the newer epoch has already +// overwritten -- so it would be admitted or rejected on numbers that do not +// belong to it. The newer epoch decides instead. +func (c *AdaptiveCache[K, V]) applySelection(snapshot epochSnapshot, newPolicy PolicyType) { + c.mu.Lock() + defer c.mu.Unlock() + + // epochID counts ticks, including the ones the capacity gate skipped and + // the ones whose decision is dropped below, and it advances here so the + // stability gates see exactly the counter they saw before the epoch was + // split into phases. + defer func() { c.epochID++ }() + + // ObserveOnly measures, reports and advises, but never acts: the cache + // keeps behaving exactly like the policy it was built with. + if !snapshot.reported || c.settings.ObserveOnly { + return + } + if c.epochsCollected != snapshot.collected { return } @@ -82,8 +186,6 @@ func (c *AdaptiveCache[K, V]) runEpoch() { c.switchLocked(c.activePolicy, newPolicy) c.lastSwitchEpoch = c.epochID } - - c.epochID++ } // hasPolicy reports whether the cache holds the named policy as one of its @@ -94,31 +196,37 @@ func (c *AdaptiveCache[K, V]) hasPolicy(policyType PolicyType) bool { return ok } -// tryChangePolicy records every policy's stats with the bandit (nothing on a -// gated epoch — see selectPolicyLocked) and returns the policy selected for -// the next epoch. It acquires the write lock and performs no migration. It -// exists as a lock-acquiring entry point; callers that already hold the lock -// must use selectPolicyLocked instead. +// tryChangePolicy takes one epoch's evidence, delivers it to the bandit and +// returns the policy selected for the next epoch. It performs no migration and +// advances no epoch counter. It exists as a lock-acquiring entry point; +// callers inside the epoch use collectEpoch and consultBandit directly. func (c *AdaptiveCache[K, V]) tryChangePolicy() PolicyType { c.mu.Lock() - defer c.mu.Unlock() + snapshot := c.snapshotEpochLocked() + c.mu.Unlock() - return c.selectPolicyLocked() + return c.consultBandit(snapshot) } -// selectPolicyLocked reports every policy's stats to the bandit -- the active -// policy included, so its posterior does not go stale -- and returns the arm -// chosen for the next epoch. +// snapshotEpochLocked reads and resets every policy's counters and returns the +// epoch's evidence -- the active policy included, so its posterior does not go +// stale. // // With EvictPartialCapacityFilling false and the active policy not yet full it -// returns early, reporting and resetting nothing; counters accumulate until -// the next reporting epoch. Otherwise counters reset after delivery, the -// active policy's folded into globalStats first so Stats() stays cumulative -// and no active-tenure count leaks into a first shadow epoch after demotion. +// collects nothing, resets nothing and reports reported=false; counters +// accumulate until the next reporting epoch. Otherwise counters reset here, +// the active policy's folded into globalStats first so Stats() stays +// cumulative and no active-tenure count leaks into a first shadow epoch after +// demotion. +// +// Counters are read and reset in this one critical section, which is what lets +// the result leave the lock: the evidence cannot then be counted twice, and +// nothing the cache does next can change it. // // Caller must hold the write lock. -func (c *AdaptiveCache[K, V]) selectPolicyLocked() PolicyType { +func (c *AdaptiveCache[K, V]) snapshotEpochLocked() epochSnapshot { currentPolicy := c.activePolicy + snapshot := epochSnapshot{epochID: c.epochID, active: currentPolicy} // The capacity gate exists to avoid switching on the strength of a // half-full cache. In ObserveOnly mode nothing switches, so the gate would @@ -128,7 +236,8 @@ func (c *AdaptiveCache[K, V]) selectPolicyLocked() PolicyType { // Nothing was measured this epoch: drop the previous epoch's numbers // so the stability gates never compare against stale evidence. clear(c.epochStats) - return currentPolicy + + return snapshot } if c.epochStats == nil { @@ -139,13 +248,9 @@ func (c *AdaptiveCache[K, V]) selectPolicyLocked() PolicyType { } c.reportingEpochs++ - // An EpochBandit is handed the whole epoch in one call, so its report is - // collected here rather than delivered arm by arm. The slice is allocated - // per epoch and never reused, so the bandit may retain it. - var report []ShadowStats - if c.epochBandit != nil { - report = make([]ShadowStats, 0, len(c.policyOrder)) - } + // The slice is allocated per epoch and never reused, so an EpochBandit + // handed the whole of it may retain it. + report := make([]ShadowStats, 0, len(c.policyOrder)) // policyOrder rather than ranging the map: a map's order is random, and an // epoch's evidence should be reproducible for anything that hashes, @@ -179,28 +284,17 @@ func (c *AdaptiveCache[K, V]) selectPolicyLocked() PolicyType { tenure.Misses += reported.Misses c.tenureStats[policy.GetType()] = tenure - armStats := ShadowStats{ + report = append(report, ShadowStats{ Policy: policy.GetType(), Hits: reported.Hits, Misses: reported.Misses, - } - - if c.epochBandit != nil { - report = append(report, armStats) - continue - } - c.bandit.RecordStats(armStats) - } - - if c.epochBandit != nil { - c.epochBandit.RecordEpoch(EpochReport{ - EpochID: c.epochID, - Active: currentPolicy, - Stats: report, - Capacity: c.nominalCap[currentPolicy], - SampleRate: c.sampler.rate, }) } - return c.bandit.SelectPolicy() + snapshot.report = report + snapshot.reported = true + snapshot.capacity = c.nominalCap[currentPolicy] + snapshot.sampleRate = c.sampler.rate + + return snapshot } diff --git a/epoch_lock_test.go b/epoch_lock_test.go index 4149ab4..7f25d85 100644 --- a/epoch_lock_test.go +++ b/epoch_lock_test.go @@ -56,18 +56,19 @@ func TestEpoch_SlowBanditDoesNotBlockReaders(t *testing.T) { lfu := newEvictingPolicy[string, int](LFU, 4) bandit := newSlowBandit(LFU) + // The ticker drives epochs, so the goroutine that parks in the bandit is + // not one of the readers. (With EpochRequests the caller completing an + // epoch runs it, and that one caller does wait for the bandit -- by + // design, and documented on the setting. Every other caller must not.) cache, err := NewAdaptiveCache[string, int]( []Policy[string, int]{lru, lfu}, bandit, - &Settings{EpochRequests: 1, EvictPartialCapacityFilling: true}, + &Settings{EpochDuration: time.Millisecond, EvictPartialCapacityFilling: true}, ) require.NoError(t, err) defer cache.Close() defer close(bandit.release) - // One Get ends an epoch, and that caller runs it. It parks in the bandit. - go func() { cache.Get("trigger") }() - select { case <-bandit.entered: case <-time.After(5 * time.Second): @@ -201,14 +202,16 @@ type steppedBandit struct { release chan struct{} once sync.Once picks atomic.Int64 - pick PolicyType + first PolicyType + rest PolicyType } -func newSteppedBandit(pick PolicyType) *steppedBandit { +func newSteppedBandit(first, rest PolicyType) *steppedBandit { return &steppedBandit{ entered: make(chan struct{}), release: make(chan struct{}), - pick: pick, + first: first, + rest: rest, } } @@ -218,9 +221,11 @@ func (b *steppedBandit) SelectPolicy() PolicyType { if b.picks.Add(1) == 1 { b.once.Do(func() { close(b.entered) }) <-b.release + + return b.first } - return b.pick + return b.rest } // TestEpoch_StaleSelectionIsDiscarded pins the re-check that releasing the @@ -234,7 +239,9 @@ func TestEpoch_StaleSelectionIsDiscarded(t *testing.T) { lru := newEvictingPolicy[string, int](LRU, 4) lfu := newEvictingPolicy[string, int](LFU, 4) - bandit := newSteppedBandit(LFU) + // The first selection asks for a switch; every later one asks for no + // change. If the stale decision is applied the cache ends on LFU. + bandit := newSteppedBandit(LFU, LRU) cache, err := NewAdaptiveCache[string, int]( []Policy[string, int]{lru, lfu}, @@ -246,7 +253,7 @@ func TestEpoch_StaleSelectionIsDiscarded(t *testing.T) { require.Equal(t, LRU, cache.ActivePolicy()) - // Epoch one parks inside the bandit holding no cache lock. + // Epoch one collects, then parks inside the bandit holding no cache lock. go func() { cache.Get("first") }() select { case <-bandit.entered: @@ -254,21 +261,23 @@ func TestEpoch_StaleSelectionIsDiscarded(t *testing.T) { require.FailNow(t, "the bandit was never consulted") } - // Epoch two runs to completion while epoch one is still parked. It is not - // blocked by epoch one, which is itself the point of the arrangement. - cache.Get("second") + // Epoch two collects while epoch one is still parked -- it takes only the + // cache lock to do so, which epoch one is not holding. From here epoch + // one's decision describes a state two epochs old. + go func() { cache.Get("second") }() + require.Eventually(t, func() bool { + cache.mu.RLock() + defer cache.mu.RUnlock() - require.Eventually(t, func() bool { return bandit.picks.Load() >= 2 }, - 5*time.Second, 5*time.Millisecond, "the second epoch never reached the bandit") + return cache.epochsCollected >= 2 + }, 5*time.Second, time.Millisecond, "the second epoch never collected") - // Releasing epoch one lets its stale decision arrive last. It must not be - // applied on top of the newer epoch's. + // Now let epoch one's decision arrive. It must be dropped. close(bandit.release) - assert.Eventually(t, func() bool { return cache.ActivePolicy() == LFU }, - 5*time.Second, 5*time.Millisecond, "the newer epoch's decision should stand") + require.Eventually(t, func() bool { return bandit.picks.Load() >= 2 }, + 5*time.Second, time.Millisecond, "the second epoch never reached the bandit") - epochsRun := bandit.picks.Load() - assert.GreaterOrEqual(t, epochsRun, int64(2), - "both epochs ran; the stale one contributed its evidence and dropped its decision") + assert.Equal(t, LRU, cache.ActivePolicy(), + "a selection superseded by a later epoch must not be applied") } diff --git a/interfaces.go b/interfaces.go index c119e51..72d3b06 100644 --- a/interfaces.go +++ b/interfaces.go @@ -52,15 +52,23 @@ type Policy[K comparable, V any] interface { // ones live in the companion github.com/sshaplygin/as-cache/bandit module: // a local Thompson sampler and a greedy control. // -// # Implementations must not block +// # What a slow implementation costs // -// Both methods are called from the epoch goroutine while it holds the cache's -// write lock, so for as long as either runs, every Get and Add in the process -// is stalled behind it. A bandit that talks to the network, reads a file, or -// waits on a channel must do it on its own goroutine and have these methods -// only exchange buffered state. This is not a performance guideline: Go's -// RWMutex queues new readers behind a waiting writer, so a multi-second -// timeout here is a multi-second outage for the whole cache. +// Both methods are called with no cache lock held, so an implementation that +// takes its time delays the switch it is deciding and the epoch after it, but +// not the cache's own operations: concurrent Get and Add are unaffected. A +// bandit that talks to a network or a disk is therefore allowed, though the +// selection it produces is applied later than the epoch that asked for it, and +// is dropped entirely if a further epoch has collected in the meantime. +// +// One caller does still wait: under Settings.EpochRequests the Get that +// completes an epoch runs it, so that one call pays for the bandit as it +// already pays for the switch and any migration. Under Settings.EpochDuration +// the work is on the background goroutine and no caller waits at all. +// +// Calls are serialised: no implementation is entered from two goroutines at +// once, and the reports for one epoch arrive before that epoch's selection is +// asked for. type Bandit interface { // RecordStats delivers one policy's performance report. On every // reporting epoch each policy reports — the active policy included — so