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
108 changes: 68 additions & 40 deletions app/ocache/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,49 +10,15 @@ func WithPrometheus(reg *prometheus.Registry, namespace, subsystem string) Optio
if reg == nil {
return nil
}
if subsystem == "" {
subsystem = "cache"
}
nameSplit := strings.Split(namespace, ".")
subSplit := strings.Split(subsystem, ".")
namespace = strings.Join(nameSplit, "_")
subsystem = strings.Join(subSplit, "_")

return func(cache *oCache) {
c := NewPrometheusCollectors(namespace, subsystem, cache.Len)
c.MustRegister(reg)
cache.metrics = &metrics{
hit: prometheus.NewCounter(prometheus.CounterOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "hit",
Help: "cache hit count",
}),
miss: prometheus.NewCounter(prometheus.CounterOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "miss",
Help: "cache miss count",
}),
gc: prometheus.NewCounter(prometheus.CounterOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "gc",
Help: "garbage collected count",
}),
size: prometheus.NewGaugeFunc(prometheus.GaugeOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "size",
Help: "cache size",
}, func() float64 {
return float64(cache.Len())
}),
hit: c.Hit,
miss: c.Miss,
gc: c.GC,
size: c.Size,
}
reg.MustRegister(
cache.metrics.hit,
cache.metrics.miss,
cache.metrics.gc,
cache.metrics.size,
)
}
}

Expand All @@ -67,6 +33,68 @@ func WithPrometheusMetrics(hit, miss, gc prometheus.Counter, size prometheus.Gau
}
}

// PrometheusCollectors are the collectors a cache reports through: the ones
// WithPrometheus builds and registers, exposed for a caller that recreates
// its cache and so must register them once and hand them to every instance.
type PrometheusCollectors struct {
Hit, Miss, GC prometheus.Counter
Size prometheus.GaugeFunc
}

// NewPrometheusCollectors builds unregistered collectors with the names
// WithPrometheus would register (<namespace>_<subsystem>_{hit,miss,gc,size},
// dots turned into underscores, subsystem defaulting to "cache"). size is
// read through sizeFn, so it can resolve whichever cache is current.
func NewPrometheusCollectors(namespace, subsystem string, sizeFn func() int) PrometheusCollectors {
if subsystem == "" {
subsystem = "cache"
}
nameSplit := strings.Split(namespace, ".")
subSplit := strings.Split(subsystem, ".")
namespace = strings.Join(nameSplit, "_")
subsystem = strings.Join(subSplit, "_")
return PrometheusCollectors{
Hit: prometheus.NewCounter(prometheus.CounterOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "hit",
Help: "cache hit count",
}),
Miss: prometheus.NewCounter(prometheus.CounterOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "miss",
Help: "cache miss count",
}),
GC: prometheus.NewCounter(prometheus.CounterOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "gc",
Help: "garbage collected count",
}),
Size: prometheus.NewGaugeFunc(prometheus.GaugeOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "size",
Help: "cache size",
}, func() float64 {
return float64(sizeFn())
}),
}
}

// MustRegister registers the collectors with reg; like prometheus it panics
// on a second registration of the same names.
func (c PrometheusCollectors) MustRegister(reg prometheus.Registerer) {
reg.MustRegister(c.Hit, c.Miss, c.GC, c.Size)
}

// Option makes a cache report through these collectors without registering
// anything.
func (c PrometheusCollectors) Option() Option {
return WithPrometheusMetrics(c.Hit, c.Miss, c.GC, c.Size)
}

