aggregate

package
v0.6.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 21, 2026 License: MIT Imports: 29 Imported by: 0

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

View Source
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.

View Source
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).

View Source
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.

View Source
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.

View Source
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.

View Source
const (
	// HistPercentilesUnavailable marks that the sketch does NOT describe the
	// 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.

View Source
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.

View Source
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.

View Source
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.

View Source
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.

View Source
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
)
View Source
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.

View Source
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.

View Source
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.

View Source
const (
	// DefaultAggregateDBPath is the default AGGREGATE_DB_PATH.
	DefaultAggregateDBPath = "./data/aggregate.db"
)

Default store tuning.

View Source
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.

View Source
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.

View Source
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.

View Source
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.

View Source
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.

View Source
const IDPlaceholder = "{id}"

IDPlaceholder replaces a path segment the segment rules classify as variable.

View Source
const MaxDimensionKeys = 16

MaxDimensionKeys bounds one metric's configured dimension tuple. A tuple longer than this is refused from identity rather than silently truncated.

View Source
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.

View Source
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).

View Source
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.

View Source
const SketchEncodingVersion uint8 = 0x01

SketchEncodingVersion is the only encoding version this package writes or accepts.

View Source
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.

View Source
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

View Source
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.

View Source
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.

View Source
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.

View Source
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.

View Source
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.

View Source
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.

View Source
var ErrSketchScale = errors.New("aggregate: unsupported sketch scale")

ErrSketchScale is returned when a scale outside the supported range is requested or decoded.

View Source
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.

View Source
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

func AppendCanonicalDims(dst []byte, pairs []DimPair) []byte

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

func DecodeSketchInto(dst *Sketch, data []byte) error

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

func FinalizeCutoff(now time.Time) int64

FinalizeCutoff returns the newest window start whose lateness horizon has expired at now — windows at or below it are ready to finalize.

func InternDimValues

func InternDimValues(c *Cache, tenantID uint32, keys, values []string) uint32

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

func IsDiskFull(err error) bool

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

func IsRequestSpan(root bool, spanKind int32) bool

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 LogGC

func LogGC(stats GCStats, err error)

LogGC prints one collection pass at info level.

func LogRecovery

func LogRecovery(stats RecoveryStats, path string)

LogRecovery emits the operator-facing recovery summary.

func MutableSince

func MutableSince(now time.Time) int64

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

func NormalizeOperation(httpRoute, urlPath, spanName string) string

NormalizeOperation resolves the operation name of an HTTP span using the precedence fixed in #159:

  1. http.route verbatim when present — the instrumentation already told us the template, and second-guessing it can only make things worse.
  2. otherwise url.path / http.target, normalized.
  3. 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

func NormalizePath(path string) string

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

func NormalizeSpanName(name string) string

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

func WindowStart(t time.Time) int64

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).
	PercentilesUnavailable       bool   `json:"percentiles_unavailable,omitempty"`
	PercentilesUnavailableReason string `json:"percentiles_unavailable_reason,omitempty"`
}

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

type Applier interface {
	Apply(DeltaMap) uint64
}

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 NewCache

func NewCache(reg Registrar) *Cache

NewCache returns a Cache in front of reg with default bounds. reg must not be nil.

func NewCacheWithBounds

func NewCacheWithBounds(reg Registrar, b Bounds) *Cache

NewCacheWithBounds returns a Cache in front of reg honouring b (#200 Q3).

func (*Cache) Bounds

func (c *Cache) Bounds() Bounds

Bounds returns the identity bounds in force.

func (*Cache) Fence

func (c *Cache) Fence(ids map[uint32]struct{})

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

func (c *Cache) Forget(ids map[uint32]struct{})

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

func (c *Cache) Intern(tenantID uint32, kind Kind, value string) uint32

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

func (c *Cache) InternBytes(tenantID uint32, kind Kind, value []byte) uint32

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

func (c *Cache) InternDims(tenantID uint32, pairs []DimPair) uint32

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

func (c *Cache) InternTenant(name string) (uint32, bool)

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) Len

func (c *Cache) Len() int

Len returns the number of cached entries across every scope. Diagnostic only.

func (*Cache) Lookup

func (c *Cache) Lookup(id uint32) (DictEntry, bool)

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

func (c *Cache) OtherID(tenantID uint32, kind Kind) uint32

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

func (c *Cache) Roots() map[uint32]struct{}

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

