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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 16 additions & 1 deletion cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
16 changes: 12 additions & 4 deletions docs/design.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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)
Expand Down
200 changes: 147 additions & 53 deletions epoch.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand All @@ -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
}

Expand All @@ -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
Expand All @@ -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
Expand All @@ -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 {
Expand All @@ -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,
Expand Down Expand Up @@ -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
}
Loading
Loading