type metrics struct {
hit prometheus.Counter
miss prometheus.Counter
Expand Down
41 changes: 39 additions & 2 deletions app/ocache/metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,11 @@ package ocache

import (
"context"
"github.com/prometheus/client_golang/prometheus"
"github.com/stretchr/testify/require"
"strings"
"testing"

"github.com/prometheus/client_golang/prometheus"
"github.com/stretchr/testify/require"
)

func TestWithPrometheus_MetricsConvertsDots(t *testing.T) {
Expand All @@ -17,3 +18,39 @@ func TestWithPrometheus_MetricsConvertsDots(t *testing.T) {
require.NoError(t, err)
require.True(t, strings.Contains(cache.metrics.hit.Desc().String(), "some_name_some_system_hit"))
}

func TestWithPrometheus_Registers(t *testing.T) {
reg := prometheus.NewRegistry()
cache := New(func(ctx context.Context, id string) (value Object, err error) {
return &testObject{}, nil
}, WithPrometheus(reg, "some.name", "some.system"))
_, err := cache.Get(context.Background(), "id")
require.NoError(t, err)
families, err := reg.Gather()
require.NoError(t, err)
values := map[string]float64{}
for _, mf := range families {
m := mf.GetMetric()[0]
if m.GetGauge() != nil {
values[mf.GetName()] = m.GetGauge().GetValue()
} else {
values[mf.GetName()] = m.GetCounter().GetValue()
}
}
require.Equal(t, map[string]float64{
"some_name_some_system_hit": 0,
"some_name_some_system_miss": 1,
"some_name_some_system_gc": 0,
"some_name_some_system_size": 1,
}, values)
// the same names cannot be registered twice: the reason a caller that
// recreates its cache goes through NewPrometheusCollectors instead
require.Panics(t, func() {
New(nil, WithPrometheus(reg, "some.name", "some.system"))
})
require.NotPanics(t, func() {
c := NewPrometheusCollectors("some.name", "some.system", func() int { return 0 })
New(nil, c.Option())
New(nil, c.Option())
})
}
101 changes: 98 additions & 3 deletions app/ocache/ocache.go
Original file line number Diff line number Diff line change
Expand Up @@ -133,8 +133,9 @@ type OCache interface {
// RemoveSame closes and removes the object only if the value currently
// stored under id is exactly the given one (pointer identity). It lets a
// caller evict a specific instance it owns without racing a newer value
// that has replaced it under the same id. Returns ok=true only when this
// call performed the removal.
// that has replaced it under the same id. A value that is not comparable
// has no identity and is never matched (store such values behind a
// pointer). Returns ok=true only when this call performed the removal.
RemoveSame(ctx context.Context, id string, value Object) (ok bool, err error)
// TryRemove tries to close and to remove the object. ok reports whether
// this call removed it; (false, nil) means the object declined to close,
Expand Down Expand Up @@ -251,6 +252,81 @@ func (c *oCache) Pick(ctx context.Context, id string) (value Object, err error)
return val.waitLoad(ctx, id)
}

// PeekState is Peek's verdict about an id.
type PeekState int

const (
// PeekMiss: no entry for id (or the cache is closed)
PeekMiss PeekState = iota
// PeekBusy: an entry exists but is still loading or is being closed; a
// Get or Pick would wait for it
PeekBusy
// PeekHit: a loaded value, returned
PeekHit
)

// Peeker is the non-blocking read the cache returned by New offers on top of
// OCache; kept off that interface so other implementations stay valid.
type Peeker interface {
// Peek returns the value for id only if it is loaded and not being
// closed, without loading, waiting or allocating: the hot path for
// callers that handle a miss themselves. With touch a hit refreshes the
// GC deadline like Get. The state tells a miss from an entry that is
// loading or closing (a caller that must not act on a false miss waits
// for the latter, see WaitClosing). Peek counts no metrics: the caller,
// which decides whether the result is used, accounts for it.
Peek(id string, touch bool) (value Object, state PeekState)
// WaitClosing blocks while the entry for id is being closed, bounded by
// ctx, and returns at once when there is no such entry or it is not
// closing. It is the wait a caller needs before it can add a replacement
// for a value whose removal is still running.
WaitClosing(ctx context.Context, id string) error
}

func (c *oCache) Peek(id string, touch bool) (value Object, state PeekState) {
c.mu.Lock()
e, exists := c.data[id]
if c.closed || !exists {
c.mu.Unlock()
return nil, PeekMiss
}
if e.isClosing() {
c.mu.Unlock()
return nil, PeekBusy
}
select {
case <-e.load:
default:
// still loading
c.mu.Unlock()
return nil, PeekBusy
}
// value and loadErr are written before load closes; a failed load deletes
// its entry under c.mu, so a non-nil loadErr here means the entry is on
// its way out
if e.loadErr != nil || e.value == nil {
c.mu.Unlock()
return nil, PeekBusy
}
if touch {
e.lastUsage = time.Now()
}
value = e.value
c.mu.Unlock()
return value, PeekHit
}

func (c *oCache) WaitClosing(ctx context.Context, id string) error {
c.mu.Lock()
e, ok := c.data[id]
c.mu.Unlock()
if !ok {
return nil
}
_, err := e.waitClose(ctx, id)
return err
}

// ctx is the cancellable load context Get created together with the entry.
func (c *oCache) load(ctx context.Context, id string, e *entry) {
defer func() {
Expand Down Expand Up @@ -364,14 +440,33 @@ func (c *oCache) RemoveSame(ctx context.Context, id string, value Object) (ok bo
// if this call is the one that transitions it to closing. If e was already
// replaced under the same id it is in a closed state and remove() is a
// no-op, so a stale caller can never close the newer value that took the id.
same := exists && value != nil && e.value == value
same := exists && value != nil && sameObject(e.value, value)
c.mu.Unlock()
if !same {
return false, ErrNotExists
}
return c.removeCtx(ctx, e)
}

// sameObject reports whether stored is the very instance given. Pointer
// implementations (the usual kind) compare by identity. A value that is not
// comparable has no identity to check and must not panic the comparison: it
// never matches, so nothing is removed by mistake (a caller that needs
// instance-safe removal of such values stores them behind a pointer). Checked
// on the values, not the types: a struct with an interface field is
// comparable as a type and still panics when that field holds a slice.
func sameObject(stored, given Object) bool {
if stored == nil || given == nil {
// a still-loading entry has no value yet; nothing matches it
return false
}
sv, gv := reflect.ValueOf(stored), reflect.ValueOf(given)
if sv.Type() != gv.Type() || !sv.Comparable() || !gv.Comparable() {
return false
}
return stored == given
}

func (c *oCache) TryRemove(id string) (ok bool, err error) {
c.mu.Lock()

Expand Down
Loading
Loading