diff --git a/tracker.go b/tracker.go index d33a8fb79..36625b7b8 100644 --- a/tracker.go +++ b/tracker.go @@ -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 { diff --git a/tracker_test.go b/tracker_test.go index 4afa1abce..fa6767e50 100644 --- a/tracker_test.go +++ b/tracker_test.go @@ -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()