func (c *Cache) SetOverflowSink(fn func(kind Kind, bound string))

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.

func (*Cache) Unfence

func (c *Cache) Unfence()

Unfence releases every fenced ID.

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.

func (Coverage) Note

func (c Coverage) Note() string

Note returns the caveat that belongs with a coverage value, or "" for full coverage, which needs none.

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 DictEntry

type DictEntry struct {
	ID       uint32
	TenantID uint32
	Kind     Kind
	Value    []byte
}

DictEntry is one dictionary row as held by MemRegistrar.

type DictRow

type DictRow struct {
	ID       uint32
	TenantID uint32
	Kind     Kind
	Value    []byte
}

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

type DimPair struct {
	KeyID   uint32
	ValueID uint32
}

DimPair is one operator-configured dimension, already reduced to dictionary IDs. The hot path never carries the strings.

type DimsConfig

type DimsConfig map[string][]string

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) Register

func (r *DurableRegistrar) Register(tenantID uint32, kind Kind, value []byte) (uint32, error)

Register implements Registrar.

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

func (e *Engine) ActiveSeriesKeys() map[SeriesKey]struct{}

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

func (e *Engine) ApplyCommitted(m DeltaMap) uint64

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

func (e *Engine) ApplyDeltas(m DeltaMap) uint64

ApplyDeltas applies one reducer's output through the configured applier.

func (*Engine) ApplyDeltasErr

func (e *Engine) ApplyDeltasErr(m DeltaMap) (uint64, error)

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

func (e *Engine) ApplyReducer(r *Reducer) uint64

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

func (e *Engine) ApplyReducerErr(r *Reducer) (uint64, error)

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) Cache

func (e *Engine) Cache() *Cache

Cache returns the dictionary cache.

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

func (e *Engine) Epoch() string

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) Limiter

func (e *Engine) Limiter() *Limiter

Limiter returns the cardinality limiter.

func (*Engine) MarkFinalized

func (e *Engine) MarkFinalized(windowStart int64)

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) Mode

func (e *Engine) Mode() string

Mode returns the configured aggregate mode.

func (*Engine) NewReducer

func (e *Engine) NewReducer(arrival time.Time) *Reducer

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

func (e *Engine) Ownership() 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

func (e *Engine) QueryBuckets(ctx context.Context, q Query) (*BucketsResult, error)

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

func (e *Engine) QueryDashboard(ctx context.Context, q Query) (*DashboardResult, error)

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

func (e *Engine) QueryTopology(ctx context.Context, q Query) (*TopologyResult, error)

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

func (e *Engine) Revision() uint64

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

func (e *Engine) Rollover(now time.Time) int

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

func (e *Engine) SetApplier(a Applier)

SetApplier replaces the apply path. Phase 2 uses it to insert the group-commit writer between the reducer and the shards.

func (*Engine) SetStore

func (e *Engine) SetStore(st Store)

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) Snapshot

func (e *Engine) Snapshot() Snapshot

Snapshot returns a copy of the mutable window set.

func (*Engine) Store

func (e *Engine) Store() Store

Store returns the wired durable store, or nil when none is configured.

func (*Engine) TenantID

func (e *Engine) TenantID(name string) (uint32, bool)

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

func (e *Engine) TopologyEpoch() uint64

TopologyEpoch returns this engine instance's epoch.

func (*Engine) TopologyHorizon

func (e *Engine) TopologyHorizon() time.Duration

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

func (e *Engine) TopologyRevision(tenant string) uint64

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

func (e *Engine) TopologyTenants() []string

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

type ExpBuckets struct {
	Offset int32
	Counts []uint64
}

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

type FieldError struct {
	Field  string
	Value  uint8
	Signal Signal
}

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.

func Collect

func Collect(cfg GCConfig) (GCStats, error)

Collect runs one mark-and-sweep pass over the aggregate identity tables.

It is safe to call concurrently with ingest. It is NOT safe to call twice concurrently — the writer runs it, and the writer is one goroutine.

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

func HTTPClassFromStatus(code int) HTTPClass

HTTPClassFromStatus maps an HTTP status code onto its family. Codes outside 100..599 yield HTTPClassNone.

func (HTTPClass) String

func (h HTTPClass) String() string

String implements fmt.Stringer.

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.
	PercentilesUnavailable bool
	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.

func NameKind

