Replace query history map with ringBuffer

This commit is contained in:
Alan Bernstein 2020-10-26 13:16:29 -05:00
parent 65de5df2a9
commit bbd147d18f
2 changed files with 79 additions and 25 deletions

View file

@ -55,16 +55,53 @@ type queryStatusUpdate struct {
type queryTracker struct {
updates chan<- queryStatusUpdate
checks chan<- chan<- []*activeQuery
history map[pastQuery]struct{} // TODO not in memory
wg sync.WaitGroup
stop chan struct{}
history *ringBuffer
wg sync.WaitGroup
stop chan struct{}
}
type ringBuffer struct {
queries []pastQuery
start int
count int
mu sync.Mutex
}
// newRingBuffer initializes an empty RingBuffer of specified capacity.
func newRingBuffer(n int) *ringBuffer {
return &ringBuffer{
queries: make([]pastQuery, n),
start: 0,
count: 0,
}
}
// add adds a new element to the queue, overwriting the oldest if it is already full.
func (b *ringBuffer) add(q pastQuery) {
b.mu.Lock()
defer b.mu.Unlock()
b.queries[(b.start+b.count)%cap(b.queries)] = q
if b.count == cap(b.queries) {
b.start = (b.start + 1) % cap(b.queries)
} else {
b.count++
}
}
// slice returns the contents of the RingBuffer, in insertion order.
func (b *ringBuffer) slice() []pastQuery {
b.mu.Lock()
defer b.mu.Unlock()
return append(b.queries[b.start:b.count], b.queries[0:b.start]...)
}
func newQueryTracker() *queryTracker {
done := make(chan struct{})
updates := make(chan queryStatusUpdate, 128)
checks := make(chan chan<- []*activeQuery)
history := make(map[pastQuery]struct{})
history := newRingBuffer(128)
tracker := &queryTracker{
updates: updates,
checks: checks,
@ -85,7 +122,7 @@ func newQueryTracker() *queryTracker {
} else {
activeQueries[update.q] = struct{}{}
pq := pastQuery{update.q.query, update.q.node, update.q.started, update.endTime.Sub(update.q.started)} // TODO move end time to api.go
tracker.history[pq] = struct{}{}
tracker.history.add(pq)
}
case check := <-checks:
out := make([]*activeQuery, len(activeQueries))
@ -142,26 +179,7 @@ func (t *queryTracker) ActiveQueries() []ActiveQueryStatus {
}
func (t *queryTracker) PastQueries() []PastQueryStatus {
queries := make([]pastQuery, 0, len(t.history))
for pq := range t.history {
queries = append(queries, pq)
}
// TODO use sort.Sort
sort.Slice(queries, func(i, j int) bool {
switch {
case queries[i].started.Before(queries[j].started):
return true
case queries[i].started.After(queries[j].started):
return false
case queries[i].query < queries[j].query:
return true
case queries[i].query > queries[j].query:
return false
default:
return false
}
})
queries := t.history.slice()
now := time.Now()
out := make([]PastQueryStatus, len(queries))
for i, v := range queries {

View file

@ -15,10 +15,46 @@
package pilosa
import (
"fmt"
"testing"
"time"
)
func TestRingBuffer(t *testing.T) {
tests := []struct {
start int
count int
queries []string
}{
{start: 0, count: 0, queries: []string{}},
{start: 0, count: 1, queries: []string{"0"}},
{start: 0, count: 2, queries: []string{"0", "1"}},
{start: 0, count: 3, queries: []string{"0", "1", "2"}},
{start: 0, count: 4, queries: []string{"0", "1", "2", "3"}},
{start: 0, count: 5, queries: []string{"0", "1", "2", "3", "4"}},
{start: 1, count: 5, queries: []string{"1", "2", "3", "4", "5"}},
{start: 2, count: 5, queries: []string{"2", "3", "4", "5", "6"}},
{start: 3, count: 5, queries: []string{"3", "4", "5", "6", "7"}},
{start: 4, count: 5, queries: []string{"4", "5", "6", "7", "8"}},
{start: 0, count: 5, queries: []string{"5", "6", "7", "8", "9"}},
{start: 1, count: 5, queries: []string{"6", "7", "8", "9", "10"}},
{start: 2, count: 5, queries: []string{"7", "8", "9", "10", "11"}},
}
buffer := newRingBuffer(5)
for k := 0; k < len(tests); k++ {
if !(buffer.start == tests[k].start && buffer.count == tests[k].count) {
t.Fatalf("expected %d %d, found %d %d", tests[k].start, tests[k].count, buffer.start, buffer.count)
}
for n, q := range buffer.slice() {
if q.query != tests[k].queries[n] {
t.Fatalf("test[%d], buffer[%d] expected querystring '%s', found '%s'", k, n, tests[k].queries[n], q.query)
}
}
buffer.add(pastQuery{query: fmt.Sprintf("%d", k)})
}
}
func TestQueryTracker(t *testing.T) {
tracker := newQueryTracker()
defer tracker.Stop()