feat: fill in NaN values when buffers dont contain requested data at the start or end of the buffers

This commit is contained in:
Aditya Ujeniya
2026-07-20 13:35:31 +02:00
parent f83f8afa86
commit 7b8668b8e4
5 changed files with 160 additions and 51 deletions
+106 -15
View File
@@ -303,43 +303,134 @@ func (b *buffer) firstWrite() int64 {
// - error: Non-nil on failure // - error: Non-nil on failure
// //
// Panics if 'data' slice is too small to hold all values in [from, to). // Panics if 'data' slice is too small to hold all values in [from, to).
func (b *buffer) read(from, to int64, data []schema.Float) ([]schema.Float, int64, int64, error) { // func (b *buffer) read(from, to int64, data []schema.Float) ([]schema.Float, int64, int64, error) {
// Walk back to the buffer that covers 'from', adjusting if we hit the oldest. // // Walk back to the buffer that covers 'from', adjusting if we hit the oldest.
for from < b.firstWrite() { // for from < b.firstWrite() {
if b.prev == nil { // if b.prev == nil {
from = b.firstWrite() // from = b.firstWrite()
break // break
} // }
// b = b.prev
// }
// i := 0
// t := from
// for ; t < to; t += b.frequency {
// idx := int((t - b.start) / b.frequency)
// if idx >= cap(b.data) {
// if b.next == nil {
// break
// }
// b = b.next
// // Recalculate idx in the new buffer; a gap between buffers may exist.
// idx = int((t - b.start) / b.frequency)
// }
// if idx >= len(b.data) {
// if b.next == nil || to <= b.next.start {
// break
// }
// data[i] += schema.NaN // NaN + anything = NaN; propagates missing data
// } else if t < b.start {
// data[i] += schema.NaN // gap before this buffer's first write
// } else {
// data[i] += b.data[idx]
// }
// i++
// }
// return data[:i], from, t, nil
// }
// read retrieves time-series data from the buffer chain for the specified time range.
//
// Traverses the buffer chain backwards (via prev links) if 'from' precedes the current
// buffer's start. Missing data points are represented as NaN. Values are accumulated
// into the provided 'data' slice (using +=, so caller must zero-initialize if needed).
//
// Points with no stored data (before the oldest write, in inter-buffer gaps, or
// past the newest write) are NaN-filled while scanning. Whether those NaN pads
// are kept is controlled by 'trim':
//
// - trim == false: the returned slice always spans the full requested window
// [from, to), NaN where missing. Use this when several scopes accumulate
// into one index-aligned slice (aggregation) or when the caller wants the
// whole window for display.
// - trim == true: the returned slice is cut down to the real data extent
// [dataFrom, dataTo), i.e. leading and trailing NaN pads are dropped so that
// result[0] aligns with dataFrom. Use this when persisting only real data
// (checkpointing).
//
// The second and third return values report the actual extent of real (non-NaN)
// data found within the window: dataFrom is the timestamp of the first stored
// sample and dataTo is one frequency past the last. When the buffer fully covers
// the request these equal from/to; when a scope holds less data they shrink,
// which is how callers detect misalignment across aggregated scopes. If no real
// data exists in the window both equal 'from' (and a trimmed read returns an
// empty slice).
//
// Parameters:
// - from: Start timestamp (Unix seconds)
// - to: End timestamp (Unix seconds, exclusive)
// - data: Pre-allocated slice to accumulate results (must be large enough)
// - trim: If true, cut the result to the real data extent; if false, keep the
// full NaN-padded [from, to) window.
//
// Returns:
// - []schema.Float: Full [from, to) window, or the [dataFrom, dataTo) slice if trim
// - int64: dataFrom — timestamp of the first real sample in the window
// - int64: dataTo — one frequency past the last real sample
// - error: Non-nil on failure
//
// Panics if 'data' slice is too small to hold all values in [from, to).
func (b *buffer) read(from, to int64, data []schema.Float, trim bool) ([]schema.Float, int64, int64, error) {
// Walk back toward the buffer that covers 'from'. Do NOT clamp 'from'
// forward when it predates the oldest write; those leading points are
// NaN-filled below so the full requested window is always available.
for from < b.firstWrite() && b.prev != nil {
b = b.prev b = b.prev
} }
freq := b.frequency // constant across the chain for this metric level
i := 0 i := 0
t := from t := from
dataFrom, dataTo := from, from // real (non-NaN) data extent within [from, to)
haveData := false
for ; t < to; t += b.frequency { for ; t < to; t += b.frequency {
idx := int((t - b.start) / b.frequency) idx := int((t - b.start) / b.frequency)
if idx >= cap(b.data) { if idx >= cap(b.data) {
if b.next == nil { if b.next == nil {
break data[i] += schema.NaN // past the newest buffer; rest is missing
i++
continue
} }
b = b.next b = b.next
// Recalculate idx in the new buffer; a gap between buffers may exist. // Recalculate idx in the new buffer; a gap between buffers may exist.
idx = int((t - b.start) / b.frequency) idx = int((t - b.start) / b.frequency)
} }
if idx >= len(b.data) { if idx >= len(b.data) || t < b.start {
if b.next == nil || to <= b.next.start { // Not yet written, or in the gap before this buffer's first write.
break
}
data[i] += schema.NaN // NaN + anything = NaN; propagates missing data data[i] += schema.NaN // NaN + anything = NaN; propagates missing data
} else if t < b.start {
data[i] += schema.NaN // gap before this buffer's first write
} else { } else {
data[i] += b.data[idx] data[i] += b.data[idx]
if !haveData {
dataFrom = t
haveData = true
}
dataTo = t + b.frequency
} }
i++ i++
} }
return data[:i], from, t, nil if trim {
// Drop leading/trailing NaN pads so result[0] aligns with dataFrom.
front := int((dataFrom - from) / freq)
back := int((dataTo - from) / freq)
return data[front:back], dataFrom, dataTo, nil
}
return data[:i], dataFrom, dataTo, nil
} }
// free removes buffers older than the specified timestamp from the chain. // free removes buffers older than the specified timestamp from the chain.
+4 -5
View File
@@ -298,15 +298,14 @@ func (l *Level) toCheckpointFile(from, to int64, m *MemoryStore) (*CheckpointFil
} }
data := make([]schema.Float, (to-from)/b.frequency+1) data := make([]schema.Float, (to-from)/b.frequency+1)
data, start, end, err := b.read(from, to, data)
// trim=true: read() returns only the real data extent, sliced so that
// Data[0] aligns with Start (the checkpoint reload contract).
data, start, _, err := b.read(from, to, data, true)
if err != nil { if err != nil {
return nil, err return nil, err
} }
for i := int((end - start) / b.frequency); i < len(data); i++ {
data[i] = schema.NaN
}
retval.Metrics[metric] = &CheckpointMetrics{ retval.Metrics[metric] = &CheckpointMetrics{
Frequency: b.frequency, Frequency: b.frequency,
Start: start, Start: start,
+21 -17
View File
@@ -708,32 +708,36 @@ func (m *MemoryStore) Read(selector util.Selector, metric string, from, to, reso
return nil, 0, 0, 0, errors.New("[METRICSTORE]> unknown metric: " + metric) return nil, 0, 0, 0, errors.New("[METRICSTORE]> unknown metric: " + metric)
} }
// data spans the full requested window; every scope's read() writes into the
// same index-aligned slice (NaN where it has no value), so aggregation never
// needs trimming. dataFrom/dataTo track the real (non-NaN) extent of the first
// scope seen; later scopes that report a different extent are misaligned. We
// no longer abort on misalignment — the full NaN-padded window is still
// returned for display — but we log it so the condition stays visible.
n, data := 0, make([]schema.Float, (to-from)/minfo.Frequency+1) n, data := 0, make([]schema.Float, (to-from)/minfo.Frequency+1)
var dataFrom, dataTo int64
err := m.root.findBuffers(selector, minfo.offset, func(b *buffer) error { err := m.root.findBuffers(selector, minfo.offset, func(b *buffer) error {
cdata, cfrom, cto, err := b.read(from, to, data) cdata, cfrom, cto, err := b.read(from, to, data, false)
if err != nil { if err != nil {
return err return err
} }
if n == 0 { if n == 0 {
from, to = cfrom, cto dataFrom, dataTo = cfrom, cto
} else if from != cfrom || to != cto || len(data) != len(cdata) { } else if cfrom != dataFrom || cto != dataTo {
missingfront, missingback := int((from-cfrom)/minfo.Frequency), int((to-cto)/minfo.Frequency) missingfront, missingback := int((dataFrom-cfrom)/minfo.Frequency), int((dataTo-cto)/minfo.Frequency)
if missingfront != 0 { switch {
return ErrDataDoesNotAlignMissingFront case missingfront != 0:
cclog.Warnf("%s", fmt.Errorf("%w: metric=%s buf#%d freq=%d ref[%d,%d] this[%d,%d] missingfront=%d pts",
ErrDataDoesNotAlignMissingFront, metric, n, minfo.Frequency, dataFrom, dataTo, cfrom, cto, missingfront))
case missingback != 0:
cclog.Warnf("%s", fmt.Errorf("%w: metric=%s buf#%d freq=%d ref[%d,%d] this[%d,%d] missingback=%d pts",
ErrDataDoesNotAlignMissingBack, metric, n, minfo.Frequency, dataFrom, dataTo, cfrom, cto, missingback))
default:
cclog.Warnf("%s", fmt.Errorf("%w: metric=%s buf#%d ref[%d,%d] this[%d,%d]",
ErrDataDoesNotAlignDataLenMismatch, metric, n, dataFrom, dataTo, cfrom, cto))
} }
newlen := len(cdata) - missingback
if newlen < 1 {
return ErrDataDoesNotAlignMissingBack
}
cdata = cdata[0:newlen]
if len(cdata) != len(data) {
return ErrDataDoesNotAlignDataLenMismatch
}
from, to = cfrom, cto
} }
data = cdata data = cdata
+24 -12
View File
@@ -141,23 +141,35 @@ func TestBufferWriteOverwrite(t *testing.T) {
// ─── Buffer read ────────────────────────────────────────────────────────────── // ─── Buffer read ──────────────────────────────────────────────────────────────
// TestBufferReadBeforeFirstWrite verifies that 'from' is clamped to firstWrite // TestBufferReadBeforeFirstWrite verifies that a range starting before any data
// when the requested range starts before any data in the chain. // in the chain returns the full requested window with the leading, unwritten
// points NaN-filled, and reports the real data extent (dataFrom) via the second
// return value (rather than clamping 'from' forward and truncating).
func TestBufferReadBeforeFirstWrite(t *testing.T) { func TestBufferReadBeforeFirstWrite(t *testing.T) {
b := newBuffer(100, 10) // firstWrite = 100 b := newBuffer(100, 10) // firstWrite = 100
b.write(100, schema.Float(1.0)) b.write(100, schema.Float(1.0))
b.write(110, schema.Float(2.0)) b.write(110, schema.Float(2.0))
data := make([]schema.Float, 10) data := make([]schema.Float, 10)
result, adjustedFrom, _, err := b.read(50, 120, data) result, dataFrom, dataTo, err := b.read(50, 120, data, false)
if err != nil { if err != nil {
t.Fatalf("read() error = %v", err) t.Fatalf("read() error = %v", err)
} }
if adjustedFrom != 100 { // Real data extent, not the requested window: first sample at 100, last at 110.
t.Errorf("adjustedFrom = %d, want 100 (clamped to firstWrite)", adjustedFrom) if dataFrom != 100 || dataTo != 120 {
t.Errorf("dataFrom/dataTo = %d/%d, want 100/120", dataFrom, dataTo)
} }
if len(result) != 2 { // Full window [50,120) step 10 = 7 points; first 5 predate firstWrite=100 -> NaN.
t.Errorf("len(result) = %d, want 2", len(result)) if len(result) != 7 {
t.Fatalf("len(result) = %d, want 7", len(result))
}
for i := 0; i < 5; i++ {
if !result[i].IsNaN() {
t.Errorf("result[%d] = %v, want NaN", i, result[i])
}
}
if result[5] != 1.0 || result[6] != 2.0 {
t.Errorf("result[5:7] = %v, want [1 2]", result[5:7])
} }
} }
@@ -177,7 +189,7 @@ func TestBufferReadChain(t *testing.T) {
b1.next = b2 b1.next = b2
data := make([]schema.Float, 6) data := make([]schema.Float, 6)
result, from, to, err := b2.read(100, 160, data) result, from, to, err := b2.read(100, 160, data, false)
if err != nil { if err != nil {
t.Fatalf("read() error = %v", err) t.Fatalf("read() error = %v", err)
} }
@@ -219,7 +231,7 @@ func TestBufferReadIdxAfterSwitch(t *testing.T) {
// from=0 triggers the walkback to b1 (from < b2.firstWrite=5). // from=0 triggers the walkback to b1 (from < b2.firstWrite=5).
// After clamping, the loop runs t=5,15,25,35. // After clamping, the loop runs t=5,15,25,35.
data := make([]schema.Float, 4) data := make([]schema.Float, 4)
result, _, _, err := b2.read(0, 36, data) result, _, _, err := b2.read(0, 36, data, false)
if err != nil { if err != nil {
t.Fatalf("read() error = %v", err) t.Fatalf("read() error = %v", err)
} }
@@ -247,7 +259,7 @@ func TestBufferReadNaNValues(t *testing.T) {
b.write(120, schema.Float(3.0)) b.write(120, schema.Float(3.0))
data := make([]schema.Float, 3) data := make([]schema.Float, 3)
result, _, _, err := b.read(100, 130, data) result, _, _, err := b.read(100, 130, data, false)
if err != nil { if err != nil {
t.Fatalf("read() error = %v", err) t.Fatalf("read() error = %v", err)
} }
@@ -274,7 +286,7 @@ func TestBufferReadAccumulation(t *testing.T) {
// Pre-populate data slice (simulates a second metric being summed in). // Pre-populate data slice (simulates a second metric being summed in).
data := []schema.Float{2.0, 1.0, 0.0} data := []schema.Float{2.0, 1.0, 0.0}
result, _, _, err := b.read(100, 120, data) result, _, _, err := b.read(100, 120, data, false)
if err != nil { if err != nil {
t.Fatalf("read() error = %v", err) t.Fatalf("read() error = %v", err)
} }
@@ -577,7 +589,7 @@ func TestBufferRead(t *testing.T) {
// Read data // Read data
data := make([]schema.Float, 3) data := make([]schema.Float, 3)
result, from, to, err := b.read(100, 130, data) result, from, to, err := b.read(100, 130, data, false)
if err != nil { if err != nil {
t.Errorf("buffer.read() error = %v", err) t.Errorf("buffer.read() error = %v", err)
} }
+5 -2
View File
@@ -7,6 +7,7 @@ package metricstore
import ( import (
"errors" "errors"
"fmt"
"math" "math"
"github.com/ClusterCockpit/cc-lib/v2/schema" "github.com/ClusterCockpit/cc-lib/v2/schema"
@@ -142,9 +143,11 @@ func (m *MemoryStore) Stats(selector util.Selector, metric string, from, to int6
if n == 0 { if n == 0 {
from, to = cfrom, cto from, to = cfrom, cto
} else if from != cfrom { } else if from != cfrom {
return ErrDataDoesNotAlignMissingFront return fmt.Errorf("%w: metric=%s buf#%d want[%d,%d] got[%d,%d]",
ErrDataDoesNotAlignMissingFront, metric, n, from, to, cfrom, cto)
} else if to != cto { } else if to != cto {
return ErrDataDoesNotAlignMissingBack return fmt.Errorf("%w: metric=%s buf#%d want[%d,%d] got[%d,%d]",
ErrDataDoesNotAlignMissingBack, metric, n, from, to, cfrom, cto)
} }
samples += stats.Samples samples += stats.Samples