func NameKind(s Signal) (Kind, bool)

NameKind returns the dictionary kind that NameID is resolved through for the given signal. It reports ok=false for an unspecified or unknown signal.

func (Kind) String

func (k Kind) String() string

String implements fmt.Stringer.

func (Kind) Valid

func (k Kind) Valid() bool

Valid reports whether k is one of the defined dictionary kinds.

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 NewLimiter

func NewLimiter(cfg LimiterConfig) *Limiter

NewLimiter returns a Limiter for cfg.

func (*Limiter) Admit

func (l *Limiter) Admit(key SeriesKey, window int64) Admission

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

func (l *Limiter) IsOverflowSeries(key SeriesKey) bool

IsOverflowSeries reports whether key is a live __other__ series.

func (*Limiter) Release

func (l *Limiter) Release(key SeriesKey, window int64)

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.

func (*MemRegistrar) OtherID

func (r *MemRegistrar) OtherID(tenantID uint32, kind Kind) uint32

OtherID implements Registrar. The overflow entry is created on first demand and bypasses the capacity cap.

func (*MemRegistrar) Register

func (r *MemRegistrar) Register(tenantID uint32, kind Kind, value []byte) (uint32, error)

Register implements Registrar.

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

func LookupMethod(s string) (Method, bool)

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

func ParseMethod(s string) Method

ParseMethod maps an HTTP method string onto the bounded enum. The empty string yields MethodNone; anything unrecognized yields MethodOther.

func (Method) String

func (m Method) String() string

String implements fmt.Stringer. Known methods render in their canonical uppercase form; MethodNone renders as the empty string.

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

func (o Ownership) OwnsInMemory(windowStart int64) bool

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

type PreloadError struct {
	Table string
	Rows  int
	Max   int
}

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) Complete

func (g *RecoveryGate) Complete()

Complete opens the gate.

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) Arrival

func (r *Reducer) Arrival() time.Time

Arrival returns the reducer's arrival time.

func (*Reducer) Deltas

func (r *Reducer) Deltas() DeltaMap

Deltas returns the reduced deltas. The map is handed to the engine as-is; the reducer must not be used afterwards.

func (*Reducer) Len

func (r *Reducer) Len() int

Len returns the number of deltas produced so far.

func (*Reducer) MergeFrom

func (r *Reducer) MergeFrom(other *Reducer)

MergeFrom folds another reducer's deltas and stats into r. Used by Export paths that reduce resource batches in parallel.

func (*Reducer) ReduceEdge

func (r *Reducer) ReduceEdge(in EdgeInput)

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

func (r *Reducer) ReduceLog(in LogInput)

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

func (r *Reducer) ReduceSpan(in SpanInput)

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.

func (ResetReason) String

func (r ResetReason) String() string

String implements fmt.Stringer.

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) Close

func (s *SQLiteStore) Close() error

Close releases both pools.

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) Path

func (s *SQLiteStore) Path() string

Path returns the database file path.

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

func (s *SQLiteStore) SumBuckets(ctx context.Context, sel Selector, by GroupBy) ([]SumRow, error)

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).

func (Selector) Validate

func (s Selector) Validate() (int, error)

Validate applies the mandatory bounds and returns the effective row limit.

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

type SeriesInfo struct {
	ID  SeriesID
	Key SeriesKey
}

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

func DecodeSeriesKey(b []byte) (SeriesKey, error)

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

func (k SeriesKey) AppendBinary(dst []byte) ([]byte, error)

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

func (k SeriesKey) MarshalBinary() ([]byte, error)

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

func (k SeriesKey) String() string

String renders a key for logs and test failures. It is diagnostic output, not a wire format — never parse it.

func (*SeriesKey) UnmarshalBinary

func (k *SeriesKey) UnmarshalBinary(b []byte) error

UnmarshalBinary implements encoding.BinaryUnmarshaler. b must be exactly EncodedSeriesKeyLen bytes.

func (SeriesKey) Validate

func (k SeriesKey) Validate() error

Validate reports whether the key is internally consistent: every enum is in range and the signal-specific constraints from #159 hold. Metric and log series carry no HTTP identity; only traces and edges carry a span kind.

type SeriesRow

type SeriesRow struct {
	ID  SeriesID
	Key SeriesKey
}

