Documentation
¶
Overview ¶
Package aggregate implements the aggregate-first accounting engine: every accepted telemetry point is reduced into a small number of series deltas before sampling, so aggregate counts describe traffic rather than the sampling rate.
This file owns series identity. A series is identified by a SeriesKey — a fixed struct of dictionary IDs plus small bounded enums (ADR 0001, issue #159). The struct is directly comparable and therefore usable as a Go map key, and it is serialized field-wise with an explicit little-endian layout and a leading version byte. Raw struct memory is never persisted: padding, field reordering and endianness would all silently corrupt identity across builds.
Package aggregate implements the aggregate accounting engine.
Index ¶
- Constants
- Variables
- func AppendCanonicalDims(dst []byte, pairs []DimPair) []byte
- func DecodeSketchInto(dst *Sketch, data []byte) error
- func FinalizeCutoff(now time.Time) int64
- func InternDimValues(c *Cache, tenantID uint32, keys, values []string) uint32
- func IsDiskFull(err error) bool
- func IsGaugeLike(temporality Temporality, monotonic bool) bool
- func IsRequestSpan(root bool, spanKind int32) bool
- func LogGC(stats GCStats, err error)
- func LogRecovery(stats RecoveryStats, path string)
- func MutableSince(now time.Time) int64
- func NormalizeOperation(httpRoute, urlPath, spanName string) string
- func NormalizePath(path string) string
- func NormalizeSpanName(name string) string
- func PercentileFromSketch(s *Sketch) *latency.Percentile
- func RestoreMiner(store Store, miner *TemplateMiner) (int, error)
- func WindowStart(t time.Time) int64
- type AccuracyMetadata
- type Admission
- type AdmissionPlan
- type AggregateDelta
- func (d *AggregateDelta) Clone() *AggregateDelta
- func (d *AggregateDelta) HistogramPercentilesAvailable() bool
- func (d *AggregateDelta) HistogramQuantile(q float64) (value float64, lowerBound, ok bool)
- func (d *AggregateDelta) Merge(other *AggregateDelta)
- func (d *AggregateDelta) ObserveCounter(delta float64, reset bool)
- func (d *AggregateDelta) ObserveGauge(value float64, ts time.Time)
- func (d *AggregateDelta) ObserveHistogram(f HistogramFold)
- func (d *AggregateDelta) ObserveLog(ts time.Time, isError bool)
- func (d *AggregateDelta) ObserveSpan(durationMicros float64, isError, isRequest bool)
- type Applier
- type BacklogStats
- type Barrier
- type Baseline
- type BaselineRow
- type BaselineStats
- type BaselineTracker
- func (t *BaselineTracker) Baseline(key SeriesKey, producer ProducerID) (Baseline, bool)
- func (t *BaselineTracker) DirtyCount() int
- func (t *BaselineTracker) DrainDirty() []DirtyBaseline
- func (t *BaselineTracker) ObserveCumulative(key SeriesKey, producer ProducerID, startTime, ts time.Time, value float64) CumulativeOutcome
- func (t *BaselineTracker) Rollback(rows []DirtyBaseline)
- func (t *BaselineTracker) Seed(key SeriesKey, producer ProducerID, b Baseline)
- func (t *BaselineTracker) Stats() BaselineStats
- type BaselineTrackerConfig
- type Bounds
- type Bucket
- type BucketCursor
- type BucketPage
- type BucketSource
- type BucketsResult
- type Cache
- func (c *Cache) Bounds() Bounds
- func (c *Cache) Fence(ids map[uint32]struct{})
- func (c *Cache) Forget(ids map[uint32]struct{})
- func (c *Cache) Intern(tenantID uint32, kind Kind, value string) uint32
- func (c *Cache) InternBytes(tenantID uint32, kind Kind, value []byte) uint32
- func (c *Cache) InternDims(tenantID uint32, pairs []DimPair) uint32
- func (c *Cache) InternTenant(name string) (uint32, bool)
- func (c *Cache) Len() int
- func (c *Cache) Lookup(id uint32) (DictEntry, bool)
- func (c *Cache) OtherID(tenantID uint32, kind Kind) uint32
- func (c *Cache) Roots() map[uint32]struct{}
- func (c *Cache) SetOverflowSink(fn func(kind Kind, bound string))
- func (c *Cache) Stats() CacheStats
- func (c *Cache) Unfence()
- type CacheStats
- type Coverage
- type CumulativeOutcome
- type DashboardResult
- type DeltaMap
- type DeltaRow
- type DictEntry
- type DictRow
- type DimPair
- type DimsConfig
- type DirtyBaseline
- type DurableRegistrar
- func (r *DurableRegistrar) ClearTouched()
- func (r *DurableRegistrar) Committed(rows []DictRow)
- func (r *DurableRegistrar) DrainPending() []DictRow
- func (r *DurableRegistrar) Fence(ids map[uint32]struct{})
- func (r *DurableRegistrar) Forget(ids map[uint32]struct{})
- func (r *DurableRegistrar) Lookup(id uint32) (DictEntry, bool)
- func (r *DurableRegistrar) Next() uint32
- func (r *DurableRegistrar) OtherID(tenantID uint32, kind Kind) uint32
- func (r *DurableRegistrar) PendingCount() int
- func (r *DurableRegistrar) Register(tenantID uint32, kind Kind, value []byte) (uint32, error)
- func (r *DurableRegistrar) Revalidate(candidates map[uint32]struct{})
- func (r *DurableRegistrar) Roots() map[uint32]struct{}
- func (r *DurableRegistrar) Unfence()
- type EdgeInput
- type EdgeResolver
- type Engine
- func (e *Engine) ActiveSeriesKeys() map[SeriesKey]struct{}
- func (e *Engine) ApplyCommitted(m DeltaMap) uint64
- func (e *Engine) ApplyDeltas(m DeltaMap) uint64
- func (e *Engine) ApplyDeltasErr(m DeltaMap) (uint64, error)
- func (e *Engine) ApplyReducer(r *Reducer) uint64
- func (e *Engine) ApplyReducerErr(r *Reducer) (uint64, error)
- func (e *Engine) Baselines() *BaselineTracker
- func (e *Engine) Cache() *Cache
- func (e *Engine) CommitAdmission(plan *AdmissionPlan) uint64
- func (e *Engine) EdgeResolver() *EdgeResolver
- func (e *Engine) Epoch() string
- func (e *Engine) Limiter() *Limiter
- func (e *Engine) MarkFinalized(windowStart int64)
- func (e *Engine) MetricDims() DimsConfig
- func (e *Engine) Miner() *TemplateMiner
- func (e *Engine) Mode() string
- func (e *Engine) NewReducer(arrival time.Time) *Reducer
- func (e *Engine) Ownership() Ownership
- func (e *Engine) PlanAdmission(m DeltaMap) *AdmissionPlan
- func (e *Engine) PruneTopology()
- func (e *Engine) QueryBuckets(ctx context.Context, q Query) (*BucketsResult, error)
- func (e *Engine) QueryDashboard(ctx context.Context, q Query) (*DashboardResult, error)
- func (e *Engine) QueryTopology(ctx context.Context, q Query) (*TopologyResult, error)
- func (e *Engine) RestoreTopology(ids map[SeriesWindowKey]topoIdentity, deltas DeltaMap) int
- func (e *Engine) Revision() uint64
- func (e *Engine) RollbackAdmission(plan *AdmissionPlan)
- func (e *Engine) Rollover(now time.Time) int
- func (e *Engine) SetApplier(a Applier)
- func (e *Engine) SetStore(st Store)
- func (e *Engine) SetTemplateFactSink(fn func(TemplateFact))
- func (e *Engine) Snapshot() Snapshot
- func (e *Engine) Store() Store
- func (e *Engine) TenantID(name string) (uint32, bool)
- func (e *Engine) TopologyEpoch() uint64
- func (e *Engine) TopologyHorizon() time.Duration
- func (e *Engine) TopologyRevision(tenant string) uint64
- func (e *Engine) TopologySnapshot(tenant string) TopologySnapshot
- func (e *Engine) TopologyTenants() []string
- type EngineConfig
- type ExpBuckets
- type ExponentialHistogramInput
- type FailableApplier
- type FieldError
- type FinalizeStats
- type FinalizedPage
- type GCConfig
- type GCSnapshot
- type GCStats
- type GCStore
- type GroupBatch
- type GroupBy
- type HTTPClass
- type HistogramCommon
- type HistogramFold
- type HistogramInput
- type Kind
- type Limiter
- type LimiterConfig
- type LimiterStats
- type LogInput
- type MemRegistrar
- type MemRegistrarOptions
- type Method
- type MetricInput
- type MetricPointOutcome
- type MetricPointResult
- type MetricsRecorder
- type OverflowReason
- type Ownership
- type PointDisposition
- type PreloadError
- type ProducerID
- type PurgeStats
- type Query
- type RecoverOptions
- type RecoveryGate
- type RecoveryStats
- type Reducer
- func (r *Reducer) Arrival() time.Time
- func (r *Reducer) Deltas() DeltaMap
- func (r *Reducer) Len() int
- func (r *Reducer) MergeFrom(other *Reducer)
- func (r *Reducer) ReduceEdge(in EdgeInput)
- func (r *Reducer) ReduceExponentialHistogramPoint(in ExponentialHistogramInput) MetricPointResult
- func (r *Reducer) ReduceHistogramPoint(in HistogramInput) MetricPointResult
- func (r *Reducer) ReduceLog(in LogInput)
- func (r *Reducer) ReduceMetricPoint(in MetricInput)
- func (r *Reducer) ReduceSpan(in SpanInput)
- func (r *Reducer) Stats() ReducerStats
- type ReducerStats
- type Registrar
- type ResetReason
- type Resolver
- type ResourceIdentity
- type SQLiteStore
- func (s *SQLiteStore) Analyze() error
- func (s *SQLiteStore) Backlog() (BacklogStats, error)
- func (s *SQLiteStore) Close() error
- func (s *SQLiteStore) CommitGroup(b *GroupBatch) error
- func (s *SQLiteStore) FinalizableWindows(cutoff int64, limit int) ([]int64, error)
- func (s *SQLiteStore) FinalizeWindow(windowStart int64) (FinalizeStats, error)
- func (s *SQLiteStore) GCSnapshot() (*GCSnapshot, error)
- func (s *SQLiteStore) LoadBaselines(max int) ([]BaselineRow, error)
- func (s *SQLiteStore) LoadDict(max int) ([]DictRow, error)
- func (s *SQLiteStore) LoadSeries(max int) ([]SeriesRow, error)
- func (s *SQLiteStore) LoadTemplates(max int) ([]TemplateRow, error)
- func (s *SQLiteStore) Path() string
- func (s *SQLiteStore) PingContext(ctx context.Context) error
- func (s *SQLiteStore) PurgeBefore(cutoff int64) (PurgeStats, error)
- func (s *SQLiteStore) ReadBuckets(ctx context.Context, sel Selector) (BucketPage, error)
- func (s *SQLiteStore) ReadFinalizedSince(since int64, signals []Signal, limit int) (FinalizedPage, error)
- func (s *SQLiteStore) ReplayMutable(since int64) ([]DeltaRow, error)
- func (s *SQLiteStore) ResolveSeries(ids []SeriesID) ([]SeriesInfo, error)
- func (s *SQLiteStore) SaveTemplateStats(rows []TemplateStatRow) error
- func (s *SQLiteStore) SumBuckets(ctx context.Context, sel Selector, by GroupBy) ([]SumRow, error)
- func (s *SQLiteStore) SweepIdentities(series []SeriesID, dict []uint32, templates []uint32) (SweepStats, error)
- func (s *SQLiteStore) UUID() string
- func (s *SQLiteStore) VisitSketches(ctx context.Context, sel Selector, ...) error
- func (s *SQLiteStore) Watermarks() (uint32, SeriesID, error)
- type SaturationError
- type SchemaError
- type Selector
- type SeriesID
- type SeriesInfo
- type SeriesKey
- type SeriesRow
- type SeriesWindowKey
- type ServiceStat
- type Signal
- type Sketch
- func (s *Sketch) AppendTo(dst []byte) []byte
- func (s *Sketch) Collapsed() bool
- func (s *Sketch) Count() uint64
- func (s *Sketch) Downscale(target uint8) error
- func (s *Sketch) Encode() []byte
- func (s *Sketch) Merge(other *Sketch)
- func (s *Sketch) Observe(value float64)
- func (s *Sketch) ObserveBucket(index int32, n uint64)
- func (s *Sketch) ObserveN(value float64, n uint64)
- func (s *Sketch) ObserveZero(n uint64)
- func (s *Sketch) PopulatedBins() int
- func (s *Sketch) Quantile(q float64) float64
- func (s *Sketch) RelativeError() float64
- func (s *Sketch) Saturations() uint64
- func (s *Sketch) Scale() uint8
- func (s *Sketch) Sum() float64
- func (s *Sketch) ZeroCount() uint64
- type SketchDropReason
- type Snapshot
- type SnapshotEdge
- type SpanInput
- type StatusClass
- type Store
- type StoreConfig
- type StoreInspection
- type StoreMetrics
- type SumRow
- type SweepStats
- type TemplateFact
- type TemplateMiner
- func (m *TemplateMiner) Committed(rows []TemplateRow)
- func (m *TemplateMiner) DrainDirtyStats() []TemplateStatRow
- func (m *TemplateMiner) DrainPending() []TemplateRow
- func (m *TemplateMiner) Mine(tenant, service, severity, body string) (id uint32, isOther bool)
- func (m *TemplateMiner) MineAt(tenant, service, severity, body string, at time.Time) (id uint32, isOther bool)
- func (m *TemplateMiner) PartitionStats(tenant, service string) (TemplatePartitionStats, bool)
- func (m *TemplateMiner) PendingCount() int
- func (m *TemplateMiner) Restore(rows []TemplateRow)
- func (m *TemplateMiner) Roots() map[uint32]struct{}
- func (m *TemplateMiner) SetFactSink(fn func(TemplateFact))
- func (m *TemplateMiner) Stats() []TemplatePartitionStats
- func (m *TemplateMiner) TemplateText(id uint32) (string, bool)
- type TemplateMinerConfig
- type TemplatePartitionStats
- type TemplateRegistrar
- type TemplateRegistrarFunc
- type TemplateRegistration
- type TemplateRow
- type TemplateStatRow
- type Temporality
- type TopologyConfig
- type TopologyEdge
- type TopologyMetric
- type TopologyOperation
- type TopologyResult
- type TopologyService
- type TopologySnapshot
- type TopologyWindow
- type TrafficPoint
- type Variant
- type VersionError
- type WatermarkStore
- type WindowSnapshot
- type Writer
- func (w *Writer) Apply(m DeltaMap) uint64
- func (w *Writer) ApplyErr(m DeltaMap) (uint64, error)
- func (w *Writer) CollectIdentities() (GCStats, error)
- func (w *Writer) FinalizeDue(now time.Time) int
- func (w *Writer) RunBarrier(fn func()) error
- func (w *Writer) SaveTemplateStats() error
- func (w *Writer) SeriesID(key SeriesKey) SeriesID
- func (w *Writer) SeriesKeyByID(id SeriesID) (SeriesKey, bool)
- func (w *Writer) Shutdown(ctx context.Context) error
- func (w *Writer) Start()
- func (w *Writer) Stats() WriterStats
- func (w *Writer) Stop()
- type WriterConfig
- type WriterStats
Constants ¶
const ( DefaultMaxSeries = 6000 DefaultMaxSeriesMetrics = 2400 DefaultMaxSeriesTraces = 2400 DefaultMaxSeriesEdges = 500 DefaultMaxSeriesLogs = 500 DefaultMaxSeriesSystem = 200 DefaultMaxOperationsPerService = 20 DefaultMaxTraceSeriesPerService = 50 DefaultMaxMetricSeriesPerService = 50 DefaultMaxProducerBaselinesPerSeries = 8 )
Platform default caps, matching the #158 resolution and the AGGREGATE_* defaults in internal/config.
const ( // DefaultMaxValueBytes is the encoded-value length cap applied to every // non-tenant dictionary kind. An over-length value is routed to // __other__ and NEVER truncated: a truncated value is a different // identity wearing the same name, which is worse than an honest // "unnamed" bucket. DefaultMaxValueBytes = 512 // DefaultMaxTenantBytes is the stricter cap on a tenant name. Tenants are // asserted by clients (header, gRPC metadata, or an OTLP resource // attribute), so the namespace that scopes every other namespace gets the // tightest bound. DefaultMaxTenantBytes = 128 // DefaultMaxTenants is the instance-wide tenant-identity cap. A shared // API key grants access to every tenant, so an authenticated but hostile // client can still assert tenant names; this is what bounds that. DefaultMaxTenants = 256 // Instance-wide dictionary backstops for the namespaces that were // uncapped before #200. The per-tenant caps bound one tenant; these bound // the instance when many tenants each stay just under their own cap. DefaultMaxServicesPerTenant = 500 DefaultMaxServices = 5000 DefaultMaxDimKeysPerTenant = 200 DefaultMaxDimKeys = 2000 DefaultMaxDimValuesPerTenant = 5000 DefaultMaxDimValues = 50000 DefaultMaxDimTuplesPerTenant = 5000 DefaultMaxDimTuples = 50000 )
Identity bound defaults (#200 Q3).
const ( // WindowSize is the tumbling window width. Windows are aligned to UTC. WindowSize = 5 * time.Minute // AllowedLateness is how long after a window closes its points are still // accepted. Later points are excluded from aggregates and counted. AllowedLateness = 10 * time.Minute // MaxFutureSkew is how far ahead of arrival a point may be timestamped. // Future windows are never created. MaxFutureSkew = 2 * time.Minute // NumShards is the shard count. Fixed at four per #160. NumShards = 4 // DefaultMaxClosedWindows bounds how many closed-but-unfinalized windows // the engine keeps in RAM (#194 blocker 6). // // A window closes one lateness horizon after it opens and the writer's // finalize pass runs every 30 s, so the steady state is one closed window // or none. Twelve is an hour of finalizer backlog: enough that a slow or // briefly failing store never costs data, small enough that a wedged // finalizer cannot grow memory without bound. Past it the engine falls // back to the old lossy eviction and counts every window it drops. DefaultMaxClosedWindows = 12 )
Engine time bounds.
const ( ModeLegacy = "legacy" ModeShadow = "aggregate-shadow" ModeAggregate = "aggregate" )
Modes. Phase 1 ships legacy and aggregate-shadow; ModeAggregate is accepted by config but the read path does not switch over until a later phase.
const ( ExpHistogramMinScale int32 = -10 ExpHistogramMaxScale int32 = 20 )
Exponential-histogram scale bounds from the OTLP data model. Scales outside this range are malformed, not merely unrepresentable.
const ( // whole distribution, so no quantile may be published from it. HistPercentilesUnavailable uint32 = 1 << 0 // HistUnboundedTail marks that observations landed in a +Inf bucket. HistUnboundedTail uint32 = 1 << 1 // HistHasMin marks that the producer reported the minimum. HistHasMin uint32 = 1 << 2 // HistHasMax marks that the producer reported the maximum. HistHasMax uint32 = 1 << 3 )
Histogram delta flag bits. They live in AggregateDelta.HistogramFlags; the low byte is a bit set and the high byte carries the SketchDropReason.
const ( // SampledCoverageNote explains a partially aggregate-derived response. SampledCoverageNote = "counts and rates are exact for accepted telemetry; " + "parts of this response are derived from retained exemplars and are not complete" // ExemplarCoverageNote explains an exemplar-only response. ExemplarCoverageNote = "results come from retained raw exemplars only; " + "a missing exemplar does not mean zero matching events occurred" )
Coverage notes. They exist because "exemplar" is only honest if the surface also says what an absent exemplar does and does not mean.
const ( // ReasonCumulativeTemporality — GA aggregates delta-temporality // histograms only; see CLAUDE.md for the collector-side conversion. ReasonCumulativeTemporality = "cumulative_temporality" // ReasonUnspecifiedTemporality — a histogram with no temporality set is // not interpretable either way. ReasonUnspecifiedTemporality = "unspecified_temporality" // ReasonMalformedPoint — the point violates the OTLP data model. ReasonMalformedPoint = "malformed_point" // ReasonUnsupportedType — the point type has no aggregate model at all // (Summary). ReasonUnsupportedType = "unsupported_type" )
Rejection reasons reported to OTLP clients.
const ( // SketchDefaultScale is the platform-wide mapping scale. SketchDefaultScale uint8 = 4 // SketchMaxScale is the finest scale the codec accepts. The OTel data model // defines positive scales up to 20. SketchMaxScale uint8 = 20 // SketchMaxBins bounds the in-memory dense bin array. 512 bins at scale 4 // span a 2^32 dynamic range (~9.6 decades); the theoretical requirement for // 1us-100s is 426 bins. SketchMaxBins int32 = 512 )
Sketch parameters. The mapping is the OTLP exponential-histogram base-2 function at a fixed platform scale: bucket i covers (base^i, base^(i+1)] with base = 2^(2^-scale). At the platform default scale of 4 the base is 2^(1/16) ~= 1.044274 and the worst-case relative error of a quantile estimate is (base-1)/(base+1) ~= 2.17%.
Choosing the OTel mapping rather than an arbitrary DDSketch gamma makes a scale change an exact integer index shift ("perfect subsetting"), which is what allows Downscale and mismatched-scale Merge to stay lossless in the mapping. See docs/research/latency-sketch.md and issue #157.
const ( MetaDictWatermark = "dict_id_high_watermark" MetaSeriesWatermark = "series_id_high_watermark" )
Meta keys that carry monotonically increasing identity high-watermarks (#200 Q1). MAX(id)+1 stopped being a safe reseed the moment GC could delete the highest ID, so allocation comes from these instead.
const ( // TemplateWildcard is the placeholder token for a generalized position. TemplateWildcard = "<*>" // TemplateOther is the rendered text of a partition's overflow template. TemplateOther = "__other__" // DefaultMaxLogTemplatesPerService is the #158 per-service log-template // cap, wired to AGGREGATE_MAX_LOG_TEMPLATES_PER_SERVICE. DefaultMaxLogTemplatesPerService = 10 )
const ( // DefaultTopologyMaxServices bounds distinct services per tenant. DefaultTopologyMaxServices = 512 // DefaultTopologyMaxOperationsPerService bounds distinct operations per // (tenant, service). DefaultTopologyMaxOperationsPerService = 64 // DefaultTopologyMaxEdges bounds distinct caller/callee pairs per tenant. DefaultTopologyMaxEdges = 4096 // DefaultTopologyMaxMetrics bounds distinct (service, metric) pairs per // tenant. Metric windows carry the baselines the anomaly detector needs. DefaultTopologyMaxMetrics = 4096 // DefaultTopologyHorizon is how much finalized history the projection // retains behind the mutable set. Six five-minute windows is enough for a // rolling mean/variance baseline without pinning a day of state in RAM. DefaultTopologyHorizon = 30 * time.Minute )
Topology projection defaults.
const ( // DefaultCommitCoalesceMs is the first-waiter coalescing window. 25 ms, // not #160's provisional 5 ms: measured at 10k pts/s on 2 vCPU, 5 ms cost // a 37.8% writer duty cycle and a 286 ms ACK p99 against 109 ms at 25 ms // (#173). Kept in step with config.AggregateCommitCoalesceMs, which is // what production actually passes in. DefaultCommitCoalesceMs = 25 // DefaultCommitMaxDeltas is the early-commit count target — the "5k dirty // series" shape the gate benchmarks. DefaultCommitMaxDeltas = 5000 // DefaultCommitMaxBytes is the early-commit size target. DefaultCommitMaxBytes = 8 << 20 // DefaultCommitMaxPendingBytes bounds un-committed delta payload. DefaultCommitMaxPendingBytes = 64 << 20 // DefaultCommitMaxWaiters bounds Export goroutines parked on a commit. DefaultCommitMaxWaiters = 512 // DefaultCommitMaxPendingDeltas is the defensive delta-count bound. DefaultCommitMaxPendingDeltas = 200000 // DefaultFinalizeIntervalSec is how often the writer looks for windows // whose lateness horizon has expired. DefaultFinalizeIntervalSec = 30 )
Writer defaults. They are the provisional numbers of #160/#162; the benchmark gate ratifies or replaces them.
const CoverageHeader = "OtelContext-Data-Coverage"
CoverageHeader is the response header that carries Coverage on bare-array endpoints, where an envelope would silently break the response shape.
const (
// DefaultAggregateDBPath is the default AGGREGATE_DB_PATH.
DefaultAggregateDBPath = "./data/aggregate.db"
)
Default store tuning.
const DefaultEdgeResolverSpans = 100_000
DefaultEdgeResolverSpans bounds the span-ID memory across all shards. Each entry is two short strings plus map and list overhead (~0.1 KiB), so the default is a few tens of MiB at worst.
const DefaultTopologyRestoreMaxRows = 20000
DefaultTopologyRestoreMaxRows caps the finalized rows a topology restore reads. At the default series budget a 30-minute horizon is six windows of at most a few thousand series, so the cap is headroom rather than a routine truncation — and when it does bind, it is reported, not hidden.
const DimsRejectUnsupportedValue = "unsupported_value_type"
DimsRejectUnsupportedValue is the metric label for an attribute value that has no canonical scalar rendering (array, kvlist, bytes). The point is still aggregated -- under DimsID 0, the "no configured dims" sentinel.
const EncodedSeriesKeyLen = 1 + 4*4 + 5
EncodedSeriesKeyLen is the exact wire size of an encoded SeriesKey: one version byte, four little-endian uint32 fields, five enum bytes.
const GlobalTenant uint32 = 0
GlobalTenant is the tenant scope that KindTenant entries themselves live in. A tenant name cannot be scoped by its own ID, so the tenant dictionary is instance-global.
const IDPlaceholder = "{id}"
IDPlaceholder replaces a path segment the segment rules classify as variable.
const MaxDimensionKeys = 16
MaxDimensionKeys bounds one metric's configured dimension tuple. A tuple longer than this is refused from identity rather than silently truncated.
const MaxReadWindowSpan = int64(168*time.Hour/time.Second) + int64(WindowSize/time.Second)
MaxReadWindowSpan bounds the window range one ReadBuckets call may cover: the 168 h retention horizon plus one window of slack, expressed in seconds.
const OtherValue = "__other__"
OtherValue is the canonical dictionary value of the per-(tenant, kind) overflow entry. It is pre-created outside the capacity cap: a quota must never prevent creation of the entry that absorbs quota violations (#158).
const SeriesKeyVersion uint8 = 1
SeriesKeyVersion is the current encoding version. It is the first byte of every encoded key; decoders reject anything else rather than guessing a layout.
const SketchEncodingVersion uint8 = 0x01
SketchEncodingVersion is the only encoding version this package writes or accepts.
const SketchMaxSerializedBins = 256
SketchMaxSerializedBins caps the populated bins written by the encoder. It bounds a serialized sketch at roughly 1.5 KiB, which is what the 7-day worst-case disk arithmetic in issue #162 is budgeted against. A sketch with more populated bins is collapsed-lowest at encode time.
const StoreSchemaVersion = 5
StoreSchemaVersion is the aggregate schema version written into aggregate_meta at creation and verified at every open. There are no automatic migrations in v1 (#162): a mismatch fails startup and the operator chooses between an older binary and AGGREGATE_ALLOW_REBUILD=true.
v2 re-keyed aggregate_delta_log on (window_start, series_id) and dropped its append sequence (#173). A v1 file cannot be read by this binary and is not migrated: the rows it holds are unfinalized deltas with a retention horizon measured in minutes, so AGGREGATE_ALLOW_REBUILD loses far less than a migration would risk getting wrong.
v3 added request_count and error_request_count to both aggregate_delta_log and aggregate_buckets (#197 Q5). A v2 file holds only span counts, and there is no way to derive a request count from them after the fact — the parent and kind of every span it summarised are long gone. Backfilling zeros would make every historical window read "0 requests, N spans", which is worse than an operator-acknowledged rebuild, so the fail-closed policy stands unchanged.
v4 added the eight hist_* columns that carry an OTLP histogram point's population statistics and its accuracy metadata (#199). A v3 file has no column to put them in and no way to reconstruct them, so the same rebuild-or-downgrade choice applies.
v5 added aggregate_log_template — the durable log-template miner state — plus the dict_id/series_id high-watermark meta keys that dictionary GC requires (#200). A v4 file has neither, and a binary that started GC against a MAX(id)+1 reseed could re-mint an ID a finalized bucket still names, so the fail-closed policy stands: run an older binary or accept the rebuild.
Variables ¶
var ( // ErrKeyTruncated reports a buffer shorter than EncodedSeriesKeyLen. ErrKeyTruncated = errors.New("aggregate: encoded series key truncated") // ErrKeyTrailingBytes reports a buffer longer than EncodedSeriesKeyLen. ErrKeyTrailingBytes = errors.New("aggregate: encoded series key has trailing bytes") )
Decoding errors.
var ( // ErrSketchVersion reports an encoding version the decoder does not know. ErrSketchVersion = errors.New("aggregate: unsupported sketch encoding version") // ErrSketchTruncated reports input that ends inside a field. ErrSketchTruncated = errors.New("aggregate: truncated sketch encoding") // ErrSketchCorrupt reports input that is structurally invalid: bad varints, // non-canonical bins, impossible totals, or trailing bytes. ErrSketchCorrupt = errors.New("aggregate: corrupt sketch encoding") )
Codec errors. All decode failures wrap one of these.
var ( MaxDictRows = 2_000_000 MaxSeriesRows = 500_000 )
MaxDictRows and MaxSeriesRows bound the startup identity warm-up. They are far above the #158 caps; they exist so a corrupted or hostile file cannot make startup allocate without bound. Exceeding one fails startup with a *PreloadError rather than truncating the load in silence (#200 Q3).
Declared as var (not const) for the same reason MaxReadRows is: a test has to exercise the fail-fast path without seeding two million rows through a race-instrumented SQLite. Nothing outside a test may assign them.
var ( // ErrSaturated reports that the group-commit writer refused admission // because one of its three bounds (pending bytes, waiters, deltas) is // full. The ingest path maps it to gRPC RESOURCE_EXHAUSTED / HTTP 429, // exactly like the raw pipeline's ErrQueueFull. Errors returned by the // writer satisfy errors.Is(err, ErrSaturated) and carry which bound // tripped via *SaturationError. ErrSaturated = errors.New("aggregate: store admission saturated") // ErrStoreClosed reports use of a closed store or writer. ErrStoreClosed = errors.New("aggregate: store closed") // ErrSelectorUnbounded reports a read whose selector is missing a // mandatory bound (window range, tenant scope) or asks for more than the // store-side row cap. ErrSelectorUnbounded = errors.New("aggregate: unbounded selector") )
Store errors.
var ErrDictFull = errors.New("aggregate: dictionary full")
ErrDictFull is returned by a Registrar when a (tenant, kind) namespace has no capacity left. The Cache translates it into the pre-created __other__ ID and never propagates it to the hot path — identity resolution never fails.
var ErrHistogramMalformed = errors.New("aggregate: malformed histogram point")
ErrHistogramMalformed marks a data point that violates the OTLP data model. A malformed point is refused entirely -- it never contributes a count, a sum or a bin -- and is reported in ExportMetricsPartialSuccess.
var ErrSketchScale = errors.New("aggregate: unsupported sketch scale")
ErrSketchScale is returned when a scale outside the supported range is requested or decoded.
var ErrTenantRejected = errors.New("aggregate: tenant identity rejected")
ErrTenantRejected is returned when a tenant identity cannot be admitted: its encoded name is over the tenant length cap, it is empty, or the instance-wide tenant-identity cap is full.
Unlike every other namespace, a tenant is NEVER collapsed into __other__ (#200 Q3). Merging two tenants into one identity is a data-isolation failure, not a degradation: the point is refused and counted instead.
var MaxReadRows = 20000
MaxReadRows is the store-side cap on rows returned by one ReadBuckets call and on IDs accepted by one ResolveSeries call.
Declared as var (not const) so tests can temporarily shrink it and exercise the truncation and paging paths without seeding tens of thousands of rows through a race-instrumented SQLite — same reason storage.sqliteP99RowCap is a var. Nothing outside a test may assign it.
Functions ¶
func AppendCanonicalDims ¶
AppendCanonicalDims appends the canonical encoding of pairs to dst: pairs sorted by KeyID (ValueID breaks ties so duplicate keys still encode deterministically), each pair written as two varints. The same set of pairs in any order produces byte-identical output.
pairs is sorted in place — the caller's slice is scratch space on the hot path, not a value to preserve.
func DecodeSketchInto ¶
DecodeSketchInto parses a serialized sketch into dst, replacing whatever dst held. It is DecodeSketch without the allocation: a Sketch is a 2 KiB dense array, and a wide-range read decodes one per stored row, so the streaming readers reuse one scratch value per stream instead of leaving hundreds of thousands of them to the collector. On error dst is unspecified.
func FinalizeCutoff ¶
FinalizeCutoff returns the newest window start whose lateness horizon has expired at now — windows at or below it are ready to finalize.
func InternDimValues ¶
InternDimValues resolves configured dimension keys and their extracted values to dictionary IDs and returns the DimsID for a SeriesKey via Cache.InternDims. IDs come from the tenant-scoped dictionary, never from hashing: a hash collision would silently merge unrelated series. keys and values must be parallel slices (as returned by ExtractDimensionValues for the configured keys); an empty or mismatched input yields 0, the "no configured dims" sentinel.
func IsDiskFull ¶
IsDiskFull reports whether err is a device-out-of-space failure.
It checks the typed errno first (errors.Is unwraps *fs.PathError and any wrapping the storage layer added) and falls back to the driver's message text, which is the only channel the pure-Go SQLite driver offers for SQLITE_FULL.
func IsGaugeLike ¶
func IsGaugeLike(temporality Temporality, monotonic bool) bool
IsGaugeLike reports whether a metric point aggregates gauge-like — last, min, max, sum-of-samples and count, with no reset detection. Gauges do, and so do cumulative non-monotonic sums (UpDownCounter): negative movement there is legitimate, not a reset (#166 case 2).
func IsRequestSpan ¶
IsRequestSpan reports whether a span is a request entry point and therefore contributes to AggregateDelta.RequestCount.
The contract frozen in #197 Q2 is "root OR server span", and it is deliberately an OR: a root span with no parent starts a request whatever its kind, and a SERVER span is the server side of one even when the caller propagated a parent. Either qualifies, and a span that is both is still one request — the caller increments once per span, never once per condition.
func LogRecovery ¶
func LogRecovery(stats RecoveryStats, path string)
LogRecovery emits the operator-facing recovery summary.
func MutableSince ¶
MutableSince returns the oldest window start still inside the mutable set at now: the current window minus the lateness horizon. Everything strictly older is finalized history and never re-enters memory.
func NormalizeOperation ¶
NormalizeOperation resolves the operation name of an HTTP span using the precedence fixed in #159:
- http.route verbatim when present — the instrumentation already told us the template, and second-guessing it can only make things worse.
- otherwise url.path / http.target, normalized.
- otherwise the span name, normalized only if it is shaped "<METHOD> /path".
Nothing here learns or infers: the same inputs always produce the same output, on every process and every restart.
func NormalizePath ¶
NormalizePath normalizes a genuine URL path value (url.path or http.target). The query string and fragment are stripped first, then each segment is tested against the segment rules and replaced with IDPlaceholder when it matches.
Input that is not a path — invalid UTF-8, or no leading '/' once the query and fragment are gone — is returned verbatim, query included. It is somebody else's string and the per-service operation cap will catch it if it is pathological.
func NormalizeSpanName ¶
NormalizeSpanName normalizes a span name shaped "<METHOD> /path", which is what OTel HTTP instrumentation emits when it has no route template. Names of any other shape — including anything whose first token is not a known HTTP method — pass through verbatim; the URL rules must never be let loose on arbitrary span names.
func PercentileFromSketch ¶
func PercentileFromSketch(s *Sketch) *latency.Percentile
PercentileFromSketch describes one sketch-derived percentile: approximate with the sketch's own bound, unavailable when the sketch holds nothing. Shared with the legacy GraphRAG ServiceStore (#291), which feeds the same sketch type per service.
func RestoreMiner ¶
func RestoreMiner(store Store, miner *TemplateMiner) (int, error)
RestoreMiner warms miner from store before ingest starts. A store that cannot carry templates (an older implementation) restores nothing, which is the pre-#200 behaviour.
func WindowStart ¶
WindowStart returns the UTC-aligned start of the window containing t, as Unix seconds. Alignment is computed arithmetically rather than with Truncate so it is obviously independent of the local zone.
Types ¶
type AccuracyMetadata ¶
type AccuracyMetadata struct {
// Approximate is true whenever a percentile came from a sketch.
Approximate bool `json:"approximate"`
// SketchScale is the mapping scale of the merged sketch.
SketchScale uint8 `json:"sketch_scale"`
// RelativeErrorBound is the worst-case relative error of a quantile
// estimate at SketchScale, as a fraction (0.0217 = 2.17%).
RelativeErrorBound float64 `json:"relative_error_bound"`
// Degraded reports that the merged sketch collapsed or saturated, which
// puts estimates in the affected range OUTSIDE RelativeErrorBound.
Degraded bool `json:"degraded,omitempty"`
// SourceBucketError is the worst-case relative error imported from the
// SOURCE histogram's bucket widths, as a fraction. Zero for a native
// sketch and for exponential histograms, whose index transfer is exact.
// When non-zero it, not RelativeErrorBound, is the number that dominates.
SourceBucketError float64 `json:"source_bucket_error,omitempty"`
// UnboundedTail reports that observations landed in the source
// histogram's +Inf bucket. A quantile that falls there is a LOWER BOUND
// (>= UnboundedTailBound), never an estimate.
UnboundedTail bool `json:"unbounded_tail,omitempty"`
UnboundedTailBound float64 `json:"unbounded_tail_bound,omitempty"`
// PercentilesUnavailable suppresses every quantile: the sketch does not
// describe the whole distribution. PercentilesUnavailableReason names why
// (negative_observations, scale_out_of_range, no_finite_boundaries).
}
AccuracyMetadata describes how accurate a percentile in the response is. It is computed from the FINAL MERGED sketch on every response and is never a hard-coded figure: merging mismatched scales downscales to the coarser one, so a merged sketch can be less accurate than the platform default scale.
func AccuracyFromHistogramDelta ¶
func AccuracyFromHistogramDelta(d *AggregateDelta) AccuracyMetadata
AccuracyFromHistogramDelta derives the accuracy metadata of a folded OTLP histogram distribution. It layers the source histogram's provenance on top of the merged sketch's own bound, so a caller cannot read a 2.17% error bound off a distribution whose source buckets were a decade wide.
func AccuracyFromSketch ¶
func AccuracyFromSketch(s *Sketch) AccuracyMetadata
AccuracyFromSketch derives the accuracy metadata of a merged sketch. A nil sketch still reports approximate: the path is approximate whether or not any duration happened to arrive in the window.
type Admission ¶
type Admission struct {
// Key is the series to record under. It differs from the requested key
// when a cap forced overflow routing.
Key SeriesKey
// Overflowed reports that a cap was hit and Key is an __other__ series.
Overflowed bool
// Reason names the cap that triggered overflow.
Reason OverflowReason
// Reserved reports that THIS call created (Key, window) presence, and so
// that a matching Release is the exact undo. It is false when the pair was
// already present, which is what lets a rolled-back group commit release
// only the occupancy it charged and never occupancy an earlier committed
// batch is still using (#194 blocker 3).
Reserved bool
}
Admission is the outcome of evaluating one series against the budget.
type AdmissionPlan ¶
type AdmissionPlan struct {
// Resolved is the batch keyed by the identity to record under.
Resolved DeltaMap
// contains filtered or unexported fields
}
AdmissionPlan is one batch's cardinality reservation: the identities the shards will hold, plus the occupancy this plan charged the limiter for.
It exists because admission has to precede the durable write — the row that becomes durable must carry the identity the shards will hold — while the charge must not survive a write that never happened (#194 blocker 3). Before it, a failed CommitGroup left the occupancy charged with no shard window to release it from, so a store that kept refusing writes ate the cardinality budget permanently and forced live telemetry into __other__.
type AggregateDelta ¶
type AggregateDelta struct {
// Count is the number of accepted points that contributed to this series:
// spans for traces, log records for logs, data points for metrics. It is
// accepted telemetry, never sampled telemetry (#153 §8).
Count uint64
// ErrorCount is the subset of Count classified as an error: span status
// ERROR for traces, severity ERROR/FATAL for logs. Metrics never set it.
ErrorCount uint64
// RequestCount is the subset of Count that qualifies as a REQUEST: a span
// that is a trace root (no parent span) or a SERVER-kind span. Either
// qualifies and a span is counted at most once (#197 Q2).
//
// It exists because Count is per SPAN, and a dashboard that labels a span
// total "traces" is lying by a factor of the average trace size. A
// distributed trace with several entry points counts once per entry point:
// a documented approximation, not a unique-trace-ID count, because trace
// IDs are on the permanent banned list (#153, #159) and cannot be counted
// without carrying them into aggregate identity.
//
// Only trace-shaped signals set it; logs and metrics leave it zero.
RequestCount uint64
// ErrorRequestCount is the error subset of RequestCount. It is the
// numerator of the headline dashboard error rate (#197 Q3); per-operation
// error rates stay span-based on ErrorCount/Count.
ErrorRequestCount uint64
// DurationCount is the number of duration observations. It tracks Count
// for trace series and stays zero elsewhere, so a merged delta can still
// tell "no spans" from "spans with zero duration".
DurationCount uint64
// DurationSum is the sum of observed durations, in microseconds.
DurationSum float64
// DurationMin and DurationMax are the extremes of the observed durations
// in microseconds. They are only meaningful when DurationCount > 0.
DurationMin float64
DurationMax float64
// Sketch holds the quantile contribution of the observed durations. It is
// allocated on the first duration observation and stays nil for series
// that carry no latency (logs, metrics), which keeps the common delta
// small — a Sketch is ~2 KiB of dense bins.
Sketch *Sketch
GaugeCount uint64
GaugeSum float64
GaugeMin float64
GaugeMax float64
// GaugeLast is the sample with the highest GaugeLastTime seen so far.
// Merging picks the later timestamp so the result does not depend on the
// order in which deltas were merged. Two samples of one series carrying
// the IDENTICAL timestamp are ambiguous about which is last, and that tie
// resolves by arrival order — no merge rule can fix a producer that
// timestamps two values the same.
GaugeLast float64
GaugeLastTime time.Time
// CounterDelta is the increase attributed to this window: the sum of
// per-point deltas computed by the baseline tracker for cumulative sums,
// or the raw values of delta-temporality points.
CounterDelta float64
// ResetCount is the number of counter resets observed in this window
// (start-time change or value regression, per #166).
ResetCount uint64
// HistogramCount is the population count reported by the folded points.
// It counts OBSERVATIONS, unlike Count which counts data points.
HistogramCount uint64
// HistogramSum is the authoritative sum reported by the producer. It is
// never derived from the sketch: bucket folding only knows bucket
// midpoints, and a sum of midpoints is an estimate.
HistogramSum float64
// HistogramMin and HistogramMax are the producer-reported extremes, valid
// only when HistHasMin / HistHasMax are set in HistogramFlags.
HistogramMin float64
HistogramMax float64
// HistogramFlags is the bit set defined in histogram.go:
// HistPercentilesUnavailable, HistUnboundedTail, HistHasMin, HistHasMax,
// plus a SketchDropReason in its high byte.
HistogramFlags uint32
// HistogramSourceError is the worst-case relative error imported from the
// SOURCE histogram's bucket widths, as a fraction. Zero for exponential
// histograms; an explicit-bounds histogram with decade-wide buckets can
// put this an order of magnitude above the sketch's own bound.
HistogramSourceError float64
// HistogramTailBound is the last finite boundary below the +Inf bucket,
// and HistogramTailCount how many observations landed above it. They are
// NOT in the sketch: a quantile that falls in the tail is answerable only
// as "at least HistogramTailBound".
HistogramTailBound float64
HistogramTailCount uint64
// LogCount is the number of log records in this series. It is redundant
// with Count for log series and zero everywhere else; it exists so a
// merged multi-signal view can still answer "how many logs" without
// carrying the SeriesKey alongside.
LogCount uint64
// FirstTimestamp and LastTimestamp bound the log records that contributed.
// Zero when LogCount is zero.
FirstTimestamp time.Time
LastTimestamp time.Time
}
AggregateDelta is the compact aggregate contribution of one Export request to one series in one window (CONTEXT.md, "Delta"). It is what request-local reduction produces and what the engine applies; finalized buckets are built from deltas, never the other way round.
Every field is additive or order-independent so two deltas for the same (series, window) merge without consulting the points that produced them. That is what lets the Phase 2 group-commit writer (#173) pre-merge deltas inside a transaction, and what lets Phase 1 apply them straight to the shards.
A delta is not safe for concurrent use. Reducers are request-local and each owns its deltas until they are handed to the engine.
func (*AggregateDelta) Clone ¶
func (d *AggregateDelta) Clone() *AggregateDelta
Clone returns a deep copy, including the sketch. Used by Snapshot so a caller can read a consistent view without holding a shard lock.
func (*AggregateDelta) HistogramPercentilesAvailable ¶
func (d *AggregateDelta) HistogramPercentilesAvailable() bool
HistogramPercentilesAvailable reports whether a quantile may be published for this delta's histogram distribution.
func (*AggregateDelta) HistogramQuantile ¶
func (d *AggregateDelta) HistogramQuantile(q float64) (value float64, lowerBound, ok bool)
HistogramQuantile estimates the q-quantile of the folded histogram distribution.
The three-way return is the whole point (#199 Q2). lowerBound=true means the quantile falls inside the source histogram's +Inf bucket: the answer is "at least value", not "approximately value", and a caller that renders it as an ordinary number is fabricating a tail it was never given. ok=false means no quantile may be published at all -- an empty window, or a point whose percentiles were suppressed.
func (*AggregateDelta) Merge ¶
func (d *AggregateDelta) Merge(other *AggregateDelta)
Merge folds other into d. Every field is additive or order-independent, so merging is associative and commutative: the result depends only on the multiset of observations, never on the order they arrived in.
func (*AggregateDelta) ObserveCounter ¶
func (d *AggregateDelta) ObserveCounter(delta float64, reset bool)
ObserveCounter records one counter increase already converted to a delta by the baseline tracker. reset marks that the increase followed a counter reset.
func (*AggregateDelta) ObserveGauge ¶
func (d *AggregateDelta) ObserveGauge(value float64, ts time.Time)
ObserveGauge records one gauge-like sample at ts.
func (*AggregateDelta) ObserveHistogram ¶
func (d *AggregateDelta) ObserveHistogram(f HistogramFold)
ObserveHistogram records one folded OTLP histogram data point.
Count counts the DATA POINT; HistogramCount counts the observations it summarizes. Both are needed: the first is the reduction denominator, the second is the population the percentiles describe.
A fold with percentiles unavailable contributes its scalars and poisons the series' percentile availability for the window. That is deliberate: once one point in a window could not be represented, no quantile computed from the rest describes the window's distribution.
func (*AggregateDelta) ObserveLog ¶
func (d *AggregateDelta) ObserveLog(ts time.Time, isError bool)
ObserveLog records one log record at ts. isError marks ERROR/FATAL severity.
func (*AggregateDelta) ObserveSpan ¶
func (d *AggregateDelta) ObserveSpan(durationMicros float64, isError, isRequest bool)
ObserveSpan records one span: its count, its error classification, whether it is a request entry point, and its duration in microseconds. Negative durations are recorded in the counters but clamped to zero for the sketch, which is what Sketch.Observe does with them anyway (latency is non-negative by construction).
isRequest is IsRequestSpan's verdict for the span. It is passed in rather than derived here because the delta never sees the span's parent or kind.
type Applier ¶
Applier applies deltas to the engine's shards. Phase 1 wires the direct applier, which mutates the shards inline. Phase 2 (#173) interposes the group-commit writer here: it batches deltas from many Export requests into one SQLite transaction and calls ApplyCommitted only after the COMMIT, making the shards a projection of committed state. No caller changes when it does.
type BacklogStats ¶
type BacklogStats struct {
// Rows is the number of delta-log rows currently awaiting finalization.
Rows int64
// OldestWindow is the oldest un-finalized window start, 0 when empty.
OldestWindow int64
// Bytes is the approximate delta-log payload size.
Bytes int64
}
BacklogStats describes the delta-log backlog — the health bound #160 requires an operator to be able to alarm on.
type Barrier ¶
type Barrier interface {
// RunBarrier executes fn with no group commit in flight and no commit
// able to start until it returns. fn must be bounded — it is on the ACK
// path of every Export parked behind the writer.
RunBarrier(fn func()) error
}
Barrier is the identity-maintenance barrier the collector runs its sweep through (#200 Q2). It is implemented by the group-commit writer: running the sweep on the writer goroutine is what makes "no commit interleaves with a delete" structural rather than a convention.
type Baseline ¶
type Baseline struct {
// StartTime is the point's start_time_unix_nano. A change means the
// producer restarted its counter.
StartTime time.Time
// LastTimestamp is the timestamp of the last accepted point. Points at or
// before it are stale or duplicate.
LastTimestamp time.Time
// Value is the last accepted cumulative value.
Value float64
}
Baseline is the per-(series, producer) record of the last cumulative point.
type BaselineRow ¶
type BaselineRow struct {
SeriesID SeriesID
Producer ProducerID
Baseline Baseline
}
BaselineRow is one durable cumulative baseline (#166). It is upserted inside the same group commit as the deltas it justifies, so a restart never has to re-seed a counter it already acknowledged points for.
type BaselineStats ¶
type BaselineStats struct {
// Entries is the number of live baseline records.
Entries int
// Series is the number of series holding at least one baseline.
Series int
// Stale counts points ignored as stale or duplicate.
Stale uint64
// Seeded counts baselines created.
Seeded uint64
// Gaps counts baselines re-seeded after a downtime gap.
Gaps uint64
// ResetsStartTime and ResetsRegression count resets by reason.
ResetsStartTime uint64
ResetsRegression uint64
// ProducerOverflow counts points routed to a degraded shared baseline
// because the per-series producer bound was exhausted.
ProducerOverflow uint64
// GlobalOverflow counts points that could not take a dedicated baseline
// because the global baseline budget was exhausted.
GlobalOverflow uint64
// Owed is the number of records currently carrying an increase whose
// group commit failed. A number that does not return to zero means the
// store is refusing writes, not that accounting is drifting.
Owed int
// Recovered counts points that re-attributed a stranded increase.
Recovered uint64
// Stranded counts stranded increases dropped by a downtime re-seed,
// which is the one place the owed ledger is deliberately not honoured.
Stranded uint64
}
BaselineStats is a snapshot of BaselineTracker counters.
type BaselineTracker ¶
type BaselineTracker struct {
// contains filtered or unexported fields
}
BaselineTracker converts cumulative monotonic points into deltas and detects resets. It is safe for concurrent use.
func NewBaselineTracker ¶
func NewBaselineTracker(cfg BaselineTrackerConfig) *BaselineTracker
NewBaselineTracker returns a tracker bounded by cfg. Zero or negative bounds take the platform defaults.
func (*BaselineTracker) Baseline ¶
func (t *BaselineTracker) Baseline(key SeriesKey, producer ProducerID) (Baseline, bool)
Baseline returns a copy of the baseline for (key, producer). Tests and diagnostics only.
func (*BaselineTracker) DirtyCount ¶
func (t *BaselineTracker) DirtyCount() int
DirtyCount reports how many baselines are awaiting a durable upsert.
func (*BaselineTracker) DrainDirty ¶
func (t *BaselineTracker) DrainDirty() []DirtyBaseline
DrainDirty returns the baselines mutated since the last drain and clears the dirty set. The writer calls it while building a group batch, so every drained record is at least as new as the deltas in that batch: on a crash a baseline can be ahead of the durable deltas (the next point under-counts by one interval) but never behind them (which would double-count).
Each drained row also carries the increase that record emitted since the previous drain, so Rollback can hand it back if the commit fails.
func (*BaselineTracker) ObserveCumulative ¶
func (t *BaselineTracker) ObserveCumulative(key SeriesKey, producer ProducerID, startTime, ts time.Time, value float64) CumulativeOutcome
ObserveCumulative applies the normative #166 evaluation order to one cumulative monotonic point and returns what to do with it.
key must be the full canonical metric-series identity BEFORE cardinality overflow routing: two series that later collapse into one __other__ series still have independent counters, and merging their baselines would invent resets.
func (*BaselineTracker) Rollback ¶
func (t *BaselineTracker) Rollback(rows []DirtyBaseline)
Rollback puts drained baselines back after a failed commit AND hands each record the increase that commit was carrying, so the amount is re-attributed to the next point instead of being lost (#194 blocker 2).
Re-dirtying alone was the bug: the in-memory baseline had already advanced at reduction time, so an identical client retry classified as stale and its delta vanished. See baselineState for why the amount is carried forward rather than the value rewound.
func (*BaselineTracker) Seed ¶
func (t *BaselineTracker) Seed(key SeriesKey, producer ProducerID, b Baseline)
Seed installs a baseline read back from the durable store at startup. It does NOT mark the record dirty: it is already durable, and re-writing every recovered baseline on the first commit would make restart the most expensive transaction the store ever runs.
func (*BaselineTracker) Stats ¶
func (t *BaselineTracker) Stats() BaselineStats
Stats returns a snapshot of the tracker counters.
type BaselineTrackerConfig ¶
type BaselineTrackerConfig struct {
// MaxProducersPerSeries is the per-series producer-baseline bound
// (AGGREGATE_MAX_PRODUCER_BASELINES_PER_SERIES, default 8). Excess
// producers share one degraded baseline and never evict the first N.
MaxProducersPerSeries int
// MaxBaselines is the global baseline-entry budget
// (the resolved AGGREGATE_MAX_BASELINES). Past it, only the per-series
// degraded slot may still be created — it is reserved capacity, exactly
// like the __other__ series, because refusing it would strand the series.
MaxBaselines int
// GapThreshold is the downtime bound: a point more than this after the
// baseline's last timestamp re-seeds instead of crediting the increase.
// Defaults to AllowedLateness.
GapThreshold time.Duration
}
BaselineTrackerConfig bounds the tracker.
type Bounds ¶
type Bounds struct {
// MaxValueBytes caps the encoded length of a non-tenant dictionary value.
MaxValueBytes int
// MaxTenantBytes caps the encoded length of a tenant name.
MaxTenantBytes int
// MaxTenants is the instance-wide tenant-identity cap.
MaxTenants int
// PerTenantKind caps entries per (tenant, kind). A missing or zero entry
// means unlimited. __other__ entries are exempt.
PerTenantKind map[Kind]int
// InstanceKind caps entries per kind across every tenant — the backstop
// behind PerTenantKind. A missing or zero entry means unlimited.
InstanceKind map[Kind]int
}
Bounds is the identity-bound configuration of a Cache and its Registrar (#200 Q3). The zero value takes every default.
type Bucket ¶
type Bucket struct {
WindowStart int64
SeriesID SeriesID
Delta *AggregateDelta
// Source says which table the row came from. A window may hold a
// materialized bucket and a not-yet-finalized delta row for one series;
// both are real contributions and both must be counted.
Source BucketSource
}
Bucket is one store-owned (window, series) row.
type BucketCursor ¶
type BucketCursor struct {
WindowStart int64
SeriesID SeriesID
Source BucketSource
}
BucketCursor is the keyset position of a paged bucket read. Treat it as opaque: obtain it from BucketPage.Next and hand it back through Selector.After.
func (BucketCursor) After ¶
func (c BucketCursor) After(window int64, id SeriesID, src BucketSource) bool
After reports whether row (window, id, src) sorts strictly after the cursor. It is the Go-side twin of the SQL keyset predicate, so an in-memory Store implementation pages identically to the SQLite one.
type BucketPage ¶
type BucketPage struct {
// Buckets are the rows of this page, ordered by (window, series, source).
Buckets []Bucket
// Limit is the row limit that was applied.
Limit int
// Truncated reports that more rows matched than this page returned. It is
// result-completeness metadata and is INDEPENDENT of Coverage: a result
// can be full-coverage and truncated at the same time (#197 Q4).
Truncated bool
// Next resumes the read immediately past the last returned row. Only
// meaningful when Truncated.
Next BucketCursor
}
BucketPage is one page of a ReadBuckets call.
type BucketSource ¶
type BucketSource uint8
BucketSource says which durable table a row came from. It exists so a paged read has a TOTAL order to resume from: (window_start, series_id) is unique within each table but a window can legitimately hold a materialized bucket AND a not-yet-finalized delta row for the same series.
const ( // SourceFinalized is a row from aggregate_buckets. SourceFinalized BucketSource = 0 // SourceDelta is a not-yet-finalized row from aggregate_delta_log. SourceDelta BucketSource = 1 )
BucketSource values, in scan order.
type BucketsResult ¶
type BucketsResult struct {
Points []TrafficPoint
Coverage Coverage
Epoch string
Revision uint64
}
BucketsResult is the answer to QueryBuckets.
type Cache ¶
type Cache struct {
// contains filtered or unexported fields
}
Cache is the hot-path canonical-value to ID map in front of a Registrar. It is safe for concurrent use.
Hits take a read lock and no allocation. Misses call the Registrar without holding any lock — a durable registrar performs I/O, and serializing every ingest goroutine behind one dictionary mutex would be the whole engine's bottleneck. The cost is that two goroutines may register the same value concurrently, which is why Registrar is required to be idempotent.
Entries are never evicted: the series caps of #158 bound how many distinct identities can exist. Overflow (__other__) resolutions are deliberately NOT cached, so a pathological-cardinality tenant cannot grow this map without bound; they pay a Registrar call per point instead.
func NewCacheWithBounds ¶
NewCacheWithBounds returns a Cache in front of reg honouring b (#200 Q3).
func (*Cache) Fence ¶
Fence takes ids out of service for the hot path without touching the maps that hold them. A fenced ID is never returned by Intern; the value resolves to __other__ for the duration, which is the same degradation a full namespace already produces.
Fencing is deliberately NOT deletion: a sweep whose DELETE fails must be able to release the fence and leave memory exactly as it was.
func (*Cache) Forget ¶
Forget removes cached entries whose ID is in ids. It runs after a sweep has COMMITTED, never before: the forward map and the reverse map must not disagree with the database in the window where the delete could still fail.
func (*Cache) Intern ¶
Intern returns the dictionary ID for a string value, registering it on a miss. It never fails: a full or failing dictionary resolves to the pre-created __other__ ID for the scope.
func (*Cache) InternBytes ¶
InternBytes is Intern for a byte-slice value (dimension tuples). The slice is never retained: it is copied when it becomes a map key. The lookup itself does not allocate.
func (*Cache) InternDims ¶
InternDims canonicalizes pairs and interns the encoding as KindDimTuple, returning the DimsID for a SeriesKey. An empty set yields 0, the "no configured dims" sentinel. pairs is sorted in place.
func (*Cache) InternTenant ¶
InternTenant returns the ID for a tenant name, which lives in the instance-global tenant namespace.
It is the ONE namespace that refuses instead of degrading (#200 Q3): an over-length, empty, or over-cap tenant returns ok=false and the caller drops the point. Collapsing two tenants onto one identity would silently merge their telemetry, and no downstream reader could tell.
One consequence worth stating plainly: while the identity-maintenance barrier is fencing a tenant ID, this refuses points for that tenant instead of routing them anywhere. That window is single-digit milliseconds, once a day, and it only ever covers a tenant GC had already determined nothing references. A refused point is the honest answer; an __other__ tenant is not.
func (*Cache) Lookup ¶
Lookup reverses a dictionary ID through the underlying registrar. It returns ok=false when the registrar cannot reverse IDs at all, so a caller that only needs presentation names degrades to "unresolved" instead of failing.
func (*Cache) OtherID ¶
OtherID returns the pre-created overflow ID for (tenantID, kind). Cardinality overflow routing needs it directly, not just as a miss fallback: an overflow series carries the __other__ entry as its NameID.
func (*Cache) Roots ¶
Roots returns every dictionary ID this cache can still hand to the hot path. They are GC roots by construction: a cached ID is one map read away from becoming a series identity, and no lock the collector can take would make that read wait.
func (*Cache) SetOverflowSink ¶
SetOverflowSink installs (or, with nil, removes) the callback that publishes bound-driven __other__ routing. Safe to call while interning is in flight.
func (*Cache) Stats ¶
func (c *Cache) Stats() CacheStats
Stats returns a snapshot of the cache counters.
type CacheStats ¶
type CacheStats struct {
// Hits counts lookups served from the in-memory map.
Hits uint64
// Misses counts lookups that reached the Registrar.
Misses uint64
// Overflows counts misses the Registrar rejected as ErrDictFull, which
// were routed to the __other__ entry.
Overflows uint64
// Errors counts misses the Registrar failed for any other reason. These
// were also routed to __other__; Phase 2 must surface this as a metric.
Errors uint64
// OverLength counts values refused for exceeding the encoded-length cap
// and routed to __other__ (#200 Q3). Never truncated.
OverLength uint64
// TenantsRejected counts tenant identities refused outright — over-length,
// empty, or past the instance-wide tenant cap. These points are DROPPED,
// not collapsed.
TenantsRejected uint64
// Fenced counts lookups that hit a dictionary ID the identity-maintenance
// barrier had fenced, and were therefore served from __other__ (#200 Q2).
Fenced uint64
}
CacheStats is a snapshot of Cache counters.
type Coverage ¶
type Coverage string
Coverage is the honesty vocabulary of a response (#164). It says what the numbers are derived from, so a surface can never imply completeness it does not have.
const ( // CoverageFull means the numbers describe every accepted event. CoverageFull Coverage = "full" // CoverageSampled means part of the response is derived from a subset of // events — exact where the aggregate engine counted, partial elsewhere. CoverageSampled Coverage = "sampled" // CoverageExemplar means the response is built only from retained raw // exemplars. An absent exemplar NEVER implies zero matching events. CoverageExemplar Coverage = "exemplar" )
Coverage values.
type CumulativeOutcome ¶
type CumulativeOutcome struct {
// Delta is the increase attributable to this point. Zero for ignored,
// seeded and gap outcomes.
Delta float64
// Ignored is true for a stale or duplicate point: it did not move the
// baseline and produced no delta (case 1).
Ignored bool
// Seeded is true when this point created a baseline. No delta is
// attributed: the accumulation predates our observation window and
// crediting it to one 5-minute window would fabricate a rate spike.
Seeded bool
// Gap is true when the point arrived more than the allowed-lateness
// window after the baseline's last timestamp. The computed delta is
// discarded and the baseline re-seeded (#166 downtime handling): totals
// under-count across long outages, rates never lie.
Gap bool
// Reset is true when a counter reset was detected; Reason says which.
Reset bool
Reason ResetReason
// Degraded is true when the point was attributed to the per-series shared
// baseline because the producer bound was exhausted.
Degraded bool
// Recovered is true when Delta also carries an increase whose group commit
// failed earlier and which this point re-attributes (#194 blocker 2). It is
// informational: the stranded amount is already folded into Delta.
Recovered bool
}
CumulativeOutcome classifies what happened to one cumulative point. Exactly one of Ignored/Seeded/Gap/Reset/Normal describes the point; Delta carries the increase to attribute to the point's window.
type DashboardResult ¶
type DashboardResult struct {
// RequestCount is accepted request entry points over the range;
// ErrorRequestCount is its error subset and RequestErrorRate is the
// headline dashboard error rate, as a PERCENT.
RequestCount int64
ErrorRequestCount int64
RequestErrorRate float64
// SpanCount is accepted spans; SpanErrorCount is its error subset and
// SpanErrorRate the corresponding PERCENT.
SpanCount int64
SpanErrorCount int64
SpanErrorRate float64
TotalLogs int64
AvgLatencyMs float64
ActiveServices int64
P99LatencyMicros float64
LatencyProvenance latency.Provenance
TopFailing []ServiceStat
Accuracy AccuracyMetadata
Coverage Coverage
Epoch string
Revision uint64
}
DashboardResult is the answer to QueryDashboard.
There is no TotalTraces field: #194 blocker 5 is precisely that the old one carried a SPAN count under a trace name, and a field that has been wrong cannot be fixed by leaving its name in place. Every count here says its basis (#197 Q3), and the headline error rate is the request-basis one.
type DeltaMap ¶
type DeltaMap map[SeriesWindowKey]*AggregateDelta
DeltaMap is one reducer's output: the deltas of one Export request, keyed by series and window.
type DeltaRow ¶
type DeltaRow struct {
// SeriesID and WindowStart identify the bucket the row contributes to.
SeriesID SeriesID
WindowStart int64
// Delta is the aggregate contribution. Its sketch is encoded with the
// versioned codec from #157 on write and decoded on read.
Delta *AggregateDelta
}
DeltaRow is one delta-log row: the accumulated, not-yet-finalized contribution to one (series, window).
On the write side it carries one group commit's contribution, which the store merges into the row already standing for that (series, window) — the log is keyed by identity, not by an append sequence, so a window's row count tracks its active series rather than the number of commits that touched it (#173).
type DictRow ¶
DictRow is one dictionary registration awaiting commit. ID is assigned by the registrar before the row is written, because the hot path needs the ID synchronously; the row itself lands in the same transaction as the first delta that references it (#162's first atomicity invariant).
type DimPair ¶
DimPair is one operator-configured dimension, already reduced to dictionary IDs. The hot path never carries the strings.
type DimsConfig ¶
DimsConfig maps metric names to their aggregation dimension keys. Keys are sorted canonically per ParseAggregateMetricDims; the lookup is exact on metric name and returns the dimension key list or nil if the metric is not configured for aggregation.
func (DimsConfig) ExtractDimensionValues ¶
func (d DimsConfig) ExtractDimensionValues(metricName string, attrs []*commonpb.KeyValue) []string
ExtractDimensionValues extracts the values of configured dimension keys from a slice of OTLP attributes. Returns nil if the metric is not configured or if any configured key is missing from the attributes. Returned values are in the order of the configured keys.
func (DimsConfig) Get ¶
func (d DimsConfig) Get(metricName string) []string
Get returns the sorted dimension keys for a metric, or nil if the metric is not in the config.
type DirtyBaseline ¶
type DirtyBaseline struct {
Key SeriesKey
Producer ProducerID
Baseline Baseline
// Inflight is the increase this record emitted as deltas since its last
// drain. It rides the drain so a failed commit can hand the amount back
// as owed instead of stranding it (#194 blocker 2). It is not part of the
// durable BaselineRow — the store never sees it.
Inflight float64
}
DirtyBaseline is one baseline awaiting durable upsert. The group-commit writer drains them into the same transaction as the deltas they justify (#166), which is what closes the restart gap durable ACK exists to close.
type DurableRegistrar ¶
type DurableRegistrar struct {
// contains filtered or unexported fields
}
DurableRegistrar is the Registrar backed by the aggregate store. It replaces MemRegistrar whenever the store is enabled — the seam in dict.go exists for this and needs no other change.
func NewDurableRegistrar ¶
func NewDurableRegistrar(store Store, limits map[Kind]int) (*DurableRegistrar, error)
NewDurableRegistrar builds a registrar warmed from the store's dictionary so IDs survive restart. limits caps entries per (tenant, kind); a missing or zero entry means unlimited, and __other__ entries are always exempt.
func NewDurableRegistrarWithBounds ¶
func NewDurableRegistrarWithBounds(store Store, b Bounds) (*DurableRegistrar, error)
NewDurableRegistrarWithBounds is NewDurableRegistrar with the full #200 Q3 bound set: encoded-value length caps, per-(tenant, kind) counts, and the instance-wide backstops behind them.
The preload is exact, not truncated. LoadDict is asked for one row more than the supported bound and a full page fails startup: a silent LIMIT is how a registrar comes up believing a value is unregistered, mints a second ID for it, and splits one series into two that no query can reunite.
func (*DurableRegistrar) ClearTouched ¶
func (r *DurableRegistrar) ClearTouched()
ClearTouched resets the "handed out since the last sweep" set without deleting anything. A sweep that collected nothing still has to clear it, or the set grows for the life of the process.
func (*DurableRegistrar) Committed ¶
func (r *DurableRegistrar) Committed(rows []DictRow)
Committed marks drained rows durable.
func (*DurableRegistrar) DrainPending ¶
func (r *DurableRegistrar) DrainPending() []DictRow
DrainPending returns the staged rows for inclusion in the next group commit. They stay staged until Committed confirms them, so a failed commit re-offers them instead of stranding an ID that no row backs.
func (*DurableRegistrar) Fence ¶
func (r *DurableRegistrar) Fence(ids map[uint32]struct{})
Fence takes candidate IDs out of service. Register refuses a fenced ID (the Cache routes the value to __other__) and Lookup reports it unknown, but no map entry moves: a failed DELETE must be able to release the fence and leave memory byte-for-byte as it was.
func (*DurableRegistrar) Forget ¶
func (r *DurableRegistrar) Forget(ids map[uint32]struct{})
Forget removes swept IDs from the forward, reverse and count maps and clears the fence. It runs only after the DELETE has committed.
func (*DurableRegistrar) Lookup ¶
func (r *DurableRegistrar) Lookup(id uint32) (DictEntry, bool)
Lookup implements Resolver. The entry is visible as soon as the ID is minted, before the row is durable: a query that resolves a name the next commit will persist is correct, and one that cannot resolve it at all is not.
func (*DurableRegistrar) Next ¶
func (r *DurableRegistrar) Next() uint32
Next reports the ID the registrar would mint next. It is the value persisted as the dictionary high-watermark.
func (*DurableRegistrar) OtherID ¶
func (r *DurableRegistrar) OtherID(tenantID uint32, kind Kind) uint32
OtherID implements Registrar. The overflow entry bypasses the capacity cap: a quota must never block creation of the entry that absorbs quota violations.
func (*DurableRegistrar) PendingCount ¶
func (r *DurableRegistrar) PendingCount() int
PendingCount returns how many registrations are staged but not yet durable.
func (*DurableRegistrar) Revalidate ¶
func (r *DurableRegistrar) Revalidate(candidates map[uint32]struct{})
Revalidate drops from candidates every ID the registrar can still hand out. It runs INSIDE the maintenance barrier, under the registrar mutex, so an ID a concurrent Register returned either lands in touched before this runs (and survives) or blocks until the sweep has decided (and gets re-minted).
func (*DurableRegistrar) Roots ¶
func (r *DurableRegistrar) Roots() map[uint32]struct{}
Roots returns the dictionary IDs the registrar itself keeps alive: every staged (not yet durable) registration, every pre-created __other__ sentinel, and every ID handed out since the last completed sweep.
The __other__ entries are unconditional roots. They absorb quota violations, so a sweep that collected one because nothing referenced it this hour would force the next overflow to mint a second sentinel for the same namespace.
func (*DurableRegistrar) Unfence ¶
func (r *DurableRegistrar) Unfence()
Unfence releases every fenced ID without changing anything else.
type EdgeInput ¶
type EdgeInput struct {
Tenant string
// Caller is the service that owns the parent span.
Caller string
// Callee is the service that owns this span.
Callee string
// HTTPRoute, URLPath and SpanName resolve the callee's operation for
// route normalization; only the callee service name enters edge identity.
HTTPRoute string
URLPath string
SpanName string
Method string
HTTPStatusCode int
SpanKind int32
StatusCode int32
// Root mirrors SpanInput.Root for the callee's span.
Root bool
Timestamp time.Time
DurationMicros float64
}
EdgeInput is one resolved caller/callee call, derived from a child span whose parent span belongs to a different service. Everything except Caller comes from the child span, matching what the edge measures: the callee's work as observed by this call.
type EdgeResolver ¶
type EdgeResolver struct {
// contains filtered or unexported fields
}
EdgeResolver recovers the caller service of a span from its parent span ID. Safe for concurrent use.
func NewEdgeResolver ¶
func NewEdgeResolver(maxSpans int) *EdgeResolver
NewEdgeResolver builds a resolver whose total span memory is capped at maxSpans across all shards. Zero or negative takes DefaultEdgeResolverSpans.
func (*EdgeResolver) Len ¶
func (r *EdgeResolver) Len() int
Len reports how many span mappings are currently held. Test and diagnostic accessor only.
func (*EdgeResolver) Observe ¶
func (r *EdgeResolver) Observe(tenant, spanID, parentSpanID, service string) (string, bool)
Observe records this span's service and resolves its caller.
It returns the parent's service and true only when the parent span is still remembered AND belongs to a different service — a same-service parent is an internal call, not a topology edge. Recording and lookup touch different shards and never hold two locks at once.
type Engine ¶
type Engine struct {
// contains filtered or unexported fields
}
Engine is the aggregate accounting engine.
func NewEngine ¶
func NewEngine(cfg EngineConfig) (*Engine, error)
NewEngine builds an Engine. It fails on a budget that cannot hold — a misconfigured cap must stop startup, not silently become policy.
func (*Engine) ActiveSeriesKeys ¶
ActiveSeriesKeys returns every series identity the mutable shards currently hold. It is a GC root set (#200 Q1): a series present in memory is live even if the durable tables have not caught up with it yet.
Bounded by AGGREGATE_MAX_SERIES plus the overflow reserve, so the map is thousands of entries, not millions.
func (*Engine) ApplyCommitted ¶
ApplyCommitted admits the deltas against the cardinality budget and folds them into the mutable windows. Phase 2's writer calls it after its COMMIT; Phase 1 calls it inline.
Admission happens first, with no shard lock held, because a series pushed past a cap is rerouted to an __other__ series that may hash to a different shard. Resolving identity before touching any shard is what keeps the rule "never hold two shard locks" trivially true.
func (*Engine) ApplyDeltas ¶
ApplyDeltas applies one reducer's output through the configured applier.
func (*Engine) ApplyDeltasErr ¶
ApplyDeltasErr applies one reducer's output and surfaces any refusal. When the configured applier cannot fail (Phase 1's direct applier) the error is always nil.
func (*Engine) ApplyReducer ¶
ApplyReducer records the reduction metrics for one Export request and applies its deltas. This is the entry point the OTLP servers call.
func (*Engine) ApplyReducerErr ¶
ApplyReducerErr is ApplyReducer for callers that must honour the durable-ACK contract: it returns the applier's refusal so an Export can answer RESOURCE_EXHAUSTED instead of acknowledging telemetry that is not durable.
func (*Engine) Baselines ¶
func (e *Engine) Baselines() *BaselineTracker
Baselines returns the cumulative baseline tracker.
func (*Engine) CommitAdmission ¶
func (e *Engine) CommitAdmission(plan *AdmissionPlan) uint64
CommitAdmission folds a committed plan into the mutable windows and returns the new revision. The reservations become ordinary occupancy, released later by MarkFinalized like any other series in the window.
func (*Engine) EdgeResolver ¶
func (e *Engine) EdgeResolver() *EdgeResolver
EdgeResolver returns the caller-resolution memory used to derive service edges. The OTLP trace receiver feeds every received span through it.
func (*Engine) Epoch ¶
Epoch returns the process generation identifier. It is stable for the life of the engine and pairs with Revision to give consumers a total order across restarts (#164).
func (*Engine) MarkFinalized ¶
MarkFinalized transitions one window from memory ownership to store ownership. The writer calls it after store.FinalizeWindow has committed, so the handover is exactly as atomic as the transaction that materialized the buckets: a query either sees the window in memory or in the store, never in both and never in neither.
It is the ONLY path that may evict a window from the shards, drop its mutable ownership and advance the watermark — with one counted exception, the closed window cap in Rollover (#194 blocker 6).
func (*Engine) MetricDims ¶
func (e *Engine) MetricDims() DimsConfig
MetricDims returns the configured metric dimension keys for aggregation.
func (*Engine) Miner ¶
func (e *Engine) Miner() *TemplateMiner
Miner returns the ingest-owned template miner.
func (*Engine) NewReducer ¶
NewReducer returns a reducer for one Export request. arrival is the single timestamp captured for that request and used to evaluate lateness and future skew for every point in it.
func (*Engine) Ownership ¶
Ownership captures {mutable set, finalized watermark, revision, epoch} in one critical section. Every read path starts here.
func (*Engine) PlanAdmission ¶
func (e *Engine) PlanAdmission(m DeltaMap) *AdmissionPlan
PlanAdmission rolls the mutable window set forward and reserves one batch of deltas against the cardinality budget WITHOUT touching the shards.
The durable writer (#173) calls it before its COMMIT so the row that becomes durable carries the same identity the shards will later hold: admitting after the write would let the store accumulate series the in-memory caps already rerouted to __other__. Exactly one of CommitAdmission or RollbackAdmission must follow.
func (*Engine) PruneTopology ¶
func (e *Engine) PruneTopology()
PruneTopology drops topology windows past the retention horizon. A visible expiry is a new full-replacement state, so it advances the same engine revision used by commits and leaves an empty tenant tombstone for consumers.
func (*Engine) QueryBuckets ¶
QueryBuckets returns per-window traffic counts. It is the traffic-chart query: one point per five-minute window, never one row per series.
The store half is a SQL GROUP BY: the result is one row per window (per service when a service filter forces it), so the 20,000-row read cap cannot reach it. That is #194 blocker 4 for this query class.
func (*Engine) QueryDashboard ¶
QueryDashboard returns the dashboard summary: totals, averages, active services, the p99 from the merged sketch, and the accuracy metadata that sketch justifies.
The store side is read TWICE, on purpose (#197 Q1). The scalar totals come from one SQL aggregation that no row cap can truncate. The p99 cannot: a quantile sketch is not SUMmable, so the sketch-bearing rows are streamed once and merged in Go (#219). Doing both through one capped read is exactly what made the old dashboard quietly wrong past 20,000 rows.
func (*Engine) QueryTopology ¶
QueryTopology returns the service topology: one node per service with its aggregate accounting, plus the caller/callee edges of the SAME tenant and range, read under the SAME ownership snapshot.
Both halves are SQL GROUP BY on the store side — one row per service for the nodes, one row per (caller, callee) for the edges, never one per series — so CoverageFull here means what it says and no row cap can truncate the answer.
A Services filter selects a SUBGRAPH: an edge survives only when both of its ends survive, so the result is never a graph with edges hanging off nodes it does not contain.
func (*Engine) RestoreTopology ¶
func (e *Engine) RestoreTopology(ids map[SeriesWindowKey]topoIdentity, deltas DeltaMap) int
RestoreTopology folds durable rows into the TOPOLOGY PROJECTION ONLY.
It is the read side of the bounded startup exception (#194 finding 15): the mutable shards are untouched, so nothing here can resurrect finalized history as mutable state or double-count a window the delta-log replay already restored. Identities come from reversed dictionary IDs rather than from a reducer, and the projection's own retention cutoff still applies — which is why the folded count, not the row count, is what recovery reports.
func (*Engine) Revision ¶
Revision returns the current revision. It increases by one on every applied batch and never decreases, so a consumer can tell "nothing changed" from "changed back to the same numbers" (#163's replacement-by-revision topology).
func (*Engine) RollbackAdmission ¶
func (e *Engine) RollbackAdmission(plan *AdmissionPlan)
RollbackAdmission releases the occupancy a plan charged after its commit failed. Nothing reached the shards, so there is no window eviction that would ever release it — this is the only undo there is.
func (*Engine) Rollover ¶
Rollover CLOSES every window whose lateness horizon has expired and returns how many windows the closed-window cap forced out of memory — which is loss, and is normally zero.
It does not evict. A closed window refuses new points (Admit drops deltas at or below the cutoff) but keeps its shard contents, its mutable ownership and its place below the watermark until a committed FinalizeWindow hands it over through MarkFinalized. Evicting here — the pre-#194 behaviour — advanced ownership to a store that had not yet materialized aggregate_buckets, so every read of that window skipped the delta log and returned nothing until the next successful finalize, and forever while the finalizer was failing.
The one thing memory cannot do is grow without bound behind a wedged finalizer, so the closed set is capped. Past the cap the oldest closed windows are evicted the old way and counted: lossy, but no longer silent.
func (*Engine) SetApplier ¶
SetApplier replaces the apply path. Phase 2 uses it to insert the group-commit writer between the reducer and the shards.
func (*Engine) SetStore ¶
SetStore wires the durable store into the READ path. The write path already reaches the store through the group-commit writer; this is what lets the query facade serve finalized windows. Call it once at startup, before the engine takes traffic.
func (*Engine) SetTemplateFactSink ¶
func (e *Engine) SetTemplateFactSink(fn func(TemplateFact))
SetTemplateFactSink installs the log-fact consumer on the ingest-owned template miner (#163). GraphRAG performs no mining of its own in shadow and aggregate modes; it consumes these facts. Passing nil detaches the sink.
func (*Engine) TenantID ¶
TenantID resolves a tenant name to its dictionary ID for the read path. ok is false when the tenant identity was REJECTED — over-length, empty, or past the instance-wide tenant cap (#200 Q3). A rejected tenant is never collapsed into a shared identity; the caller reads nothing rather than another tenant's data.
func (*Engine) TopologyEpoch ¶
TopologyEpoch returns this engine instance's epoch.
func (*Engine) TopologyHorizon ¶
TopologyHorizon is how much finalized history the projection retains behind the mutable set. Startup restore reads no further back than this: a window the projection would immediately prune is a window there is no point paying to read.
func (*Engine) TopologyRevision ¶
TopologyRevision returns a tenant's topology revision without rendering a snapshot, so an unchanged tenant costs one map lookup.
func (*Engine) TopologySnapshot ¶
func (e *Engine) TopologySnapshot(tenant string) TopologySnapshot
TopologySnapshot renders one tenant's replacement-by-revision topology.
func (*Engine) TopologyTenants ¶
TopologyTenants returns the tenants the projection currently holds.
type EngineConfig ¶
type EngineConfig struct {
// Mode is the aggregate mode (AGGREGATE_MODE). An Engine is only
// constructed for a non-legacy mode.
Mode string
// Limiter holds the cardinality budget.
Limiter LimiterConfig
// MaxProducerBaselinesPerSeries and MaxBaselines bound the cumulative
// baseline tracker.
MaxProducerBaselinesPerSeries int
MaxBaselines int
// Registrar mints dictionary IDs. Defaults to an in-memory registrar
// whose IDs are provisional and vanish on restart (#173 replaces it).
Registrar Registrar
// Bounds are the identity bounds of #200 Q3: encoded-value length caps,
// per-(tenant, kind) and instance-wide dictionary counts, and the tenant
// cap. The zero value takes every default.
Bounds Bounds
// Miner is the ingest-owned template miner (#163). Defaults to one built
// from the log-template cap.
Miner *TemplateMiner
// Metrics is the Prometheus recorder. nil disables metric recording.
Metrics MetricsRecorder
// MetricDims maps metric names to their configured aggregation dimension keys.
// Nil or empty means no metrics are configured for custom dimensions.
MetricDims DimsConfig
// Topology bounds the per-tenant topology projection GraphRAG consumes in
// aggregate mode (#174). Zero values take the package defaults.
Topology TopologyConfig
// EdgeResolverSpans bounds the span-ID memory used to recover a call
// edge's caller service. Zero takes DefaultEdgeResolverSpans.
EdgeResolverSpans int
// MaxClosedWindows bounds the closed-but-unfinalized windows held in
// memory. Zero takes DefaultMaxClosedWindows; negative is unbounded and
// is for tests only.
MaxClosedWindows int
// Epoch identifies this engine instance. A consumer that sees a new epoch
// knows the revision counter restarted and must replace its state rather
// than reconcile against it. Zero derives one from the clock.
Epoch uint64
// Now is the clock, injectable for tests. Defaults to time.Now.
Now func() time.Time
}
EngineConfig configures an Engine. Zero values take platform defaults.
type ExpBuckets ¶
ExpBuckets is one side (positive or negative) of an OTLP exponential histogram: Counts[i] belongs to bucket index Offset+i.
type ExponentialHistogramInput ¶
type ExponentialHistogramInput struct {
HistogramCommon
Scale int32
ZeroCount uint64
Positive ExpBuckets
Negative ExpBuckets
}
ExponentialHistogramInput is one OTLP ExponentialHistogramDataPoint.
type FailableApplier ¶
type FailableApplier interface {
Applier
// ApplyErr applies the deltas, returning the revision and any refusal.
ApplyErr(DeltaMap) (uint64, error)
}
FailableApplier is an Applier whose apply path can refuse. The durable group-commit writer (#173) implements it: admission can be saturated (ErrSaturated, mapped to RESOURCE_EXHAUSTED / 429) and a COMMIT can fail, and under the durable-ACK contract neither may be acknowledged as success. The Phase 1 direct applier does not implement it, which is what keeps legacy and shadow behaviour identical when no store is wired.
type FieldError ¶
FieldError reports a SeriesKey field whose value is outside the range the enum (or the signal) permits.
func (*FieldError) Error ¶
func (e *FieldError) Error() string
type FinalizeStats ¶
type FinalizeStats struct {
// WindowStart is the finalized window.
WindowStart int64
// Buckets is the number of bucket rows written or merged.
Buckets int
// DeltaRows is the number of delta-log rows incorporated and deleted.
DeltaRows int
// Duration is the wall time of the transaction.
Duration time.Duration
}
FinalizeStats is the outcome of finalizing one window.
type FinalizedPage ¶
type FinalizedPage struct {
// Buckets are the rows, ordered by (window DESC, series). Newest first, so
// a cap drops the OLDEST part of the horizon rather than the freshest.
Buckets []Bucket
// Truncated reports that more rows matched than the cap allowed.
Truncated bool
}
FinalizedPage is one bounded, cross-tenant read of finalized bucket rows.
It carries no resume cursor on purpose: its caller restores a bounded horizon, and a caller that could page past the cap would no longer have a bounded startup. Truncated is the honest end of the answer, not an invitation to ask again.
type GCConfig ¶
type GCConfig struct {
// Store is the durable aggregate store. Must satisfy GCStore.
Store Store
// Registrar owns the dictionary identity maps.
Registrar *DurableRegistrar
// Cache is the hot-path intern cache in front of Registrar.
Cache *Cache
// Miner is the log-template miner whose templates are GC roots.
Miner *TemplateMiner
// Engine supplies the live shard keys.
Engine *Engine
// Series is the writer's series registry.
Series *seriesRegistry
// Barrier serializes the sweep with the commit path.
Barrier Barrier
// Metrics records the pass. nil disables recording.
Metrics StoreMetrics
}
GCConfig configures one collection pass.
type GCSnapshot ¶
type GCSnapshot struct {
// Referenced is every series ID named by aggregate_buckets,
// aggregate_delta_log or aggregate_baseline.
Referenced map[SeriesID]struct{}
// Series, Dict and Templates are the identity tables themselves.
Series []SeriesRow
Dict []DictRow
Templates []TemplateRow
}
GCSnapshot is one consistent read of every identity table, taken inside a single read transaction (#200 Q1).
One transaction, not four queries: a series marked live from the bucket scan and a dictionary row read a moment later have to describe the same instant. Read them separately and a commit landing in between makes a brand-new series invisible to the reference set while its brand-new name is visible to the candidate set — and GC deletes the name of a series that exists.
type GCStats ¶
type GCStats struct {
// SeriesScanned, DictScanned and TemplatesScanned are the row counts the
// mark phase examined.
SeriesScanned, DictScanned, TemplatesScanned int
// SeriesSwept, DictSwept and TemplatesSwept are the rows actually deleted.
SeriesSwept, DictSwept, TemplatesSwept int64
// SeriesRetained and DictRetained are the rows the mark phase kept.
SeriesRetained, DictRetained int
// Revalidated counts candidates the barrier rescued after the lock-free
// scan — a non-zero value is normal under load, a large one means the scan
// is running too long.
Revalidated int
// MarkDuration is the lock-free scan; BarrierDuration is the part that
// serializes with the writer. Only the second belongs in an ACK-latency
// budget.
MarkDuration, BarrierDuration time.Duration
// Duration is the wall time of the whole pass.
Duration time.Duration
}
GCStats is the outcome of one identity garbage-collection pass.
type GCStore ¶
type GCStore interface {
// GCSnapshot reads the reference set and all three identity tables inside
// ONE read transaction. It runs on the READ pool, without the writer
// lock: the full scan is the part of GC that must never become an
// ACK-latency incident.
GCSnapshot() (*GCSnapshot, error)
// LoadTemplates returns every durable log-template row, capped at max.
LoadTemplates(max int) ([]TemplateRow, error)
// SweepIdentities deletes the given series, dictionary and template rows
// in ONE transaction, series first. It runs under the writer lock. A
// partial sweep is never observable: either every row named here is gone
// or none of them are.
SweepIdentities(series []SeriesID, dict []uint32, templates []uint32) (SweepStats, error)
// SaveTemplateStats refreshes the non-identity template counters. Rows
// whose template no longer exists are ignored, not created.
SaveTemplateStats(rows []TemplateStatRow) error
}
GCStore is the Store capability the dictionary/series collector needs. It is separate from Store for the same reason WatermarkStore is.
type GroupBatch ¶
type GroupBatch struct {
Dicts []DictRow
Series []SeriesRow
Deltas []DeltaRow
Baselines []BaselineRow
// Templates are identity-critical log-template mutations: a new template,
// a pattern generalization, or an alias change (#200 Q4). They ride the
// same transaction as the delta that used the resulting identity, because
// a 15-minute snapshot alone lets acknowledged identity state vanish in a
// crash — the bucket would survive naming a template ID that the reloaded
// miner has never heard of.
Templates []TemplateRow
}
GroupBatch is one group commit: everything that must become durable together.
The three atomicity invariants of #162 are structural here — they are not a convention the implementation is asked to honour, they are the shape of the only write entry point:
Dicts + Series + Deltas in one struct -> registration and the first
referencing delta cannot split.
Baselines in the same struct -> a baseline cannot become durable
without the delta it justifies
(and vice versa).
FinalizeWindow's materialize+delete -> the third invariant, one method.
func (*GroupBatch) Empty ¶
func (b *GroupBatch) Empty() bool
Empty reports whether the batch would commit nothing.
type GroupBy ¶
type GroupBy uint8
GroupBy selects the grouping of a SumBuckets aggregation. Zero groups everything into a single row.
const ( // GroupByWindow emits one row per window start. GroupByWindow GroupBy = 1 << iota // GroupByService emits one row per service dictionary ID. GroupByService // GroupByName emits one row per name dictionary ID. The namespace of a // NameID is the signal's (see NameKind), so this flag is only meaningful // on a single-signal selector — joining NameIDs across signals is exactly // the mistake the namespace split exists to prevent. GroupByName // GroupBySignal emits one row per signal. GroupBySignal )
GroupBy flags, combinable.
type HTTPClass ¶
type HTTPClass uint8
HTTPClass is the HTTP status family of an operation. It carries the 4xx/5xx triage split separately from StatusClass because a 4xx is legitimately not a span ERROR.
const ( HTTPClassNone HTTPClass = 0 HTTPClass1xx HTTPClass = 1 HTTPClass2xx HTTPClass = 2 HTTPClass3xx HTTPClass = 3 HTTPClass4xx HTTPClass = 4 HTTPClass5xx HTTPClass = 5 )
HTTPClass values.
func HTTPClassFromStatus ¶
HTTPClassFromStatus maps an HTTP status code onto its family. Codes outside 100..599 yield HTTPClassNone.
type HistogramCommon ¶
type HistogramCommon struct {
Tenant string
Service string
Name string
Resource ResourceIdentity
Timestamp time.Time
StartTime time.Time
Temporality Temporality
Attributes []*commonpb.KeyValue
// ResourceAttributes are the RAW OTLP resource attributes; a configured
// dimension key the point lacks falls back to them (#279).
ResourceAttributes []*commonpb.KeyValue
Count uint64
Sum float64
HasSum bool
Min float64
HasMin bool
Max float64
HasMax bool
}
HistogramCommon is the identity and scalar payload shared by both histogram point shapes. Attributes are the RAW OTLP point attributes: dimension extraction happens in the reducer, against a request-local scratch, so no per-point map is allocated (#199 Q4).
type HistogramFold ¶
type HistogramFold struct {
Sketch *Sketch
Count uint64
Sum float64
HasSum bool
Min float64
HasMin bool
Max float64
HasMax bool
// PercentilesUnavailable suppresses every quantile derived from this
// point, and DropReason says why.
DropReason SketchDropReason
// SourceBucketError is the worst-case relative error contributed by the
// SOURCE histogram's own bucket widths, as a fraction. It is 0 for an
// exponential histogram (index transfer is exact) and can dwarf the
// scale-4 sketch's 2.17% for a coarse explicit-bounds histogram.
SourceBucketError float64
// UnboundedTail reports that observations landed in the +Inf bucket.
// Those observations are deliberately NOT folded into the sketch: their
// only known property is "greater than UnboundedTailBound", and placing
// them at any finite value would turn a lower bound into a fabricated
// estimate. UnboundedTailCount is how many.
UnboundedTail bool
UnboundedTailBound float64
UnboundedTailCount uint64
}
HistogramFold is the result of folding one histogram data point.
Scalars are always populated for an accepted point. Sketch is nil exactly when PercentilesUnavailable is set: the two never disagree, because a caller that saw a non-nil sketch and an "unavailable" flag would eventually publish the sketch.
func FoldExponentialHistogram ¶
func FoldExponentialHistogram(in ExponentialHistogramInput) (HistogramFold, error)
FoldExponentialHistogram converts one OTLP ExponentialHistogramDataPoint into the platform sketch (#199 Q1).
Scale handling:
- s > 4: downscale exactly to 4 by an arithmetic right shift of each bucket index, merging the counts that collapse together. OTel's perfect-subsetting property makes this exact at the bucket-count level.
- s == 4: direct index transfer.
- 0 <= s < 4: the ACCUMULATED sketch is explicitly downscaled to s before any bucket lands, so the result carries the source's coarser mapping honestly. Relying on the sketch's incidental bin-collapse instead would leave the advertised relative error a lie.
- s < 0: not representable by the sketch (its scale is unsigned). Scalars are kept, percentiles are marked unavailable.
zero_count folds into the sketch's zero bucket; count, sum, min and max fold normally. Negative buckets holding observations suppress percentiles for the WHOLE point: publishing the positive side's p99 as the distribution's p99 would be a lie, and dropping the point entirely would lose a legitimate count.
func FoldHistogram ¶
func FoldHistogram(in HistogramInput) (HistogramFold, error)
FoldHistogram converts one OTLP explicit-bounds HistogramDataPoint into the platform sketch (#199 Q2).
Each finite bucket's count enters the sketch ONCE, as a weighted synthetic observation at the bucket's geometric midpoint. The geometric midpoint is only defined for a strictly positive bucket, so:
- A bucket whose lower boundary is below zero can only be folded when the point's own min proves the population is non-negative; the effective lower boundary then becomes max(0, boundary, min). Without that proof the point keeps its scalars and loses its percentiles.
- A bucket whose effective lower boundary is exactly 0 has no geometric midpoint. Its observations are placed at upper/2 and the bucket contributes 100% source error, which is the truth: the observation could have been anywhere in (0, upper].
The +Inf bucket is never folded. Its count and the last finite boundary are carried out separately so a quantile that lands in the tail can be answered as a LOWER BOUND (p99 >= boundary) instead of an invented number.
type HistogramInput ¶
type HistogramInput struct {
HistogramCommon
Bounds []float64
BucketCounts []uint64
}
HistogramInput is one OTLP explicit-bounds HistogramDataPoint.
Bounds are the explicit_bounds array and BucketCounts the bucket_counts array; the OTLP data model requires len(BucketCounts) == len(Bounds)+1, with the last bucket unbounded above.
type Kind ¶
type Kind uint8
Kind names a dictionary namespace. Uniqueness is scoped to (tenant, kind, value), matching the durable UNIQUE(tenant_id, kind, value) constraint the Phase 2 store will carry (#159, #173).
const ( KindTenant Kind = 1 KindService Kind = 2 KindOperation Kind = 3 KindMetricName Kind = 4 KindDimKey Kind = 5 KindDimValue Kind = 6 KindDimTuple Kind = 7 KindLogTemplate Kind = 8 )
Dictionary kinds. Zero is reserved so an un-set Kind is never a valid namespace.
type Limiter ¶
type Limiter struct {
// contains filtered or unexported fields
}
Limiter owns the active-series census and enforces the budget. It is safe for concurrent use.
The engine calls Admit before taking any shard lock and Release during window rollover. The limiter never takes a shard lock, so the lock order is always engine -> limiter and can never invert.
func (*Limiter) Admit ¶
Admit evaluates key against the budget for the given window and returns the series to record under. The full materialized key is evaluated before admission; nothing is ever refused, only redirected.
func (*Limiter) IsOverflowSeries ¶
IsOverflowSeries reports whether key is a live __other__ series.
func (*Limiter) Release ¶
Release drops one series' presence in one window. When the series leaves its last mutable window it stops consuming budget: historical series are free (#158), only active ones are charged.
func (*Limiter) Stats ¶
func (l *Limiter) Stats() LimiterStats
Stats returns a snapshot of occupancy and overflow counters.
type LimiterConfig ¶
type LimiterConfig struct {
MaxSeries int
MaxSeriesMetrics int
MaxSeriesTraces int
MaxSeriesEdges int
MaxSeriesLogs int
MaxSeriesSystem int
MaxOperationsPerService int
MaxLogTemplatesPerService int
MaxTraceSeriesPerService int
MaxMetricSeriesPerService int
// SeriesPerTenantFraction is the fraction of MaxSeries one tenant may
// hold. 0 disables the tenant cap (the single-tenant default).
SeriesPerTenantFraction float64
// OtherNameID resolves the dictionary __other__ ID for the name kind of a
// signal, which is what an overflow series uses as its NameID. Required.
OtherNameID func(tenantID uint32, signal Signal) uint32
}
LimiterConfig holds the cardinality budget. Zero values take the platform defaults; a negative value disables that cap.
func (LimiterConfig) Validate ¶
func (c LimiterConfig) Validate() error
Validate reports a budget that cannot hold: sum of sub-caps above the global cap, or a tenant fraction outside [0,1]. config.Load() performs the same checks at startup (fail-closed); this exists so an engine constructed directly in a test cannot quietly run an impossible budget.
type LimiterStats ¶
type LimiterStats struct {
// Active is the number of distinct budgeted series present in a mutable
// window. __other__ series are excluded: they are the reserve the caps
// spend, not occupancy the caps admit, and counting them here is what let
// the observed census exceed its own sub-cap (#173).
Active int
// ActiveBySignal breaks Active down per signal. Each entry is bounded by
// that signal's sub-cap.
ActiveBySignal map[Signal]int
// Overflow counts admissions routed to an __other__ series, per reason.
Overflow map[OverflowReason]uint64
// OverflowSeries is the number of live __other__ series.
OverflowSeries int
// OverflowSeriesBySignal breaks OverflowSeries down per signal. This is
// the reserve's real occupancy — unbudgeted, bounded by
// (services x signals x status classes).
OverflowSeriesBySignal map[Signal]int
}
LimiterStats is a snapshot of Limiter occupancy.
type LogInput ¶
type LogInput struct {
Tenant string
Service string
// Severity is the severity text; SeverityNumber is the OTLP number and
// wins when non-zero.
Severity string
SeverityNumber int32
// Body is mined into a template. It never enters series identity itself.
Body string
Timestamp time.Time
}
LogInput is one log record.
type MemRegistrar ¶
type MemRegistrar struct {
// contains filtered or unexported fields
}
MemRegistrar is the Phase 1 in-memory Registrar. IDs are provisional: they are minted from a process-local counter and do not survive a restart, which is exactly why nothing may persist a SeriesKey until the durable registrar of #173 lands. It is safe for concurrent use.
func NewMemRegistrar ¶
func NewMemRegistrar(opts *MemRegistrarOptions) *MemRegistrar
NewMemRegistrar returns an in-memory registrar. A nil opts means no limits.
func (*MemRegistrar) Count ¶
func (r *MemRegistrar) Count(tenantID uint32, kind Kind) int
Count returns the number of capacity-consuming entries in a (tenant, kind) namespace, excluding the __other__ entry.
func (*MemRegistrar) Lookup ¶
func (r *MemRegistrar) Lookup(id uint32) (DictEntry, bool)
Lookup resolves an ID back to its entry. Presentation and test use only — the hot path never reverses a dictionary ID.
type MemRegistrarOptions ¶
type MemRegistrarOptions struct {
// Limits caps the number of entries per (tenant, kind). A missing or
// zero entry means unlimited. __other__ entries are exempt.
Limits map[Kind]int
}
MemRegistrarOptions configures a MemRegistrar.
type Method ¶
type Method uint8
Method is the bounded HTTP method enum. It is bounded, not closed: unrecognized methods map to MethodOther so a hostile client cannot mint series identity.
const ( MethodNone Method = 0 MethodGet Method = 1 MethodPost Method = 2 MethodPut Method = 3 MethodDelete Method = 4 MethodPatch Method = 5 MethodHead Method = 6 MethodOptions Method = 7 MethodTrace Method = 8 MethodConnect Method = 9 MethodOther Method = 10 )
Method values.
func LookupMethod ¶
LookupMethod resolves s against the known methods only. It reports ok=false for the empty string and for anything unrecognized; callers that want the bounded-enum degradation should use ParseMethod instead. Matching is case-insensitive so lowercase "get" from sloppy instrumentation still lands on MethodGet.
func ParseMethod ¶
ParseMethod maps an HTTP method string onto the bounded enum. The empty string yields MethodNone; anything unrecognized yields MethodOther.
type MetricInput ¶
type MetricInput struct {
Tenant string
Service string
Name string
Value float64
// Timestamp selects the window; StartTime is the cumulative start time
// used for reset detection.
Timestamp time.Time
StartTime time.Time
// Temporality and Monotonic select the aggregation model (#166).
Temporality Temporality
Monotonic bool
// Resource carries the stable identity used to derive the ProducerID.
Resource ResourceIdentity
// Attributes are the RAW OTLP point attributes. The configured dimension
// tuple is extracted from them here, against a request-local scratch, so
// the hot path never allocates a per-point map (#199 Q4).
Attributes []*commonpb.KeyValue
// ResourceAttributes are the RAW OTLP resource attributes; a configured
// dimension key the point lacks falls back to them (#279).
ResourceAttributes []*commonpb.KeyValue
}
MetricInput is one metric data point.
type MetricPointOutcome ¶
type MetricPointOutcome uint8
MetricPointOutcome is the reducer's verdict on one metric data point.
const ( // MetricPointAccepted means the point contributed to a delta. MetricPointAccepted MetricPointOutcome = iota // MetricPointExcluded means the point fell outside the mutable-window horizon. // It is counted as late or future, NOT as an OTLP rejection: the client // sent well-formed telemetry and retrying would not help. MetricPointExcluded // MetricPointRejectedTemporality means a histogram point arrived with a // temporality the GA engine does not support (#199 Q3). MetricPointRejectedTemporality // MetricPointRejectedMalformed means the point violates the OTLP data model. MetricPointRejectedMalformed )
MetricPointOutcome values.
type MetricPointResult ¶
type MetricPointResult struct {
Outcome MetricPointOutcome
// Reason is the metric label for a rejection, "" when accepted.
Reason string
// Err carries the validation failure behind PointRejectedMalformed.
Err error
// SketchDropped reports that the point's scalars were accepted but its
// percentiles suppressed, and DropReason says why. This is NOT a
// rejection: the point still contributes count, sum, min and max.
SketchDropped bool
DropReason SketchDropReason
}
MetricPointResult reports what the reducer did with one metric data point, so the OTLP Export path can build an honest partial-success response.
func (MetricPointResult) Rejected ¶
func (r MetricPointResult) Rejected() bool
Rejected reports whether the point must be counted in ExportMetricsPartialSuccess.rejected_data_points.
type MetricsRecorder ¶
type MetricsRecorder interface {
// RecordReduction publishes one Export request's reduction accounting:
// input points, emitted deltas, the reduction ratio, late/future
// exclusions, and the shadow-comparison counters.
RecordReduction(stats ReducerStats, deltas map[Signal]uint64)
// RecordOverflow counts one admission rerouted to an __other__ series.
RecordOverflow(signal Signal, reason OverflowReason)
// SetActiveSeries publishes the budgeted active-series census and the
// __other__ reserve occupancy, both per signal. They are separate numbers
// because only the first is bounded by the AGGREGATE_MAX_SERIES* caps.
SetActiveSeries(active, overflow map[Signal]int)
// SetClosedWindows publishes how many closed-but-unfinalized windows the
// engine is holding. It is the finalizer-backlog signal: a value that does
// not return to zero means finalization is failing.
SetClosedWindows(held int)
// RecordClosedWindowEvicted counts one closed window the cap forced out of
// memory before it was finalized. Every one is lost data.
RecordClosedWindowEvicted()
// RecordTenantRejected counts one point dropped because its tenant
// identity was refused (#200 Q3). Unlike every other identity cap this is
// a DROP, so it deserves its own counter rather than a reason label on
// the overflow one.
RecordTenantRejected(signal Signal)
// RecordIdentityOverflow counts one identity routed to __other__ by a
// #200 Q3 bound, labeled by dictionary kind and the bound that tripped
// ("length" or "count").
RecordIdentityOverflow(kind Kind, bound string)
}
MetricsRecorder is the engine's view of the metric surface.
func NewPrometheusRecorder ¶
func NewPrometheusRecorder(m *telemetry.Metrics) MetricsRecorder
NewPrometheusRecorder returns a recorder backed by the platform metrics. A nil *telemetry.Metrics yields a no-op recorder rather than a panic.
type OverflowReason ¶
type OverflowReason uint8
OverflowReason names the cap that forced a series into its __other__ series. The strings are metric label values and are part of the contract.
const ( OverflowNone OverflowReason = iota // OverflowTenant — the tenant's fraction of the global budget is full. OverflowTenant // OverflowServiceNames — the service has too many distinct names // (operations for traces, log templates for logs). OverflowServiceNames // OverflowServiceSeries — the service has too many materialized series. OverflowServiceSeries // OverflowSignal — the signal's sub-cap is full. OverflowSignal // OverflowGlobal — the instance-wide backstop is full. OverflowGlobal )
OverflowReason values, in enforcement order.
func (OverflowReason) String ¶
func (r OverflowReason) String() string
String implements fmt.Stringer.
type Ownership ¶
type Ownership struct {
// Epoch is the process generation.
Epoch string
// Revision is the engine revision at capture time.
Revision uint64
// Mutable is the memory-owned window starts, oldest first.
Mutable []int64
// FinalizedWatermark is the newest store-owned window start. Zero means
// nothing has been handed over yet.
FinalizedWatermark int64
}
Ownership is an atomic snapshot of {mutable set, finalized watermark, revision, epoch}. A query captures one and reads every window through it.
func (Ownership) OwnsInMemory ¶
OwnsInMemory reports whether windowStart is memory-owned in this snapshot.
type PointDisposition ¶
type PointDisposition uint8
PointDisposition classifies a point against the mutable-window horizon.
const ( // PointAccepted — the point belongs to a mutable window. PointAccepted PointDisposition = iota // PointLate — the point is older than the lateness horizon. It is excluded // from aggregates and counted; the raw/exemplar path still sees it, because // a late error trace is still evidence (#160). PointLate // PointFuture — the point is timestamped beyond the tolerated skew. PointFuture )
PointDisposition values. The non-accepted strings are metric label values.
func Classify ¶
func Classify(arrival, pointTime time.Time) (int64, PointDisposition)
Classify places a point in a window relative to one Export's arrival time. A single arrivalTime is captured per Export request and used for every point in it (#160): lateness must not depend on where in the batch a point sits.
func (PointDisposition) String ¶
func (d PointDisposition) String() string
String implements fmt.Stringer.
type PreloadError ¶
PreloadError reports a warm-up load whose retained row count exceeds the bound this build supports. It is fatal at startup by design: the alternative is a silently truncated identity map.
func (*PreloadError) Error ¶
func (e *PreloadError) Error() string
type ProducerID ¶
type ProducerID uint64
ProducerID discriminates the concrete emitter of a cumulative series — an instance, pod or process. It is internal baseline state only: it never enters a SeriesKey and is never an aggregate dimension (#159's allowlist, #166's producer keying). Two producers sharing one canonical series would otherwise look like perpetual resets.
func ResolveProducerID ¶
func ResolveProducerID(id ResourceIdentity) ProducerID
ResolveProducerID returns the producer discriminator for a resource. service.instance.id wins when present; otherwise the ID is a deterministic FNV-1a hash of the stable tuple {service.namespace, service.name, host.id|host.name, k8s.pod.uid|container.id|process.pid}.
The hash is never zero: zero is the degraded shared slot.
type PurgeStats ¶
type PurgeStats struct {
// Buckets and Deltas count deleted rows.
Buckets int64
Deltas int64
// Baselines counts baselines dropped because their series has no
// remaining data.
Baselines int64
// Duration is the wall time of the whole purge.
Duration time.Duration
}
PurgeStats is the outcome of one retention purge.
type Query ¶
type Query struct {
// Tenant is the tenant name. Required.
Tenant string
// Start and End bound the read. Required, Start before End.
Start, End time.Time
// Services, when non-empty, restricts the read to these service names.
Services []string
// Signal, when non-zero, restricts the read to one signal.
Signal Signal
}
Query bounds one read. Tenant, Start and End are mandatory; the rest narrow the scan.
type RecoverOptions ¶
type RecoverOptions struct {
// TopologyHorizon is how much FINALIZED history is rebuilt into the
// engine's topology projection. Zero disables the restore; anything past
// the projection's own horizon is clamped down to it, because a window the
// projection would prune on arrival is not worth reading.
TopologyHorizon time.Duration
// TopologyMaxRows caps the finalized rows one restore may read. Zero takes
// DefaultTopologyRestoreMaxRows. It is the second bound on startup cost:
// the horizon says how far back, this says how much.
TopologyMaxRows int
}
RecoverOptions configures the parts of recovery that are policy rather than correctness. The zero value is the pre-#194 behaviour: replay only.
type RecoveryGate ¶
type RecoveryGate struct {
// contains filtered or unexported fields
}
RecoveryGate reports whether startup recovery has completed. The readiness probe holds /ready at 503 until Done() is true.
func NewRecoveryGate ¶
func NewRecoveryGate() *RecoveryGate
NewRecoveryGate returns a gate that starts closed.
func (*RecoveryGate) Done ¶
func (g *RecoveryGate) Done() bool
Done reports whether recovery has completed.
type RecoveryStats ¶
type RecoveryStats struct {
// FinalizedWindows counts windows finalized because their lateness
// horizon expired during downtime.
FinalizedWindows int
// ReplayedRows and ReplayedSeries count delta-log rows read back and the
// distinct (series, window) pairs they folded into.
ReplayedRows int
ReplayedSeries int
// SeededBaselines counts cumulative baselines restored.
SeededBaselines int
// RestoredTopologyRows counts the durable rows the topology restore read:
// finalized bucket rows inside the horizon plus the replayed mutable rows,
// which are in the shards but not in the projection.
RestoredTopologyRows int
// RestoredTopologyWindows counts the (series, window) pairs that actually
// LANDED in the projection. It is lower than RestoredTopologyRows whenever
// the projection's cutoff or one of its caps refused a row.
RestoredTopologyWindows int
// RestoredTopologyTruncated reports that the row cap cut the horizon
// short. The restored topology is real but incomplete at its oldest end.
RestoredTopologyTruncated bool
// SkippedSeries counts delta rows whose series id no longer resolves —
// structurally impossible while registration and first delta share a
// transaction, so a non-zero value here is a corruption signal.
SkippedSeries int
// Duration is the wall time of the whole recovery.
Duration time.Duration
}
RecoveryStats is the outcome of one startup recovery.
func Recover ¶
func Recover(store Store, engine *Engine, w *Writer, now time.Time, opts RecoverOptions) (RecoveryStats, error)
Recover replays the durable store into the engine. It must run before the writer starts accepting Exports and before readiness flips.
type Reducer ¶
type Reducer struct {
// contains filtered or unexported fields
}
Reducer collapses one Export request into deltas.
func (*Reducer) Deltas ¶
Deltas returns the reduced deltas. The map is handed to the engine as-is; the reducer must not be used afterwards.
func (*Reducer) MergeFrom ¶
MergeFrom folds another reducer's deltas and stats into r. Used by Export paths that reduce resource batches in parallel.
func (*Reducer) ReduceEdge ¶
ReduceEdge folds one resolved cross-service call into its service-edge series.
#183 shipped SignalServiceEdge but emitted nothing into it: a single span does not know its caller. #174 supplies the caller through the engine's EdgeResolver and this is where it lands. Edge identity is (caller service, callee service) plus the callee's status/HTTP/kind dimensions — never a span ID, never an operation of the caller.
func (*Reducer) ReduceExponentialHistogramPoint ¶
func (r *Reducer) ReduceExponentialHistogramPoint(in ExponentialHistogramInput) MetricPointResult
ReduceExponentialHistogramPoint folds one OTLP ExponentialHistogram data point (#199 Q1).
func (*Reducer) ReduceHistogramPoint ¶
func (r *Reducer) ReduceHistogramPoint(in HistogramInput) MetricPointResult
ReduceHistogramPoint folds one OTLP explicit-bounds Histogram data point (#199 Q2).
func (*Reducer) ReduceLog ¶
ReduceLog folds one log record into its log series. The template is mined synchronously by the ingest-owned miner (#163); the template ID is the series' NameID, resolved through the log_template dictionary namespace.
func (*Reducer) ReduceMetricPoint ¶
func (r *Reducer) ReduceMetricPoint(in MetricInput)
ReduceMetricPoint folds one metric data point into its metric series, applying the #166 aggregation model for its temporality and monotonicity.
func (*Reducer) ReduceSpan ¶
ReduceSpan folds one span into its trace-operation series.
func (*Reducer) Stats ¶
func (r *Reducer) Stats() ReducerStats
Stats returns the reduction accounting.
type ReducerStats ¶
type ReducerStats struct {
// InputPoints counts points offered to the reducer, per signal, including
// the ones excluded as late or future.
InputPoints [signalCount]uint64
// LatePoints counts points older than the lateness horizon. They are
// excluded from aggregates and counted here — never dropped silently
// (#160). The raw/exemplar path still sees them.
LatePoints [signalCount]uint64
// FuturePoints counts points timestamped beyond the tolerated skew.
FuturePoints [signalCount]uint64
// Accepted counts points that contributed to a delta. This is the
// shadow-comparison numerator: it must be identical for the same input
// stream at any sampling rate.
Accepted [signalCount]uint64
// StaleCumulative counts cumulative points ignored as stale or duplicate
// against their baseline (#166 case 1).
StaleCumulative uint64
// DimsRejected counts metric points whose configured dimension tuple was
// refused from series identity because an attribute value had no scalar
// rendering (#199 Q4). The point is still aggregated, under DimsID 0.
DimsRejected uint64
// TenantsRejected counts points DROPPED because their tenant identity was
// refused — over-length, empty, or past the instance-wide tenant cap
// (#200 Q3). Every other namespace degrades into __other__; the tenant
// namespace refuses instead, because a shared overflow tenant is exactly
// the cross-tenant merge the cap exists to prevent.
TenantsRejected uint64
// ErrorsByService counts errors per service — the cheap invariant #165
// asks for on the aggregate side, and nothing more expensive.
ErrorsByService map[string]uint64
}
ReducerStats is one Export request's reduction accounting.
type Registrar ¶
type Registrar interface {
// Register returns the ID for (tenantID, kind, value), minting one if the
// value is new. It returns ErrDictFull when the namespace is at capacity.
Register(tenantID uint32, kind Kind, value []byte) (uint32, error)
// OtherID returns the pre-created overflow ID for (tenantID, kind). It
// never fails: the entry is created outside the capacity cap.
OtherID(tenantID uint32, kind Kind) uint32
}
Registrar mints dictionary IDs. Phase 1 ships MemRegistrar, whose IDs are provisional and vanish on restart; Phase 2 (#173) plugs in a SQLite-backed registrar that mints durable IDs atomically with the first delta referencing them, without any other change to this package.
Implementations MUST be safe for concurrent use, MUST be idempotent (a repeat Register for the same (tenant, kind, value) returns the same ID, because the Cache deliberately allows concurrent duplicate registrations rather than holding a lock across the call), MUST return IDs greater than zero, and MUST NOT retain the value slice after returning.
type ResetReason ¶
type ResetReason uint8
ResetReason names why a counter reset was recorded.
const ( ResetNone ResetReason = iota // ResetStartTimeChange — the producer reported a new start_time. ResetStartTimeChange // ResetValueRegression — same start_time, value went backwards. ResetValueRegression )
ResetReason values. The strings double as metric label values.
type Resolver ¶
type Resolver interface {
// Lookup returns the entry for id. ok is false when the ID is unknown.
Lookup(id uint32) (DictEntry, bool)
}
Resolver reverses a dictionary ID back to its entry. It is a READ-PATH capability: the ingest hot path never reverses an ID, but the query facade has to turn a SeriesKey back into service, operation and metric names.
It is deliberately separate from Registrar so an implementation that cannot reverse (a write-only registrar) stays usable; the Cache degrades to "unresolved" rather than failing the query.
type ResourceIdentity ¶
type ResourceIdentity struct {
// ServiceInstanceID is service.instance.id. When present it IS the
// producer identity — the semantic convention exists for exactly this.
ServiceInstanceID string
// ServiceNamespace is service.namespace.
ServiceNamespace string
// ServiceName is service.name.
ServiceName string
// Host is host.id, else host.name.
Host string
// Workload is k8s.pod.uid, else container.id, else process.pid.
Workload string
}
ResourceIdentity is the stable slice of resource attributes used to derive a ProducerID. Every field is optional; the first present value per slot is used. This is deliberately NOT the full resource attribute set: attributes that vary per export would fragment baselines into uselessness.
type SQLiteStore ¶
type SQLiteStore struct {
// contains filtered or unexported fields
}
SQLiteStore is the Store implementation on a dedicated SQLite file.
func OpenSQLiteStore ¶
func OpenSQLiteStore(cfg StoreConfig) (*SQLiteStore, error)
OpenSQLiteStore opens (creating if absent) the aggregate database.
Schema handling per #162: absent schema is created at the current version; a partial schema, missing meta, or any version mismatch fails startup with an explicit error. There are no automatic migrations in v1.
func (*SQLiteStore) Analyze ¶
func (s *SQLiteStore) Analyze() error
Analyze refreshes the planner statistics. It is the only maintenance the aggregate file gets: #162 excludes routine VACUUM outright, and ANALYZE is cheap enough to ride the existing daily maintenance tick.
func (*SQLiteStore) Backlog ¶
func (s *SQLiteStore) Backlog() (BacklogStats, error)
Backlog implements Store.
func (*SQLiteStore) CommitGroup ¶
func (s *SQLiteStore) CommitGroup(b *GroupBatch) error
CommitGroup implements Store. One transaction carries the registrations, the pre-merged delta rows and the baseline upserts, so none of #162's three atomicity invariants can be violated by a partial write.
func (*SQLiteStore) FinalizableWindows ¶
func (s *SQLiteStore) FinalizableWindows(cutoff int64, limit int) ([]int64, error)
FinalizableWindows implements Store.
func (*SQLiteStore) FinalizeWindow ¶
func (s *SQLiteStore) FinalizeWindow(windowStart int64) (FinalizeStats, error)
FinalizeWindow implements Store: it materializes the window's buckets and deletes exactly the delta rows it incorporated, in one transaction.
"Exactly the incorporated rows" is a structural property of the predicate, not of the scheduling: the SELECT that materializes and the DELETE that clears the log carry the identical `window_start = ?` predicate inside one transaction, with the writer lock held so nothing can add to the window in between. The earlier sequence bound existed because the log was append-only; with one row per (window, series) there is no partial-append to fence off.
The lock hold is O(active series in the window). Before the delta log was pre-merged this was O(commits x dirty series) — ~1.6M rows and 16 s of blocked ingestion every five minutes at 10k pts/s (#173).
The delta rows are still streamed from the READ pool while the writes go to the writer transaction: bounded is not the same as small, and there is no reason to hold the whole window in memory to satisfy one transaction.
func (*SQLiteStore) GCSnapshot ¶
func (s *SQLiteStore) GCSnapshot() (*GCSnapshot, error)
GCSnapshot implements GCStore. Every scan runs inside ONE deferred read transaction, which in WAL mode pins a single snapshot of the file for its whole life — the reference set and the identity tables therefore describe the same instant. No writer lock is taken: a commit may land during the scan, and the barrier's revalidation is what accounts for it.
func (*SQLiteStore) LoadBaselines ¶
func (s *SQLiteStore) LoadBaselines(max int) ([]BaselineRow, error)
LoadBaselines implements Store.
func (*SQLiteStore) LoadDict ¶
func (s *SQLiteStore) LoadDict(max int) ([]DictRow, error)
LoadDict implements Store.
func (*SQLiteStore) LoadSeries ¶
func (s *SQLiteStore) LoadSeries(max int) ([]SeriesRow, error)
LoadSeries implements Store.
func (*SQLiteStore) LoadTemplates ¶
func (s *SQLiteStore) LoadTemplates(max int) ([]TemplateRow, error)
LoadTemplates implements GCStore.
func (*SQLiteStore) PingContext ¶
func (s *SQLiteStore) PingContext(ctx context.Context) error
PingContext verifies the aggregate database is reachable on the READ pool.
The read pool, never the writer: the writer is MaxOpenConns(1) behind the group commit, so a ping issued there would queue behind whatever commit is in flight and report "unreachable" for a database that is merely busy. A readiness probe has to answer "can this process still read its aggregates", which is exactly what the read pool answers.
func (*SQLiteStore) PurgeBefore ¶
func (s *SQLiteStore) PurgeBefore(cutoff int64) (PurgeStats, error)
PurgeBefore implements Store. Deletion is per-window so retention never opens an unbounded transaction, and there is NO VACUUM: rolling retention on this file relies on range deletes, WAL checkpointing and free-page reuse (#162).
func (*SQLiteStore) ReadBuckets ¶
func (s *SQLiteStore) ReadBuckets(ctx context.Context, sel Selector) (BucketPage, error)
ReadBuckets implements Store. The selector's bounds are mandatory and the row cap is enforced here, not by the caller: the worst case is 6,000 series x 2,016 windows and a dashboard is not trusted with it.
The cap is no longer allowed to be silent (#194 blocker 4): the query asks for limit+1 rows and the extra row, if it exists, becomes Truncated plus a resume cursor. A caller that needs completeness pages with Selector.After until Truncated is false; a caller that wants a scalar total should not be here at all and should call SumBuckets.
func (*SQLiteStore) ReadFinalizedSince ¶
func (s *SQLiteStore) ReadFinalizedSince(since int64, signals []Signal, limit int) (FinalizedPage, error)
ReadFinalizedSince implements Store. It reads MATERIALIZED bucket rows only: the delta log holds the mutable set, which recovery replays through its own path, so reading both here would fold the same contribution twice.
Newest window first. The cap is the store's own row cap, applied with the limit+1 probe every read in this file uses, so a caller learns that its horizon was cut instead of silently receiving a partial service map.
func (*SQLiteStore) ReplayMutable ¶
func (s *SQLiteStore) ReplayMutable(since int64) ([]DeltaRow, error)
ReplayMutable implements Store. Only mutable windows are returned: finalized history never hydrates into RAM (#160).
func (*SQLiteStore) ResolveSeries ¶
func (s *SQLiteStore) ResolveSeries(ids []SeriesID) ([]SeriesInfo, error)
ResolveSeries implements Store. The input count is capped at MaxReadRows.
func (*SQLiteStore) SaveTemplateStats ¶
func (s *SQLiteStore) SaveTemplateStats(rows []TemplateStatRow) error
SaveTemplateStats implements GCStore. Statistics only: it never creates a row, so a template swept between the dirty mark and this write stays swept.
func (*SQLiteStore) SumBuckets ¶
SumBuckets implements Store. It is the scalar-totals path of the #197 read contract: the database does the SUM/COUNT and returns one row per group, so there is no row cap to truncate and no arithmetic for the caller to get wrong on a partial page.
The result size is bounded by the GROUPING — at most (windows x services x signals), all three of which are already bounded by the retention horizon and the cardinality limiter — and never by the number of rows scanned. That is the structural difference from ReadBuckets, whose result size IS the row count.
func (*SQLiteStore) SweepIdentities ¶
func (s *SQLiteStore) SweepIdentities(series []SeriesID, dict []uint32, templates []uint32) (SweepStats, error)
SweepIdentities implements GCStore: one transaction, series first, then the dictionary rows the swept series released, then the template rows those dictionary IDs backed.
The order matters for exactly one reason: a crash between two transactions could leave a series row naming a deleted dictionary ID. Inside one transaction there is no such window, and the ordering is kept anyway so the intent survives a future implementation that cannot span tables.
func (*SQLiteStore) UUID ¶
func (s *SQLiteStore) UUID() string
UUID returns the store's identity, minted at creation. It is what tells an operator "this is a different store", not "this is the same store rebuilt".
func (*SQLiteStore) VisitSketches ¶
func (s *SQLiteStore) VisitSketches(ctx context.Context, sel Selector, visit func(serviceID uint32, sk *Sketch) error) error
VisitSketches implements Store. Both durable tables are streamed once, in their natural PRIMARY KEY (window_start, series_id) range order — no ORDER BY, no LIMIT, no keyset resume. ReadBuckets must maintain a total (window, series, source) order across a UNION of both tables, which forces a sort per page and a re-scan per keyset resume; a wide dashboard range paid that price hundreds of times over (#219). A sketch merge is order-independent, so this path pays it zero times.
func (*SQLiteStore) Watermarks ¶
func (s *SQLiteStore) Watermarks() (uint32, SeriesID, error)
Watermarks implements WatermarkStore.
type SaturationError ¶
type SaturationError struct {
// Bound is "bytes", "waiters" or "deltas".
Bound string
// Current is the value at the moment of refusal; Limit is the cap.
Current, Limit int64
}
SaturationError reports which admission bound refused a CommitGroup.
func (*SaturationError) Error ¶
func (e *SaturationError) Error() string
func (*SaturationError) Is ¶
func (e *SaturationError) Is(target error) bool
Is makes errors.Is(err, ErrSaturated) true for every saturation refusal.
type SchemaError ¶
type SchemaError struct {
// Reason is a short machine-ish description ("missing_meta",
// "partial_schema", "version_mismatch").
Reason string
// Key names the meta key when Reason is "version_mismatch".
Key string
// Got and Want are the mismatched values.
Got, Want string
// Detail carries anything else worth printing.
Detail string
}
SchemaError reports a database whose aggregate schema this build cannot use: a partial schema, missing meta, or a version mismatch. It is fatal at startup by design — silently adopting an unknown layout is how a store starts returning numbers nobody can explain.
func (*SchemaError) Error ¶
func (e *SchemaError) Error() string
type Selector ¶
type Selector struct {
// TenantID scopes the read. Required.
TenantID uint32
// Start and End bound the window range, inclusive of Start and exclusive
// of End, in UTC-aligned Unix seconds. Required, Start < End.
Start, End int64
// Signal, when non-zero, restricts the read to one signal.
Signal Signal
// Signals, when non-empty, restricts the read to this set of signals. It
// is the multi-signal form of Signal for a query that folds several
// signals in one scan (the dashboard reads trace ops and logs); the two
// are mutually exclusive.
Signals []Signal
// SeriesIDs, when non-empty, restricts the read to these series.
SeriesIDs []SeriesID
// Limit caps returned rows. Zero takes MaxReadRows; anything above it is
// clamped down, never up.
Limit int
// After resumes a paged read immediately past this cursor. The zero value
// starts at the beginning of the range.
After BucketCursor
// SketchOnly restricts the read to rows that carry a sketch. It is the
// percentile path's filter: a row with no sketch cannot move a quantile,
// so paging past it is wasted work.
SketchOnly bool
}
Selector bounds a bucket read. Both window bounds and a tenant are mandatory: an unbounded scan of seven days of history is not a query, it is an outage (#162).
type SeriesID ¶
type SeriesID int64
SeriesID is the durable identity of a series. IDs are minted by the database (ADR 0001) so a recovered bucket always resolves to an identity that exists.
type SeriesInfo ¶
SeriesInfo is a series' durable identity, resolved back from its ID.
type SeriesKey ¶
type SeriesKey struct {
TenantID uint32
ServiceID uint32
NameID uint32
DimsID uint32
Signal Signal
StatusClass StatusClass
HTTPClass HTTPClass
Method Method
Variant Variant
}
SeriesKey is the complete identity of one aggregate series. Every field is a dictionary ID or a bounded enum, which makes the struct comparable and therefore directly usable as a Go map key with zero collision risk.
DimsID is 0 when no operator-configured dimensions apply. NameID is resolved through the dictionary kind named by NameKind(Signal).
func DecodeSeriesKey ¶
DecodeSeriesKey decodes exactly one encoded SeriesKey. It rejects short buffers (ErrKeyTruncated), long buffers (ErrKeyTrailingBytes), unknown versions (*VersionError) and out-of-range enum values (*FieldError).
func (SeriesKey) AppendBinary ¶
AppendBinary appends the field-wise encoding of k to dst and returns the extended slice. It never allocates when dst has EncodedSeriesKeyLen bytes of spare capacity. It implements encoding.BinaryAppender.
func (SeriesKey) MarshalBinary ¶
MarshalBinary implements encoding.BinaryMarshaler. It never returns an error; validation is the decoder's job so a caller can round-trip a key it built itself without paying for the check twice.
func (SeriesKey) String ¶
String renders a key for logs and test failures. It is diagnostic output, not a wire format — never parse it.
func (*SeriesKey) UnmarshalBinary ¶
UnmarshalBinary implements encoding.BinaryUnmarshaler. b must be exactly EncodedSeriesKeyLen bytes.
type SeriesRow ¶
SeriesRow is one series registration awaiting commit. Identity columns are explicit — there is no generic "variant" column (#162).
type SeriesWindowKey ¶
SeriesWindowKey identifies one series inside one mutable window. Window starts are UTC-aligned Unix seconds.
type ServiceStat ¶
type ServiceStat struct {
Service string
Count int64
ErrorCount int64
ErrorRate float64
AvgLatencyMs float64
P99LatencyMicros float64
LatencyProvenance latency.Provenance
RequestCount int64
ErrorRequestCount int64
}
ServiceStat is one service's aggregate accounting over the queried range.
Count/ErrorCount/ErrorRate stay SPAN-based: per-service and per-operation diagnostics are about the work done, not about how many requests entered. RequestCount/ErrorRequestCount are carried alongside for the surfaces that want the entry-point basis.
type Signal ¶
type Signal uint8
Signal names the telemetry stream a series belongs to. It also selects the dictionary namespace that SeriesKey.NameID is resolved through; readers must never join NameIDs across namespaces.
const ( SignalUnspecified Signal = 0 // SignalTraceOp is a per-operation trace series. NameID resolves through // KindOperation. SignalTraceOp Signal = 1 // SignalServiceEdge is a caller/callee service edge series. NameID // resolves through KindOperation. SignalServiceEdge Signal = 2 // SignalLog is a log series. NameID resolves through KindLogTemplate. SignalLog Signal = 3 // SignalMetric is a native metric series. NameID resolves through // KindMetricName. SignalMetric Signal = 4 )
Signal values. Zero is reserved for "unspecified" so a zero-valued SeriesKey is never mistaken for a real series.
type Sketch ¶
type Sketch struct {
// contains filtered or unexported fields
}
Sketch is a fixed-size relative-error quantile sketch for latency values.
The zero value is not usable; construct with NewSketch or NewSketchAtScale. A Sketch contains no pointers and no slices: it is copyable by assignment and Observe performs no allocation.
Sketch is not safe for concurrent use. Callers own the synchronization.
func DecodeSketch ¶
DecodeSketch parses a serialized sketch. Unknown versions and scales are rejected with ErrSketchVersion and ErrSketchScale; anything structurally invalid is rejected with ErrSketchTruncated or ErrSketchCorrupt. A decoded sketch re-encodes to exactly the input bytes.
The saturation counter is not part of the format and decodes to zero: it is an operational signal about a live sketch, not part of its value.
func NewSketch ¶
func NewSketch() *Sketch
NewSketch returns an empty sketch at the platform default scale.
func NewSketchAtScale ¶
NewSketchAtScale returns an empty sketch at an explicit scale. Scales above SketchMaxScale are rejected.
func NewSketchAtScaleUnchecked ¶
NewSketchAtScaleUnchecked returns an empty sketch at scale, clamping an out-of-range scale to the platform default. It exists for the read path, which builds a merge target from an already-validated sketch and has no error to return.
func (*Sketch) AppendTo ¶
AppendTo appends the serialized form of the sketch to dst and returns the extended buffer, allowing callers to reuse a scratch buffer.
The sketch is never mutated: if it holds more than SketchMaxSerializedBins populated bins, a copy is collapsed-lowest to the cap and that copy is written. Totals survive the collapse; only resolution below the new floor is lost, and the collapsed flag records it.
func (*Sketch) Collapsed ¶
Collapsed reports whether counts were merged into the lowest retained bin, either by range overflow or by the serialized bin cap.
func (*Sketch) Count ¶
Count returns the total number of observations, including zero-bucket values.
func (*Sketch) Downscale ¶
Downscale converts the sketch to a coarser scale. The conversion is exact: every bucket at the current scale maps onto exactly one bucket at the target scale via an arithmetic right shift of the index, so no count crosses a boundary. The relative-error bound becomes that of the coarser scale.
Downscaling to the current scale is a no-op; a finer target is rejected.
func (*Sketch) Merge ¶
Merge folds other into s. Totals add, bins add bin-wise, and the result is independent of the order in which sketches are merged.
Scales are aligned before merging by downscaling the finer sketch, which is exact; the result carries the coarser of the two scales. Merging a nil or empty sketch is a no-op.
func (*Sketch) Observe ¶
Observe records one value. Non-finite values are rejected and do not affect any counter. Values at or below zero go to the zero bucket: latency is non-negative by construction, so a negative value is a caller bug that must not be allowed to poison the mapping.
Observe never allocates.
func (*Sketch) ObserveBucket ¶
ObserveBucket records n observations already resolved to a bucket index in THIS sketch's mapping. It is the exact-transfer primitive for OTLP exponential histograms, whose bucket indexes use the identical (base^i, base^(i+1)] convention: no value is reconstructed, so no mapping error is introduced beyond the source histogram's own bucket width.
The sketch's running sum is advanced by the bucket's representative value times n. That makes Sum() an ESTIMATE for bucket-folded observations; the authoritative sum of a histogram point is carried separately on the delta (AggregateDelta.HistogramSum) and never derived from here.
func (*Sketch) ObserveN ¶
ObserveN records n observations of the same value. It is the weighted form of Observe and exists for histogram folding (#199): a source bucket holding a million counts must cost one bin add, never a million calls to Observe.
ObserveN never allocates.
func (*Sketch) ObserveZero ¶
ObserveZero records n observations that belong in the zero bucket. It is the landing slot for an OTLP exponential histogram's zero_count, which by definition holds values at (or within zero_threshold of) zero and has no representation in the log mapping.
func (*Sketch) PopulatedBins ¶
PopulatedBins returns the number of bins holding a non-zero count.
func (*Sketch) Quantile ¶
Quantile returns the estimated q-quantile of the observed values, for q in [0,1]. An empty sketch returns 0; an out-of-range or NaN q returns NaN.
The returned value is within RelativeError of the exact q-quantile of the observed sample, provided the sketch has not collapsed below that quantile and no bin has saturated.
func (*Sketch) RelativeError ¶
RelativeError returns the worst-case relative error of a quantile estimate at the sketch's current scale, ignoring degradation from collapse or saturation. This is the accuracy the sketch advertises to callers.
func (*Sketch) Saturations ¶
Saturations returns the number of bin adds that clamped at MaxUint32. Quantile estimates from a saturated sketch are degraded, not corrupted.
type SketchDropReason ¶
type SketchDropReason uint8
SketchDropReason names why a histogram point's percentiles are unavailable while its scalar statistics were still accepted. It is persisted with the delta so a read can say WHY it has no percentile rather than implying the window was empty.
const ( // SketchDropNone means the sketch describes the whole distribution. SketchDropNone SketchDropReason = 0 // SketchDropNegativeObservations means the point recorded values below // zero. The positive-only sketch cannot hold them, and reporting the // positive side's p99 as the distribution's p99 would be a lie. SketchDropNegativeObservations SketchDropReason = 1 // SketchDropScaleOutOfRange means an ExponentialHistogram arrived at a // scale below 0. Downscaling the sketch to a negative scale is not // representable, and upscaling the point's buckets would not be exact. SketchDropScaleOutOfRange SketchDropReason = 2 // SketchDropNoFiniteBoundaries means an explicit-bounds Histogram carried // observations but no finite boundary to place them between. SketchDropNoFiniteBoundaries SketchDropReason = 3 )
SketchDropReason values. The numbering is durable: it is stored in the delta's histogram flags word.
func (SketchDropReason) String ¶
func (r SketchDropReason) String() string
String renders the reason as the metric label value.
type Snapshot ¶
type Snapshot struct {
// Revision is the revision at the end of the walk.
Revision uint64
// Windows are the mutable windows, oldest first.
Windows []WindowSnapshot
// ActiveSeries and ActiveBySignal mirror the limiter's census.
ActiveSeries int
ActiveBySignal map[Signal]int
// Overflow counts admissions rerouted to an __other__ series, per reason.
Overflow map[OverflowReason]uint64
// WindowsDiscarded and SeriesDiscarded count rollover loss.
WindowsDiscarded uint64
SeriesDiscarded uint64
// ClosedWindowsForced counts closed-but-unfinalized windows the cap
// forced out of memory. Every one is data the finalizer never
// materialized, so a non-zero value means the finalizer is behind.
ClosedWindowsForced uint64
// ClosedWindows is how many closed-but-unfinalized windows memory holds.
ClosedWindows int
}
Snapshot is a consistent-enough view of the engine for tests and metrics.
It walks the shards one at a time, so it is not a single atomic instant across shards. That is deliberate: an atomic cross-shard snapshot would need all four locks held at once, which is the one thing the shard design forbids.
type SnapshotEdge ¶
type SnapshotEdge struct {
Caller string `json:"caller"`
Callee string `json:"callee"`
FirstSeen time.Time `json:"first_seen"`
LastSeen time.Time `json:"last_seen"`
Windows []TopologyWindow `json:"windows"`
}
SnapshotEdge is one caller/callee service pair and its retained windows.
type SpanInput ¶
type SpanInput struct {
Tenant string
Service string
// SpanName is the raw span name, used for operation naming when no route
// attribute is present.
SpanName string
// HTTPRoute is http.route, URLPath is url.path or http.target. Route
// normalization precedence is fixed in #159.
HTTPRoute string
URLPath string
// Method is the raw HTTP method string; it collapses onto the bounded
// Method enum.
Method string
// HTTPStatusCode is http.response.status_code, 0 when absent.
HTTPStatusCode int
// SpanKind and StatusCode are the OTLP numeric values.
SpanKind int32
StatusCode int32
// Root reports that the span has no parent span, i.e. it starts a trace.
// The parent span ID itself is NOT carried: it is on the permanent banned
// list (#153, #159) and only its emptiness affects an aggregate.
Root bool
// Timestamp is the span start time; it selects the window.
Timestamp time.Time
DurationMicros float64
}
SpanInput is one span, already parsed out of OTLP. Only the fields that can affect series identity or aggregates are carried: IDs, URLs and messages are on the permanent banned list (#153, #159) and never reach this struct.
type StatusClass ¶
type StatusClass uint8
StatusClass is the per-signal status dimension of a SeriesKey. Its meaning depends on the Signal:
SignalTraceOp, SignalServiceEdge — OTLP span status: StatusUnset/OK/Error. SignalLog — severity tier: SeverityTierTrace..Fatal. SignalMetric — always 0.
const ( StatusUnset StatusClass = 0 StatusOK StatusClass = 1 StatusError StatusClass = 2 )
Span-status values of StatusClass, mirroring the OTLP status code numbering.
const ( SeverityTierUnspecified StatusClass = 0 SeverityTierTrace StatusClass = 1 SeverityTierDebug StatusClass = 2 SeverityTierInfo StatusClass = 3 SeverityTierWarn StatusClass = 4 SeverityTierError StatusClass = 5 SeverityTierFatal StatusClass = 6 )
Severity-tier values of StatusClass. OTLP severity numbers collapse into six tiers; the raw number never reaches series identity.
func SeverityTier ¶
func SeverityTier(text string, number int32) StatusClass
SeverityTier maps an OTLP severity onto the log StatusClass tier. The numeric severity wins when present; otherwise the text is classified with the same substring rules the legacy ingest gate uses, so a shadow-mode comparison is not thrown off by a disagreement about what "WARNING" means.
func SeverityTierFromNumber ¶
func SeverityTierFromNumber(sev int32) StatusClass
SeverityTierFromNumber maps an OTLP severity number (1..24) onto a severity tier. Numbers outside the range yield SeverityTierUnspecified.
func TraceStatusFromCode ¶
func TraceStatusFromCode(code int32) StatusClass
TraceStatusFromCode maps an OTLP span status code onto a StatusClass. Unrecognized codes degrade to StatusUnset rather than inventing identity.
type Store ¶
type Store interface {
// CommitGroup writes one group batch inside exactly one transaction.
// Either everything in the batch is durable when it returns nil, or
// nothing in it is.
CommitGroup(b *GroupBatch) error
// FinalizeWindow materializes the window's buckets from the delta log and
// deletes exactly the delta rows it incorporated, in one transaction.
FinalizeWindow(windowStart int64) (FinalizeStats, error)
// FinalizableWindows returns the un-finalized window starts at or below
// cutoff, oldest first, capped at limit.
FinalizableWindows(cutoff int64, limit int) ([]int64, error)
// PurgeBefore deletes finalized history older than cutoff.
PurgeBefore(cutoff int64) (PurgeStats, error)
// ReadBuckets returns the store-owned rows matching sel: materialized
// buckets AND any not-yet-finalized delta rows for the same range. Both
// are store-owned data once the engine has handed the window over, and
// omitting either is exactly the silent omission #194 blocker 4 is about.
//
// The selector's bounds are mandatory and the row cap is enforced
// store-side. The read asks the database for limit+1 rows and reports
// Truncated rather than trimming in silence; a caller that needs every row
// pages with Selector.After until Truncated is false.
ReadBuckets(ctx context.Context, sel Selector) (BucketPage, error)
// SumBuckets aggregates the same row set as ReadBuckets in SQL, grouped by
// by. NO row cap applies: the result is bounded by the grouping — windows,
// services, signals — not by the number of rows scanned, which is what
// makes a scalar dashboard total structurally impossible to truncate.
SumBuckets(ctx context.Context, sel Selector, by GroupBy) ([]SumRow, error)
// VisitSketches streams every sketch-bearing row matching sel — from
// BOTH durable tables, like ReadBuckets — to visit, in no particular
// order. It is the percentile path of the #197 read contract: a quantile
// sketch cannot be SUMmed in SQL, but its merge is commutative and
// associative, so no total order and no pagination are needed — one
// unordered pass replaces the sorted, keyset-paged drain that made a
// wide dashboard range cost O(pages) sorted re-scans (#219).
//
// The selector's bounds are mandatory (Selector.Validate applies) and
// SketchOnly is implied. The row cap is NOT applied: it bounds RESULT
// sizes, and this read's result is whatever the caller folds the
// visited sketches into; the scan itself is bounded by the validated
// window span, the same bound SumBuckets relies on.
//
// The visited sketch is valid only for the duration of the call: the
// store decodes every row into one scratch value, so a visitor must
// merge or copy it, never retain the pointer.
VisitSketches(ctx context.Context, sel Selector, visit func(serviceID uint32, sk *Sketch) error) error
// ReplayMutable returns the delta-log rows for windows at or after since —
// the mutable set only. Finalized history never hydrates into RAM (#160).
ReplayMutable(since int64) ([]DeltaRow, error)
// ReadFinalizedSince returns MATERIALIZED bucket rows at or after since
// for the given signals, across every tenant, NEWEST window first, capped
// at limit rows.
//
// It is the one bounded exception to "finalized history never hydrates"
// (#194 finding 15): the topology projection is rebuilt from it at
// startup so a restart does not erase the recent service map. It feeds
// the projection ONLY — never the mutable shards — and the two bounds
// that keep it a startup cost rather than a startup stall are the caller's
// horizon (since) and the row cap. Truncated reports that the cap cut the
// answer, so a restored topology is never presented as complete when it
// is not.
ReadFinalizedSince(since int64, signals []Signal, limit int) (FinalizedPage, error)
// LoadBaselines returns every durable cumulative baseline, capped at max.
LoadBaselines(max int) ([]BaselineRow, error)
// ResolveSeries resolves series IDs back to their identity. The input
// count is capped at MaxReadRows.
ResolveSeries(ids []SeriesID) ([]SeriesInfo, error)
// LoadDict returns every dictionary row, capped at max. It is how the
// durable registrar warms its cache at startup so IDs stay stable.
LoadDict(max int) ([]DictRow, error)
// LoadSeries returns every series row, capped at max.
LoadSeries(max int) ([]SeriesRow, error)
// Backlog reports the delta-log backlog health bounds.
Backlog() (BacklogStats, error)
// Close releases the store's connections.
Close() error
}
Store is the durable aggregate store. Implementations must be safe for concurrent use; CommitGroup and FinalizeWindow serialize internally on the single writer.
type StoreConfig ¶
type StoreConfig struct {
// Path is the database file (AGGREGATE_DB_PATH). Empty takes
// DefaultAggregateDBPath. ":memory:" and file: URIs are honoured for
// tests.
Path string
// AllowRebuild permits DESTROYING and recreating the aggregate tables
// when the on-disk schema is partial or version-mismatched
// (AGGREGATE_ALLOW_REBUILD). Off by default: silent data loss is worse
// than a refused startup.
AllowRebuild bool
// Synchronous is the SQLite synchronous mode ("NORMAL" or "FULL").
// Empty takes NORMAL. See the stanza comment in open() for the durability
// argument.
Synchronous string
// ReadPoolSize is the read-connection count. Zero takes the default.
ReadPoolSize int
// CacheSizeKB is the per-connection page cache in KB, applied as a
// negative cache_size. Zero takes 32 MB — the aggregate working set is
// the mutable delta log, not seven days of buckets.
CacheSizeKB int
// Metrics is the durable path's metric surface. nil disables recording.
Metrics StoreMetrics
}
StoreConfig configures the SQLite aggregate store.
type StoreInspection ¶
type StoreInspection struct {
State string
ExpectedSchemaVersion int
ActualSchemaVersion int
ExpectedSeriesVersion int
ActualSeriesVersion int
ExpectedSketchVersion int
ActualSketchVersion int
StoreUUID string
MigrationResult string
Detail string
}
StoreInspection is a read-only aggregate compatibility result.
func InspectSQLiteStore ¶
func InspectSQLiteStore(path string) (StoreInspection, error)
InspectSQLiteStore reads aggregate tables and metadata without creating, rebuilding, or changing the configured file.
func (StoreInspection) Description ¶
func (s StoreInspection) Description() string
Description returns a stable one-line representation for operator output.
func (StoreInspection) Usable ¶
func (s StoreInspection) Usable() bool
Usable reports whether this binary can safely open the aggregate store.
type StoreMetrics ¶
type StoreMetrics interface {
// RecordCommit publishes one group commit: its wall time, how many
// deltas it carried and how many bytes it wrote.
RecordCommit(d time.Duration, deltas int, bytes int64, err error)
// RecordAdmissionRejected counts one ErrSaturated refusal by bound.
RecordAdmissionRejected(bound string)
// RecordFinalize publishes one window finalization.
RecordFinalize(stats FinalizeStats, err error)
// RecordPurge publishes one retention purge.
RecordPurge(stats PurgeStats, err error)
// SetBacklog publishes the delta-log backlog health bounds.
SetBacklog(rows int64, ageSeconds float64)
// RecordRecovery publishes one startup recovery: its duration and every
// row class it moved.
RecordRecovery(stats RecoveryStats)
// RecordGC publishes one identity garbage-collection pass.
RecordGC(stats GCStats, err error)
}
StoreMetrics is the durable path's metric surface. It mirrors the MetricsRecorder pattern: an interface so tests need no live registry, and a nil-safe no-op default.
func NewPrometheusStoreMetrics ¶
func NewPrometheusStoreMetrics(m *telemetry.Metrics) StoreMetrics
NewPrometheusStoreMetrics returns a StoreMetrics backed by the platform metrics. A nil *telemetry.Metrics yields a no-op recorder rather than a panic, which is what every store unit test passes.
type SumRow ¶
type SumRow struct {
WindowStart int64
ServiceID uint32
// NameID is set only when GroupByName was requested. It resolves through
// NameKind(Signal), never through the service namespace.
NameID uint32
Signal Signal
Count uint64
ErrorCount uint64
RequestCount uint64
ErrorRequestCount uint64
DurationCount uint64
DurationSum float64
LogCount uint64
}
SumRow is one grouped aggregation over the store's rows. Fields outside the requested grouping are zero.
The counters are the SUMmable subset of AggregateDelta: everything a scalar dashboard total is built from. Sketches are deliberately absent — a quantile sketch cannot be merged by SQL, and pretending otherwise is how a p99 turns into an average (#197 Q1).
type SweepStats ¶
type SweepStats struct {
// Series, Dict and Templates count deleted rows per table.
Series, Dict, Templates int64
// Duration is the wall time of the delete transaction.
Duration time.Duration
}
SweepStats is the outcome of one identity sweep.
type TemplateFact ¶
type TemplateFact struct {
Tenant string
Service string
Severity string
TemplateID uint32
Template string
Timestamp time.Time
IsOther bool
}
TemplateFact is the per-log-line record handed to GraphRAG. GraphRAG does no mining of its own in shadow and aggregate modes; it consumes these.
type TemplateMiner ¶
type TemplateMiner struct {
// contains filtered or unexported fields
}
TemplateMiner mines log templates, partitioned by (tenant, service). The zero value is not usable; construct one with NewTemplateMiner.
func NewTemplateMiner ¶
func NewTemplateMiner(cfg TemplateMinerConfig) *TemplateMiner
NewTemplateMiner builds a miner from cfg, applying defaults.
func (*TemplateMiner) Committed ¶
func (m *TemplateMiner) Committed(rows []TemplateRow)
Committed marks drained rows durable. A row whose pattern moved on since the drain stays staged: the version guard is what makes that safe.
func (*TemplateMiner) DrainDirtyStats ¶
func (m *TemplateMiner) DrainDirtyStats() []TemplateStatRow
DrainDirtyStats returns the statistics-only updates for the periodic write and clears the dirty set. These are fire-and-forget: a lost batch costs a count that the next line restores, never an identity.
func (*TemplateMiner) DrainPending ¶
func (m *TemplateMiner) DrainPending() []TemplateRow
DrainPending returns the staged identity mutations for the next group commit, oldest ID first. They stay staged until Committed confirms them, so a failed commit re-offers them rather than acknowledging a delta whose template identity never became durable.
func (*TemplateMiner) Mine ¶
func (m *TemplateMiner) Mine(tenant, service, severity, body string) (id uint32, isOther bool)
Mine clusters one log body and returns its template ID. isOther is true when the line was absorbed by the partition's overflow template — the caller must still count it, only the identity detail is gone. Mine never fails; a registrar that cannot allocate yields (0, true).
func (*TemplateMiner) MineAt ¶
func (m *TemplateMiner) MineAt(tenant, service, severity, body string, at time.Time) (id uint32, isOther bool)
MineAt is Mine with an explicit arrival time — the reducer captures one timestamp per Export request and passes it for every point in that request.
func (*TemplateMiner) PartitionStats ¶
func (m *TemplateMiner) PartitionStats(tenant, service string) (TemplatePartitionStats, bool)
PartitionStats returns statistics for one partition, if it exists.
func (*TemplateMiner) PendingCount ¶
func (m *TemplateMiner) PendingCount() int
PendingCount reports how many identity mutations are staged but not durable.
func (*TemplateMiner) Restore ¶
func (m *TemplateMiner) Restore(rows []TemplateRow)
Restore rebuilds the miner's partitions, prefix trees, alias index and cap accounting from durable rows. It must run BEFORE ingest starts: a line mined against an empty miner would mint a second identity for a pattern that already has one, and both would be live.
Restoring stages nothing: every row handed here is already durable.
func (*TemplateMiner) Roots ¶
func (m *TemplateMiner) Roots() map[uint32]struct{}
Roots returns every template ID the miner keeps alive: live templates, the per-partition overflow sentinels, both ends of every alias, and every staged mutation (#200 Q5).
Both ends, deliberately. A historical series that named a retired template keeps that ID, its alias row, AND the survivor alive: collecting the survivor would leave the alias pointing at nothing, and collecting the alias would leave the series unresolvable.
func (*TemplateMiner) SetFactSink ¶
func (m *TemplateMiner) SetFactSink(fn func(TemplateFact))
SetFactSink installs (or, with nil, removes) the log-fact consumer. Safe to call while mining is in flight.
func (*TemplateMiner) Stats ¶
func (m *TemplateMiner) Stats() []TemplatePartitionStats
Stats returns per-partition statistics, sorted by tenant then service.
func (*TemplateMiner) TemplateText ¶
func (m *TemplateMiner) TemplateText(id uint32) (string, bool)
TemplateText returns the current rendered text for a template ID, following convergence aliases. IDs issued in the past always resolve: a retired ID resolves to its survivor's text, never to nothing.
type TemplateMinerConfig ¶
type TemplateMinerConfig struct {
// MaxTemplatesPerService caps distinct mined templates per
// (tenant, service). The overflow template is reserved capacity and does
// not count against it. Default DefaultMaxLogTemplatesPerService.
MaxTemplatesPerService int
// Depth is the prefix-tree depth below the token-length layer.
// Default 4, minimum 1.
Depth int
// SimilarityThreshold is Drain's simSeq threshold in (0, 1]. Default 0.4.
SimilarityThreshold float64
// MaxChildren caps distinct children per tree node; overflow tokens route
// through the wildcard child. Default 100.
MaxChildren int
// MaxTokens bounds tokens taken from one log body, keeping Mine() latency
// independent of body size. Default 128.
MaxTokens int
// Registrar allocates template IDs. Default NewInMemoryTemplateRegistrar().
Registrar TemplateRegistrar
// OnFact, when set, is called synchronously for every mined line —
// including overflow lines — outside the partition lock. Keep it cheap.
OnFact func(TemplateFact)
}
TemplateMinerConfig configures a TemplateMiner. Zero values take defaults.
type TemplatePartitionStats ¶
type TemplatePartitionStats struct {
Tenant string
Service string
Templates int // live mined templates (excludes __other__)
OtherID uint32 // 0 when the overflow identity is not allocated
Overflow uint64 // lines absorbed by __other__
Converged uint64 // templates retired into a surviving twin
RegistrarFailures uint64
}
TemplatePartitionStats is a point-in-time view of one (tenant, service) partition.
type TemplateRegistrar ¶
type TemplateRegistrar interface {
RegisterTemplate(TemplateRegistration) (uint32, error)
}
TemplateRegistrar allocates immutable surrogate template IDs. Phase 1 uses an in-memory allocator; Phase 2 replaces it with the durable aggregate dictionary, where registration commits atomically with the first referencing delta.
Implementations must be safe for concurrent use and bounded in latency: the miner calls RegisterTemplate while holding a partition lock. A returned error, or a zero ID, is treated as "no identity available" — the miner then routes the line to the partition's overflow template and never fails.
func NewInMemoryTemplateRegistrar ¶
func NewInMemoryTemplateRegistrar() TemplateRegistrar
NewInMemoryTemplateRegistrar returns the Phase 1 provisional allocator. IDs are process-local and do not survive a restart.
type TemplateRegistrarFunc ¶
type TemplateRegistrarFunc func(TemplateRegistration) (uint32, error)
TemplateRegistrarFunc adapts a plain function to TemplateRegistrar.
func (TemplateRegistrarFunc) RegisterTemplate ¶
func (f TemplateRegistrarFunc) RegisterTemplate(r TemplateRegistration) (uint32, error)
RegisterTemplate implements TemplateRegistrar.
type TemplateRegistration ¶
type TemplateRegistration struct {
Tenant string
Service string
Template string // rendered template text at registration time
IsOther bool // true for a partition's pre-created __other__ identity
}
TemplateRegistration is the request handed to a TemplateRegistrar when the miner mints a new template identity. It mirrors the dictionary registration contract for kind=log_template (#159, amended by #163): the value written to the dictionary is a surrogate identity whose presentation text may later change without the ID changing.
type TemplateRow ¶
type TemplateRow struct {
// ID is the immutable surrogate template identity — the same uint32 the
// log_template dictionary minted, and the NameID of every log series that
// used it.
ID uint32
// Tenant and Service name the partition. They are the strings, not
// dictionary IDs: the miner is keyed by them and a reload must rebuild
// its partitions before any dictionary lookup is warm.
Tenant, Service string
// PatternVersion increments on every generalization of Tokens. It makes a
// stale periodic stats write unable to overwrite a newer pattern.
PatternVersion uint32
// Tokens is the token pattern, NUL-joined. Tokens never contain
// whitespace (the tokenizer splits on it) and never contain NUL.
Tokens string
// Seq is the partition-local creation ordinal. Convergence keeps the
// lower one, so it has to survive a restart or two restarts could pick
// different survivors for the same pair.
Seq uint64
// IsOther marks the partition's pre-created overflow identity.
IsOther bool
// AliasOf is the surviving template this ID forwards to, or 0.
AliasOf uint32
// Count, FirstSeen and LastSeen are the non-identity statistics. They are
// refreshed by the periodic dirty-partition write, not by the identity
// commit path.
Count uint64
FirstSeen, LastSeen int64
}
TemplateRow is one durable log-template record (#200 Q4).
It carries exactly what a reload needs to rebuild the prefix tree, resolve a historical template ID, and re-enforce the per-partition cap. It carries NO raw log sample: the miner keeps one in memory for diagnostics, and persisting it would turn the aggregate file into a credential and PII sink for the sake of a field exemplars already provide.
type TemplateStatRow ¶
TemplateStatRow is the non-identity half of a template: the counters a periodic dirty-partition write refreshes. Losing one costs a count, not an identity, so it does not need to ride a group commit.
type Temporality ¶
type Temporality uint8
Temporality is the OTLP aggregation temporality of a metric point.
const ( TemporalityUnspecified Temporality = 0 TemporalityDelta Temporality = 1 TemporalityCumulative Temporality = 2 )
Temporality values, mirroring the OTLP numbering.
type TopologyConfig ¶
type TopologyConfig struct {
MaxServices int
MaxOperationsPerService int
MaxEdges int
MaxMetrics int
// Horizon is the finalized history retained behind the mutable windows.
Horizon time.Duration
}
TopologyConfig bounds the projection. Zero values take the defaults above.
type TopologyEdge ¶
type TopologyEdge struct {
Source, Target string
CallCount int64
AvgLatencyMs float64
ErrorRate float64
}
TopologyEdge is one caller/callee edge of the topology.
type TopologyMetric ¶
type TopologyMetric struct {
Service string `json:"service"`
Metric string `json:"metric"`
FirstSeen time.Time `json:"first_seen"`
LastSeen time.Time `json:"last_seen"`
Windows []TopologyWindow `json:"windows"`
}
TopologyMetric is one (service, metric) series and its retained windows.
type TopologyOperation ¶
type TopologyOperation struct {
Service string `json:"service"`
Operation string `json:"operation"`
FirstSeen time.Time `json:"first_seen"`
LastSeen time.Time `json:"last_seen"`
Windows []TopologyWindow `json:"windows"`
}
TopologyOperation is one (service, operation) and its retained windows.
type TopologyResult ¶
type TopologyResult struct {
Nodes []ServiceStat
Edges []TopologyEdge
Coverage Coverage
Epoch string
Revision uint64
}
TopologyResult is the answer to QueryTopology.
Nodes and Edges come from ONE query over ONE tenant and ONE range, read under ONE ownership snapshot. That is the whole point (#194 finding 15): before it, edges were supplemented into API responses from GraphRAG, a store with a different range, a different tenant scope and a different retention rule, and nothing in the response said so.
type TopologyService ¶
type TopologyService struct {
Name string `json:"name"`
FirstSeen time.Time `json:"first_seen"`
LastSeen time.Time `json:"last_seen"`
Windows []TopologyWindow `json:"windows"`
}
TopologyService is one service and its retained windows, oldest first.
type TopologySnapshot ¶
type TopologySnapshot struct {
Tenant string `json:"tenant"`
Epoch uint64 `json:"epoch"`
Revision uint64 `json:"revision"`
Now time.Time `json:"now"`
Horizon time.Duration `json:"horizon"`
Services []TopologyService `json:"services"`
Operations []TopologyOperation `json:"operations"`
Edges []SnapshotEdge `json:"edges"`
Metrics []TopologyMetric `json:"metrics"`
// Dropped* count facts refused by the projection caps since startup.
// Non-zero means the topology is truncated and must be presented as such.
DroppedServices uint64 `json:"dropped_services,omitempty"`
DroppedOperations uint64 `json:"dropped_operations,omitempty"`
DroppedEdges uint64 `json:"dropped_edges,omitempty"`
DroppedMetrics uint64 `json:"dropped_metrics,omitempty"`
}
TopologySnapshot is one tenant's replacement-by-revision topology view.
Revision is the engine revision at the last fold that CHANGED this tenant. A consumer that already applied the same (Epoch, Revision) pair must do nothing. Epoch changes when the engine's identity resets — a restart — after which Revision starts over and the consumer must replace rather than reconcile.
func (TopologySnapshot) Empty ¶
func (s TopologySnapshot) Empty() bool
Empty reports whether the snapshot carries no topology at all.
func (TopologySnapshot) Truncated ¶
func (s TopologySnapshot) Truncated() bool
Truncated reports whether any projection cap has refused a fact for this tenant. A consumer must not present a truncated topology as complete.
type TopologyWindow ¶
type TopologyWindow struct {
Start time.Time `json:"start"`
End time.Time `json:"end"`
Closed bool `json:"closed"`
Final bool `json:"final"`
Elapsed time.Duration `json:"elapsed"`
Count uint64 `json:"count"`
ErrorCount uint64 `json:"error_count"`
DurationCount uint64 `json:"duration_count,omitempty"`
DurationSumMicros float64 `json:"duration_sum_micros,omitempty"`
DurationMinMicros float64 `json:"duration_min_micros,omitempty"`
DurationMaxMicros float64 `json:"duration_max_micros,omitempty"`
// P95Micros and P99Micros come from the merged latency sketch. They are
// zero when the entity carries no sketch (operations, edges and metrics
// deliberately do not, so a 2 KiB sketch is not multiplied by every
// operation of every service of every window).
P95Micros float64 `json:"p95_micros,omitempty"`
P99Micros float64 `json:"p99_micros,omitempty"`
LatencyProvenance *latency.Provenance `json:"latency_provenance,omitempty"`
// Value* carry metric samples: gauge samples plus counter deltas, which
// is what a rolling mean/variance baseline is computed over.
ValueCount uint64 `json:"value_count,omitempty"`
ValueSum float64 `json:"value_sum,omitempty"`
ValueMin float64 `json:"value_min,omitempty"`
ValueMax float64 `json:"value_max,omitempty"`
}
TopologyWindow is one five-minute window of one topology entity.
Closed and Final are the partial-window guard the anomaly detector needs: Closed means the wall clock has passed the window's end, Final means the lateness horizon has expired too and no further point can land in it. A detector that compares the current, still-filling window against a baseline without consulting Elapsed and Count is the "anomaly storm from a nearly empty window" #163 removed.
func (TopologyWindow) AvgLatencyMs ¶
func (w TopologyWindow) AvgLatencyMs() float64
AvgLatencyMs returns the mean duration in milliseconds, or 0 when the window carries no duration observations.
func (TopologyWindow) ErrorRate ¶
func (w TopologyWindow) ErrorRate() float64
ErrorRate returns errors/count, or 0 for an empty window.
func (TopologyWindow) Mean ¶
func (w TopologyWindow) Mean() float64
Mean returns the mean observed value, or 0 when the window holds none.
type TrafficPoint ¶
type TrafficPoint struct {
// WindowStart is the UTC start of the five-minute window.
WindowStart time.Time
// RequestCount is accepted request entry points — root or SERVER spans.
// ErrorRequestCount is its error subset.
RequestCount, ErrorRequestCount int64
// SpanCount is accepted spans; SpanErrorCount is its error subset.
SpanCount, SpanErrorCount int64
}
TrafficPoint is one window of the traffic series.
BOTH bases are carried, explicitly named (#197 Q3). Traffic is plotted on the request basis; the span basis stays available because it is what per-operation diagnostics and latency are counted in.
type Variant ¶
type Variant uint8
Variant is the signal-specific variant dimension. For traces and edges it is the OTLP SpanKind; every other signal pins it to zero.
const ( SpanKindUnspecified Variant = 0 SpanKindInternal Variant = 1 SpanKindServer Variant = 2 SpanKindClient Variant = 3 SpanKindProducer Variant = 4 SpanKindConsumer Variant = 5 )
SpanKind values of Variant, mirroring the OTLP SpanKind numbering.
func VariantFromSpanKind ¶
VariantFromSpanKind maps an OTLP SpanKind number onto a Variant. Out-of-range kinds degrade to SpanKindUnspecified.
type VersionError ¶
VersionError reports an encoded key whose version byte is not one this build understands.
func (*VersionError) Error ¶
func (e *VersionError) Error() string
type WatermarkStore ¶
type WatermarkStore interface {
// Watermarks returns the persisted dictionary and series high-watermarks:
// the next ID each allocator may mint. Zero means "not yet stamped".
Watermarks() (uint32, SeriesID, error)
}
WatermarkStore is the optional Store capability that carries the identity high-watermarks. It is separate from Store so an implementation predating #200 still satisfies Store; the registrars degrade to MAX(id)+1 without it, which is correct exactly as long as nothing collects.
type WindowSnapshot ¶
type WindowSnapshot struct {
// Start is the window's UTC start.
Start time.Time
// Series maps each active series in the window to a copy of its delta.
Series map[SeriesKey]*AggregateDelta
}
WindowSnapshot is one mutable window's contents.
type Writer ¶
type Writer struct {
// contains filtered or unexported fields
}
Writer is the group-commit writer. It is the engine's Applier once the durable store is enabled.
func NewWriter ¶
func NewWriter(cfg WriterConfig) (*Writer, error)
NewWriter builds a writer over store and engine. It does not start it.
func (*Writer) Apply ¶
Apply implements Applier. It exists so a caller that cannot act on an error still gets the durable path; every ingest caller uses ApplyErr instead.
func (*Writer) ApplyErr ¶
ApplyErr implements FailableApplier: it blocks until the deltas are durable and applied, or until admission refuses them.
func (*Writer) CollectIdentities ¶
CollectIdentities runs one identity garbage-collection pass. It is the entry point retention's daily maintenance tick calls.
func (*Writer) FinalizeDue ¶
FinalizeDue finalizes every window whose lateness horizon expired at now and returns how many it finalized. Exported so recovery and tests can drive the same state machine the loop drives.
func (*Writer) RunBarrier ¶
RunBarrier implements Barrier. It parks fn on the commit goroutine, so no group commit is in flight while it runs and none can start until it returns.
fn must be bounded: every Export waiting on the writer is waiting on it too. The collector honours that by keeping its full table scan OUTSIDE the barrier and passing only the revalidate-fence-delete tail in here.
func (*Writer) SaveTemplateStats ¶
SaveTemplateStats writes the miner's dirty non-identity counters. Identity mutations do NOT come through here — they ride the group commit — so a lost batch costs a count that the next line restores.
func (*Writer) SeriesID ¶
SeriesID exposes the writer's series resolution for tests and for the read path that has to turn a SeriesKey into a durable ID.
func (*Writer) SeriesKeyByID ¶
SeriesKeyByID resolves a durable ID back to its identity from memory.
func (*Writer) Shutdown ¶
Shutdown drains in-flight submissions and joins both writer loops within the shared application deadline. The final template-statistics write is part of the durability result rather than a best-effort log line.
func (*Writer) Start ¶
func (w *Writer) Start()
Start launches the commit loop and, unless disabled, the finalize loop.
func (*Writer) Stats ¶
func (w *Writer) Stats() WriterStats
Stats returns a snapshot of the writer counters.
type WriterConfig ¶
type WriterConfig struct {
// Store is the durable store. Required.
Store Store
// Engine is the engine whose shards receive committed deltas. Required.
Engine *Engine
// Registrar is the durable dictionary registrar whose staged rows ride
// each commit. Required when the engine was built with one.
Registrar *DurableRegistrar
// CoalesceWindow, MaxBatchDeltas and MaxBatchBytes set the commit cadence.
CoalesceWindow time.Duration
MaxBatchDeltas int
MaxBatchBytes int64
// MaxPendingBytes, MaxWaiters and MaxPendingDeltas are the triple
// admission bound.
MaxPendingBytes int64
MaxWaiters int
MaxPendingDeltas int
// FinalizeInterval is the window-finalization tick. Zero takes the
// default; negative disables the loop (tests drive it by hand).
FinalizeInterval time.Duration
// Metrics is the durable path's metric surface.
Metrics StoreMetrics
// Now is the clock, injectable for tests.
Now func() time.Time
}
WriterConfig configures the group-commit writer.
type WriterStats ¶
type WriterStats struct {
// Commits and CommitErrors count group commits attempted and failed.
Commits, CommitErrors uint64
// Deltas counts delta rows written.
Deltas uint64
// Rejections counts ErrSaturated refusals.
Rejections uint64
// PendingBytes, PendingDeltas and Waiters are the live admission
// occupancy.
PendingBytes int64
PendingDeltas int
Waiters int
// Finalized counts windows finalized by the writer's state machine.
Finalized uint64
// The runtime-readiness surface (#194 finding 18). Cumulative counters
// say what happened since boot; a readiness probe needs to know what is
// happening NOW, which is what the streaks and the cached backlog sample
// carry.
//
// CommitFailureStreak and FinalizeFailureStreak count CONSECUTIVE
// failures: any success resets them to zero. A handful of failures spread
// over a week is a log line; three in a row is a store that has stopped
// accepting writes.
CommitFailureStreak uint64
FinalizeFailureStreak uint64
// FinalizeErrors counts finalization failures since boot.
FinalizeErrors uint64
// MaxPendingBytes, MaxPendingDeltas and MaxWaiters are the configured
// admission bounds the live occupancy above is measured against.
MaxPendingBytes int64
MaxPendingDeltas int
MaxWaiters int
// DeltaLogRows, DeltaLogAgeSeconds and BacklogSampledAt are the last
// delta-log backlog sample the finalize loop published. Cached rather
// than queried on demand so a readiness probe never stacks behind the
// single SQLite writer. BacklogSampledAt is zero before the first sample.
DeltaLogRows int64
DeltaLogAgeSeconds float64
BacklogSampledAt time.Time
}
WriterStats is a snapshot of writer counters.
func (WriterStats) AdmissionRatio ¶
func (s WriterStats) AdmissionRatio() float64
AdmissionRatio is the writer's admission occupancy as a fraction in [0.0, 1.0+]: the fullest of the three bounds, because admission is refused as soon as ANY of them is breached. Bounds that are unset (non-positive) contribute nothing.
func (WriterStats) DeltaLogAge ¶
func (s WriterStats) DeltaLogAge(now time.Time) float64
DeltaLogAge is the age of the oldest un-finalized window at now, carrying the staleness of the sample it is derived from.
Adding the time since the sample is not cosmetic: a wedged finalize loop stops REFRESHING the sample, and an age frozen at its last healthy value is exactly the reading that would let a stuck writer look ready forever. An empty backlog ages at zero — there is no oldest window to get older.