SeriesRow is one series registration awaiting commit. Identity columns are explicit — there is no generic "variant" column (#162).

type SeriesWindowKey

type SeriesWindowKey struct {
	Key         SeriesKey
	WindowStart int64
}

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.

func (Signal) String

func (s Signal) String() string

String implements fmt.Stringer. The lowercase forms double as metric label values, so they are part of the exported contract.

func (Signal) Valid

func (s Signal) Valid() bool

Valid reports whether s is one of the four defined signals.

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

func DecodeSketch(data []byte) (*Sketch, error)

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

func NewSketchAtScale(scale uint8) (*Sketch, error)

NewSketchAtScale returns an empty sketch at an explicit scale. Scales above SketchMaxScale are rejected.

func NewSketchAtScaleUnchecked

func NewSketchAtScaleUnchecked(scale uint8) *Sketch

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

func (s *Sketch) AppendTo(dst []byte) []byte

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

func (s *Sketch) Collapsed() bool

Collapsed reports whether counts were merged into the lowest retained bin, either by range overflow or by the serialized bin cap.

func (*Sketch) Count

func (s *Sketch) Count() uint64

Count returns the total number of observations, including zero-bucket values.

func (*Sketch) Downscale

func (s *Sketch) Downscale(target uint8) error

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) Encode

func (s *Sketch) Encode() []byte

Encode returns the serialized form of the sketch.

func (*Sketch) Merge

func (s *Sketch) Merge(other *Sketch)

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

func (s *Sketch) Observe(value float64)

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

func (s *Sketch) ObserveBucket(index int32, n uint64)

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

func (s *Sketch) ObserveN(value float64, n uint64)

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

func (s *Sketch) ObserveZero(n uint64)

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

func (s *Sketch) PopulatedBins() int

PopulatedBins returns the number of bins holding a non-zero count.

func (*Sketch) Quantile

func (s *Sketch) Quantile(q float64) float64

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

func (s *Sketch) RelativeError() float64

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

func (s *Sketch) Saturations() uint64

Saturations returns the number of bin adds that clamped at MaxUint32. Quantile estimates from a saturated sketch are degraded, not corrupted.

func (*Sketch) Scale

func (s *Sketch) Scale() uint8

Scale returns the mapping scale of the sketch.

func (*Sketch) Sum

func (s *Sketch) Sum() float64

Sum returns the sum of all observed values.

func (*Sketch) ZeroCount

func (s *Sketch) ZeroCount() uint64

ZeroCount returns the number of observations that landed in the zero bucket.

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.

func (Snapshot) String

func (s Snapshot) String() string

String implements fmt.Stringer for test failures.

func (Snapshot) Totals

func (s Snapshot) Totals(signal Signal) (count, errors uint64)

Totals sums one signal's counters across every mutable window. It is the shadow-mode comparison primitive: the same input stream must produce the same totals regardless of the sampling rate.

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

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

type TemplateStatRow struct {
	ID                  uint32
	Count               uint64
	FirstSeen, LastSeen int64
}

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

func VariantFromSpanKind(kind int32) Variant

VariantFromSpanKind maps an OTLP SpanKind number onto a Variant. Out-of-range kinds degrade to SpanKindUnspecified.

type VersionError

type VersionError struct {
	Got  uint8
	Want uint8
}

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

func (w *Writer) Apply(m DeltaMap) uint64

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

func (w *Writer) ApplyErr(m DeltaMap) (uint64, error)

ApplyErr implements FailableApplier: it blocks until the deltas are durable and applied, or until admission refuses them.

func (*Writer) CollectIdentities

func (w *Writer) CollectIdentities() (GCStats, error)

CollectIdentities runs one identity garbage-collection pass. It is the entry point retention's daily maintenance tick calls.

func (*Writer) FinalizeDue

func (w *Writer) FinalizeDue(now time.Time) int

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

func (w *Writer) RunBarrier(fn func()) error

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

func (w *Writer) SaveTemplateStats() error

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

func (w *Writer) SeriesID(key SeriesKey) 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

func (w *Writer) SeriesKeyByID(id SeriesID) (SeriesKey, bool)

SeriesKeyByID resolves a durable ID back to its identity from memory.

func (*Writer) Shutdown

func (w *Writer) Shutdown(ctx context.Context) error

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.

func (*Writer) Stop

func (w *Writer) Stop()

Stop preserves the prior internal lifecycle surface. Production shutdown uses Shutdown so deadline and persistence failures reach the process result.

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.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL