Documentation
¶
Index ¶
- Constants
- Variables
- func NewGRPCAuthInterceptors(o GRPCAuthOptions) (grpc.UnaryServerInterceptor, grpc.StreamServerInterceptor)
- func ParseSeverity(level string) int
- type Batch
- type BatchSink
- type DLQBatchEnvelope
- type DLQBatchPayload
- type ExemplarConfig
- type ExemplarMetrics
- type ExemplarPolicy
- func (p *ExemplarPolicy) AdmitSpan(in ExemplarSpan) bool
- func (p *ExemplarPolicy) DLQDisabled() bool
- func (p *ExemplarPolicy) GlobalTraceSlots(ts time.Time) int
- func (p *ExemplarPolicy) NewReservation() *ExemplarReservation
- func (p *ExemplarPolicy) ReserveLog(res *ExemplarReservation, tenant, service, severity string, ts time.Time, ...) bool
- func (p *ExemplarPolicy) ReserveSpan(res *ExemplarReservation, tenant, service, traceID string, ts time.Time, ...) bool
- func (p *ExemplarPolicy) ReserveSynthesizedLog(res *ExemplarReservation, tenant, service, traceID, spanID, severity string, ...) bool
- func (p *ExemplarPolicy) SelectedTraces(tenant, service string, ts time.Time) []string
- func (p *ExemplarPolicy) ServiceWindowBytes(tenant, service string, ts time.Time) (committed, reserved int64)
- func (p *ExemplarPolicy) SetShedding(s storage.SheddingState)
- func (p *ExemplarPolicy) Shedding() storage.SheddingState
- func (p *ExemplarPolicy) SynthesizedLogEligible(severity string) bool
- func (p *ExemplarPolicy) TraceStats(tenant, service, traceID string, ts time.Time) (ExemplarTraceStats, bool)
- func (p *ExemplarPolicy) WindowBytes(ts time.Time) (committed, reserved int64)
- type ExemplarReservation
- type ExemplarSpan
- type ExemplarTraceStats
- type GRPCAuthOptions
- type HTTPHandler
- type LogsServer
- func (s *LogsServer) Export(ctx context.Context, req *collogspb.ExportLogsServiceRequest) (*collogspb.ExportLogsServiceResponse, error)
- func (s *LogsServer) SetAggregateEngine(e *aggregate.Engine)
- func (s *LogsServer) SetExemplarPolicy(p *ExemplarPolicy)
- func (s *LogsServer) SetLogCallback(cb func(storage.Log))
- func (s *LogsServer) SetPipeline(p *Pipeline)
- func (s *LogsServer) SetResourceRegistry(r *topology.Registry)
- type MetricsServer
- func (s *MetricsServer) Export(ctx context.Context, req *colmetricspb.ExportMetricsServiceRequest) (*colmetricspb.ExportMetricsServiceResponse, error)
- func (s *MetricsServer) SetAggregateEngine(e *aggregate.Engine)
- func (s *MetricsServer) SetMetricCallback(cb func(tsdb.RawMetric))
- func (s *MetricsServer) SetResourceRegistry(r *topology.Registry)
- type Pipeline
- func (p *Pipeline) OfferToDLQ(b *Batch) error
- func (p *Pipeline) SetDLQ(sink BatchSink)
- func (p *Pipeline) SetPerTenantCap(n int)
- func (p *Pipeline) SetStoreMinSeverity(level int)
- func (p *Pipeline) Shutdown(ctx context.Context) error
- func (p *Pipeline) Start(ctx context.Context)
- func (p *Pipeline) Stats() PipelineStats
- func (p *Pipeline) Stop()
- func (p *Pipeline) Submit(b *Batch) (SubmitOutcome, error)
- func (p *Pipeline) TenantDropped() int64
- type PipelineConfig
- type PipelineStats
- type Sampler
- type SignalType
- type SubmitOutcome
- type TraceServer
- func (s *TraceServer) Export(ctx context.Context, req *coltracepb.ExportTraceServiceRequest) (*coltracepb.ExportTraceServiceResponse, error)
- func (s *TraceServer) SetAggregateEngine(e *aggregate.Engine)
- func (s *TraceServer) SetExemplarPolicy(p *ExemplarPolicy)
- func (s *TraceServer) SetLogCallback(cb func(storage.Log))
- func (s *TraceServer) SetPipeline(p *Pipeline)
- func (s *TraceServer) SetResourceRegistry(r *topology.Registry)
- func (s *TraceServer) SetSampler(sm *Sampler)
- func (s *TraceServer) SetSpanCallback(cb func(storage.Span))
- func (s *TraceServer) SetTopologyObserver(cb func(tenant, traceID, spanID, parentSpanID, service string))
Constants ¶
const ( DefaultExemplarTracesPerServiceWindow = 25 DefaultExemplarTracesGlobalWindow = 1500 DefaultExemplarBytesPerServiceWindow = 512 * 1024 DefaultExemplarBytesGlobalWindow = 3 * 1024 * 1024 DefaultExemplarHealthyRate = 0.005 DefaultExemplarStratumTopK = 5 DefaultExemplarLogsErrorPerServiceWindow = 50 DefaultExemplarLogsWarnPerServiceWindow = 20 DefaultExemplarMaxSpansPerTrace = 500 DefaultExemplarMaxBytesPerTrace = 256 * 1024 DefaultExemplarSynthLogsPerSpan = 8 DefaultExemplarSynthLogsPerTrace = 64 )
Exemplar budget defaults, frozen in #161's resolution.
const DLQBatchType = "batch"
DLQBatchType is the `type` discriminator of the typed DLQ envelope used for a whole failed pipeline batch. It extends the existing {"type":"logs|spans|traces|metrics","data":[...]} family with a fifth shape whose `data` is an object rather than an array, because the three signal slices of one Batch must be replayed through a single BatchCreateAll transaction to preserve the Trace->Span->Log FK ordering the pipeline enforces on the primary path.
Variables ¶
var ErrDLQFull = errors.New("dead letter queue at capacity")
ErrDLQFull is the sentinel a BatchSink returns when it refuses a batch because its bounded budget is exhausted, as opposed to the write itself failing. *queue.DeadLetterQueue evicts FIFO rather than refusing, so this is produced only by sinks that choose strict bounding — it exists so the exemplar loss metric can tell "nowhere to put it" from "the write broke".
var ErrNoDLQ = errors.New("no dead letter queue configured")
ErrNoDLQ is returned by OfferToDLQ when no DLQ sink is wired. The batch is gone; the caller must count it as permanent loss.
var ErrQueueFull = errors.New("ingest pipeline at capacity")
ErrQueueFull is returned by Submit when the queue is at hard capacity (100% full). Callers should map this to gRPC RESOURCE_EXHAUSTED or HTTP 429 with a Retry-After hint so OTLP clients back off cleanly.
Functions ¶
func NewGRPCAuthInterceptors ¶ added in v0.5.0
func NewGRPCAuthInterceptors(o GRPCAuthOptions) (grpc.UnaryServerInterceptor, grpc.StreamServerInterceptor)
NewGRPCAuthInterceptors returns the unary and stream interceptors enforcing bearer authentication on the OTLP gRPC listener.
A tenant key binds absolutely: the authenticated tenant is pinned onto the call context, and an `x-tenant-id` metadata value that disagrees is ignored and counted rather than honoured. An operator key authenticates only — tenant precedence stays (metadata → trusted resource attribute → DEFAULT_TENANT), so an existing single-tenant deployment that adds API_KEY sees no change in where its rows land.
Both results are nil when authentication is not configured, so the caller installs nothing and pays nothing.
func ParseSeverity ¶
ParseSeverity is the exported wrapper for parseSeverity. Used by main.go to translate the STORE_MIN_SEVERITY env value into the integer rank the pipeline's second-tier filter expects.
Types ¶
type Batch ¶
type Batch struct {
Type SignalType
Tenant string
Traces []storage.Trace
Spans []storage.Span
Logs []storage.Log
// Priority flags. Errors and slow traces are protected from soft
// backpressure drops — they may still be rejected at hard capacity.
HasError bool
HasSlow bool
// Optional per-record callbacks invoked after a successful DB write.
// In production these feed GraphRAG ingestion. Nil callbacks are
// skipped silently.
SpanCallback func(storage.Span)
LogCallback func(storage.Log)
// Reservation carries the exemplar bytes reserved for the rows in this
// batch (#201 Q4). submitExemplars commits it when a destination accepts
// the batch and releases it when none does. Nil outside aggregate mode,
// and every method on it is nil-safe.
Reservation *ExemplarReservation
// contains filtered or unexported fields
}
Batch is the unit of work flowing through the async ingest Pipeline. One Batch corresponds to the persistable output of a single OTLP Export() call. Trace insertion ordering (Traces → Spans → Logs) is honored by the worker that processes the batch — packaging the three slices together preserves the FK invariant the synchronous path already enforces.
type BatchSink ¶ added in v0.5.0
BatchSink is the slice of the Dead Letter Queue the Pipeline depends on. *queue.DeadLetterQueue satisfies it. Declaring it here (rather than importing internal/queue) keeps the package layering one-directional and lets tests inject a fake without touching the filesystem.
type DLQBatchEnvelope ¶ added in v0.5.0
type DLQBatchEnvelope struct {
Type string `json:"type"`
Data DLQBatchPayload `json:"data"`
}
DLQBatchEnvelope is the on-disk form of a batch the pipeline could not persist. main.go's replay handler decodes it and re-runs BatchCreateAll.
type DLQBatchPayload ¶ added in v0.5.0
type DLQBatchPayload struct {
Tenant string `json:"tenant,omitempty"`
Signal string `json:"signal,omitempty"`
Traces []storage.Trace `json:"traces,omitempty"`
Spans []storage.Span `json:"spans,omitempty"`
Logs []storage.Log `json:"logs,omitempty"`
}
DLQBatchPayload is the `data` member of a DLQBatchType envelope.
type ExemplarConfig ¶ added in v0.5.0
type ExemplarConfig struct {
// TracesPerServiceWindow is the unified per-service/window trace budget
// filled by priority ERROR/FATAL > slow > healthy. Default 25.
TracesPerServiceWindow int
// TracesGlobalWindow bounds selected traces across all services in one
// window. Default 1500.
TracesGlobalWindow int
// BytesPerServiceWindow / BytesGlobalWindow bound the bytes actually handed
// to persistence. Counts and bytes both bind; first breach wins.
BytesPerServiceWindow int64 // Default 512 KiB
BytesGlobalWindow int64 // Default 8 MiB
// HealthyRate is the stateless hash-threshold eligibility target for
// healthy traces. Caps dominate under load. Default 0.005 (0.5%).
HealthyRate float64
// StratumTopK bounds how many exemplars one (operation × status class)
// stratum may hold, so a single repeated failure cannot monopolize the
// per-service budget. Default 5.
StratumTopK int
// LatencyThresholdMs is the shared definition of "slow"
// (SAMPLING_LATENCY_THRESHOLD_MS). A predicate, not a policy — it survives
// the sampler's retirement in every mode (#161).
LatencyThresholdMs float64
// LogsErrorPerServiceWindow is the raw ERROR/FATAL log budget per
// service/window. Default 50.
LogsErrorPerServiceWindow int
// LogsWarnEnabled opts WARN logs into raw retention. Off by default.
LogsWarnEnabled bool
// LogsWarnPerServiceWindow is the WARN budget when enabled. Default 20.
LogsWarnPerServiceWindow int
// MaxSpansPerTrace / MaxBytesPerTrace bound one retained trace so the
// complete-retained-trace contract (#163) cannot be turned into an
// unbounded write by a pathological trace. Breaching either forces
// truncation, which is persisted as truncated=true plus retained/observed
// span counts.
MaxSpansPerTrace int // Default 500
MaxBytesPerTrace int64 // Default 256 KiB
// SynthLogsPerSpan / SynthLogsPerTrace bound the logs synthesized from
// span events and span status (#201 Q3). Before this they were
// unmetered: a span carrying two hundred exception events wrote two
// hundred log rows that no budget had ever seen. Defaults 8 and 64.
SynthLogsPerSpan int
SynthLogsPerTrace int
// WindowSize is the tumbling budget window. Defaults to the aggregate
// engine's window so exemplar budgets and aggregate buckets share edges.
WindowSize time.Duration
// Metrics receives the policy's counters. nil = no-op.
Metrics ExemplarMetrics
}
ExemplarConfig carries the frozen #161 budgets. Zero fields fall back to the defaults so tests can construct a policy with one or two overrides.
type ExemplarMetrics ¶ added in v0.5.0
type ExemplarMetrics interface {
RecordExemplarEligible(signal, class string)
RecordExemplarDropped(signal, reason string)
RecordExemplarEviction()
RecordExemplarTruncation()
}
ExemplarMetrics is the policy's view of the metric surface. It mirrors the aggregate package's recorder indirection so tests (and any caller passing a nil *telemetry.Metrics) do not need a live Prometheus registry.
type ExemplarPolicy ¶ added in v0.5.0
type ExemplarPolicy struct {
// contains filtered or unexported fields
}
ExemplarPolicy is the bounded raw-retention gate for aggregate mode.
Concurrency: state is sharded by service so the 100–200-service target does not funnel every span through one mutex. Global (instance-wide) budgets live behind their own small mutex, taken only when a shard actually wants to reserve or release a slot — never on the common "healthy trace fails the hash threshold" path.
func NewExemplarPolicy ¶ added in v0.5.0
func NewExemplarPolicy(cfg ExemplarConfig) *ExemplarPolicy
NewExemplarPolicy builds a policy from cfg, filling unset fields with the #161 defaults.
func (*ExemplarPolicy) AdmitSpan ¶ added in v0.5.0
func (p *ExemplarPolicy) AdmitSpan(in ExemplarSpan) bool
AdmitSpan reports whether this span's raw row may be built and offered to persistence. It is the count-budget half of admission; ChargeSpan meters the bytes actually produced.
A trace already selected in this window always admits — that is the complete-retained-trace contract, and it is what makes multi-batch traces and at-least-once retries cohere with zero coordination.
func (*ExemplarPolicy) DLQDisabled ¶ added in v0.5.0
func (p *ExemplarPolicy) DLQDisabled() bool
DLQDisabled reports whether the exemplar DLQ fallback is closed. True at raw-off: deferring an exemplar to the DLQ at 95% is still writing to the disk that is about to fill (#201 Q5).
func (*ExemplarPolicy) GlobalTraceSlots ¶ added in v0.5.0
func (p *ExemplarPolicy) GlobalTraceSlots(ts time.Time) int
GlobalTraceSlots returns how many instance-wide trace slots the window containing ts is holding. Test and diagnostic surface.
func (*ExemplarPolicy) NewReservation ¶ added in v0.5.0
func (p *ExemplarPolicy) NewReservation() *ExemplarReservation
NewReservation opens a reservation. A nil policy yields a nil reservation, and every method below is nil-safe, so the legacy and shadow paths need no branches.
func (*ExemplarPolicy) ReserveLog ¶ added in v0.5.0
func (p *ExemplarPolicy) ReserveLog(res *ExemplarReservation, tenant, service, severity string, ts time.Time, size int) bool
ReserveLog reports whether a raw client log row may be persisted in aggregate mode. INFO/DEBUG are aggregate-only and never raw; ERROR/FATAL get a per-service window budget; WARN is opt-in (#161). Bytes are reserved against the same service/global byte budgets as spans — one pool, so a log flood cannot buy itself past the trace budget.
size is the variable-length payload (body + serialized attributes); the fixed row overhead is added here so callers do not have to know it.
func (*ExemplarPolicy) ReserveSpan ¶ added in v0.5.0
func (p *ExemplarPolicy) ReserveSpan(res *ExemplarReservation, tenant, service, traceID string, ts time.Time, size int) bool
ReserveSpan meters the bytes a retained span will hand to persistence and reports whether it fits. A false return means the row must be dropped: the byte budget bound before the count budget did, and the trace is marked truncated so the gap is visible in the persisted data rather than inferred.
res may be nil, in which case the charge commits immediately — the caller is asserting it has no submission boundary. Production paths always pass one.
func (*ExemplarPolicy) ReserveSynthesizedLog ¶ added in v0.5.0
func (p *ExemplarPolicy) ReserveSynthesizedLog(res *ExemplarReservation, tenant, service, traceID, spanID, severity string, ts time.Time, size int) bool
ReserveSynthesizedLog gates and METERS logs synthesized from span events and span status (#201 Q3).
These used to pass on a severity check alone, on the theory that they ride a span the policy had already budgeted. They do ride it — and they are not weightless. A span carrying two hundred exception events wrote two hundred log rows that no budget had ever seen, which is precisely the kind of unmetered write the 4.5 GiB main tier cannot absorb.
So a synthesized log now reserves len(body) + len(attributesJSON) + logRowFixedBytes against the selected trace's per-trace budget AND the shared per-service and global window budgets, under its own per-span and per-trace count caps. It does NOT consume the ordinary log-exemplar quota: that budget exists for logs a client sent, and charging synthesized rows against it would silently evict real ones.
A refusal drops the log, counts a reasoned drop, and marks the trace truncated so the gap appears in the persisted data.
func (*ExemplarPolicy) SelectedTraces ¶ added in v0.5.0
func (p *ExemplarPolicy) SelectedTraces(tenant, service string, ts time.Time) []string
SelectedTraces returns the trace IDs currently selected in the window containing ts, for a (tenant, service). Test/diagnostic surface: this is the set the determinism property is stated over.
func (*ExemplarPolicy) ServiceWindowBytes ¶ added in v0.5.0
func (p *ExemplarPolicy) ServiceWindowBytes(tenant, service string, ts time.Time) (committed, reserved int64)
ServiceWindowBytes returns the (committed, reserved) bytes of one (tenant, service, window) budget cell. Test and diagnostic surface.
func (*ExemplarPolicy) SetShedding ¶ added in v0.5.0
func (p *ExemplarPolicy) SetShedding(s storage.SheddingState)
SetShedding publishes the disk watchdog's current state to the policy. Called from the watchdog goroutine; safe on a nil policy.
func (*ExemplarPolicy) Shedding ¶ added in v0.5.0
func (p *ExemplarPolicy) Shedding() storage.SheddingState
Shedding reports the current shedding state. Safe on a nil policy.
func (*ExemplarPolicy) SynthesizedLogEligible ¶ added in v0.5.0
func (p *ExemplarPolicy) SynthesizedLogEligible(severity string) bool
SynthesizedLogEligible is the CHEAP half of ReserveSynthesizedLog: the severity floor plus the shedding ladder, with no locks and no accounting.
It exists so the OTLP path can refuse an INFO span event before marshaling its attributes. Passing it is necessary, not sufficient — the caller must still ReserveSynthesizedLog once it knows the row's size, and that call re-checks everything this one did.
func (*ExemplarPolicy) TraceStats ¶ added in v0.5.0
func (p *ExemplarPolicy) TraceStats(tenant, service, traceID string, ts time.Time) (ExemplarTraceStats, bool)
TraceStats returns the accounting for a selected trace so the caller can stamp truncated / retained / observed onto the persisted trace row (#163). The second return is false when the trace was never selected.
func (*ExemplarPolicy) WindowBytes ¶ added in v0.5.0
func (p *ExemplarPolicy) WindowBytes(ts time.Time) (committed, reserved int64)
WindowBytes returns the instance-wide (committed, reserved) bytes for the window containing ts. Test and diagnostic surface: the monotonicity property of #201 Q4 is stated over the committed number.
type ExemplarReservation ¶ added in v0.5.0
type ExemplarReservation struct {
// contains filtered or unexported fields
}
ExemplarReservation accumulates one submission unit's reserved bytes.
func (*ExemplarReservation) Bytes ¶ added in v0.5.0
func (r *ExemplarReservation) Bytes() int64
Bytes reports the total reserved so far. Diagnostic surface.
func (*ExemplarReservation) Commit ¶ added in v0.5.0
func (r *ExemplarReservation) Commit()
Commit converts every reserved byte into a committed byte. Call it exactly when a destination has accepted the rows. Idempotent.
func (*ExemplarReservation) Len ¶ added in v0.5.0
func (r *ExemplarReservation) Len() int
Len reports how many rows are reserved.
func (*ExemplarReservation) Merge ¶ added in v0.5.0
func (r *ExemplarReservation) Merge(other *ExemplarReservation)
Merge folds other into r and neutralizes other, so exactly one of them can ever settle the charges.
func (*ExemplarReservation) Release ¶ added in v0.5.0
func (r *ExemplarReservation) Release()
Release gives every reserved byte back. Legitimate only while the rows have NOT been accepted anywhere: dropped before submission, or refused by both the primary queue and the DLQ. Idempotent.
type ExemplarSpan ¶ added in v0.5.0
type ExemplarSpan struct {
Tenant string
Service string
TraceID string
Operation string
Status string
DurationMs float64
Timestamp time.Time
}
ExemplarSpan is one span offered to the policy. Everything here is already computed on the hot path, so admission adds no parsing.
type ExemplarTraceStats ¶ added in v0.5.0
ExemplarTraceStats is the persisted-side view of a retained trace.
type GRPCAuthOptions ¶ added in v0.5.0
type GRPCAuthOptions struct {
// Auth resolves credentials. A disabled Authenticator yields nil
// interceptors and the server keeps its current unauthenticated behaviour
// — which is the development default.
Auth *authn.Authenticator
// ExternalTenantMetadataKey is the proxy-injected identity key, honoured
// only under AUTH_TRUST_EXTERNAL. Compared lower-cased.
ExternalTenantMetadataKey string
// OnAuthFailure receives the failure reason for metrics. Optional.
OnAuthFailure func(reason string)
}
GRPCAuthOptions configures the OTLP gRPC authentication interceptors.
type HTTPHandler ¶
type HTTPHandler struct {
// contains filtered or unexported fields
}
HTTPHandler provides HTTP OTLP endpoints that delegate to the existing gRPC Export() methods.
func NewHTTPHandler ¶
func NewHTTPHandler(traces *TraceServer, logs *LogsServer, metrics *MetricsServer) *HTTPHandler
NewHTTPHandler creates an HTTP OTLP handler wrapping the existing gRPC servers.
func (*HTTPHandler) RegisterRoutes ¶
func (h *HTTPHandler) RegisterRoutes(mux *http.ServeMux)
RegisterRoutes registers the HTTP OTLP endpoints on the given mux.
func (*HTTPHandler) SetMaxBodyBytes ¶
func (h *HTTPHandler) SetMaxBodyBytes(n int64)
SetMaxBodyBytes configures the maximum pre-decompression request body size.
func (*HTTPHandler) SetMaxDecompressedBytes ¶
func (h *HTTPHandler) SetMaxDecompressedBytes(n int64)
SetMaxDecompressedBytes configures the maximum post-gzip body size. Rejecting larger inputs defends against compression bombs from authed clients.
func (*HTTPHandler) SetThrottleCallback ¶
func (h *HTTPHandler) SetThrottleCallback(fn func(signal string))
SetThrottleCallback wires a per-signal counter that increments every time a 429 is returned because the async ingest pipeline is at capacity. Used by main.go to feed `otelcontext_http_otlp_throttled_total{signal=…}`.
type LogsServer ¶
type LogsServer struct {
collogspb.UnimplementedLogsServiceServer
// contains filtered or unexported fields
}
func NewLogsServer ¶
func NewLogsServer(repo *storage.Repository, metrics *telemetry.Metrics, cfg *config.Config) *LogsServer
func (*LogsServer) Export ¶
func (s *LogsServer) Export(ctx context.Context, req *collogspb.ExportLogsServiceRequest) (*collogspb.ExportLogsServiceResponse, error)
Export handles incoming OTLP log data.
func (*LogsServer) SetAggregateEngine ¶ added in v0.5.0
func (s *LogsServer) SetAggregateEngine(e *aggregate.Engine)
SetAggregateEngine — see TraceServer.SetAggregateEngine.
func (*LogsServer) SetExemplarPolicy ¶ added in v0.5.0
func (s *LogsServer) SetExemplarPolicy(p *ExemplarPolicy)
SetExemplarPolicy — see TraceServer.SetExemplarPolicy.
func (*LogsServer) SetLogCallback ¶
func (s *LogsServer) SetLogCallback(cb func(storage.Log))
SetLogCallback sets the function to call when a new log is received.
func (*LogsServer) SetPipeline ¶
func (s *LogsServer) SetPipeline(p *Pipeline)
SetPipeline enables the async ingest pipeline for log export. Same semantics as TraceServer.SetPipeline.
func (*LogsServer) SetResourceRegistry ¶ added in v0.5.0
func (s *LogsServer) SetResourceRegistry(r *topology.Registry)
SetResourceRegistry — see TraceServer.SetResourceRegistry.
type MetricsServer ¶
type MetricsServer struct {
colmetricspb.UnimplementedMetricsServiceServer
// contains filtered or unexported fields
}
func NewMetricsServer ¶
func NewMetricsServer(repo *storage.Repository, metrics *telemetry.Metrics, aggregator *tsdb.Aggregator, cfg *config.Config) *MetricsServer
func (*MetricsServer) Export ¶
func (s *MetricsServer) Export(ctx context.Context, req *colmetricspb.ExportMetricsServiceRequest) (*colmetricspb.ExportMetricsServiceResponse, error)
Export handles incoming OTLP metrics data.
func (*MetricsServer) SetAggregateEngine ¶ added in v0.5.0
func (s *MetricsServer) SetAggregateEngine(e *aggregate.Engine)
SetAggregateEngine — see TraceServer.SetAggregateEngine.
func (*MetricsServer) SetMetricCallback ¶
func (s *MetricsServer) SetMetricCallback(cb func(tsdb.RawMetric))
SetMetricCallback sets the function to call when a new metric point is received.
func (*MetricsServer) SetResourceRegistry ¶ added in v0.5.0
func (s *MetricsServer) SetResourceRegistry(r *topology.Registry)
SetResourceRegistry — see TraceServer.SetResourceRegistry.
type Pipeline ¶
type Pipeline struct {
// contains filtered or unexported fields
}
Pipeline decouples OTLP Export() from synchronous DB writes. It holds a bounded buffered channel of Batches, a worker pool that drains the channel into the Repository, and Prometheus instruments that surface queue depth, drop counts, and rejection counts.
Lifecycle:
p := NewPipeline(repo, metrics, cfg) p.Start(ctx) defer p.Stop() // drains in-flight before returning p.Submit(batch)
func NewPipeline ¶
func NewPipeline(writer pipelineWriter, metrics *telemetry.Metrics, cfg PipelineConfig) *Pipeline
NewPipeline constructs a Pipeline with the given config, falling back to DefaultPipelineConfig() values for non-positive fields. The Pipeline does NOT start workers — call Start(ctx) when ready.
func (*Pipeline) OfferToDLQ ¶ added in v0.5.0
OfferToDLQ hands a batch the primary queue REFUSED (ErrQueueFull) to the DLQ sink, so a hard rejection can degrade to deferred storage instead of loss. It is the second half of the aggregate-mode ACK contract (#196 Q2): once the durable aggregate commit has landed, the Export is already acknowledged and the raw exemplars must go somewhere other than the client's retry loop.
Returns nil when the sink accepted the batch, ErrNoDLQ when no sink is wired, ErrDLQFull when the sink refused for capacity, and the sink's own error otherwise. A non-nil return means the batch is permanently lost.
The batch never held a queue slot or a byte reservation (Submit released both before returning ErrQueueFull), so there is nothing to release here.
func (*Pipeline) SetDLQ ¶ added in v0.5.0
SetDLQ wires the Dead Letter Queue that receives batches whose persist transaction failed. Without it the pipeline logs the failure and drops the batch (the pre-#194 behaviour, still the default so tests and the synchronous fallback path are unaffected).
Replay is at-least-once, not exactly-once. Traces and spans collapse duplicates on their composite unique indexes, so replaying an already-persisted batch is a no-op for them; logs have no stable OTLP identifier and a replay after a partially-visible commit can duplicate log rows. That trade is deliberate — losing the batch outright is worse.
Startup-only — call before Start().
func (*Pipeline) SetPerTenantCap ¶
SetPerTenantCap configures the maximum in-flight batches one tenant may hold in the queue (and currently being processed). 0 disables the cap. Once a tenant hits the cap, further healthy submissions from that tenant are dropped at Submit() time with reason "tenant_backpressure". Priority batches (errors/slow traces) bypass the cap.
Sized as a fraction of Capacity, e.g. Capacity/4 keeps any single tenant to 25% of queue capacity. Operators tune via INGEST_PIPELINE_PER_TENANT_CAP. Startup-only — call before Start().
func (*Pipeline) SetStoreMinSeverity ¶
SetStoreMinSeverity configures the second-tier severity gate applied at persist time. Logs below `level` are dropped from the BatchCreateAll write but still flow through the LogCallback so in-memory consumers (vectordb, GraphRAG Drain mining, anomaly correlation) keep working. 0 disables the second tier — every log surviving IngestMinSeverity at the receiver is also persisted (legacy behavior).
`level` is the integer rank from parseSeverity ("DEBUG"=10 .. "FATAL"=50). Startup-only — call before Start().
func (*Pipeline) Shutdown ¶ added in v0.5.0
Shutdown signals workers to exit, drains accepted batches, and reports any state that did not reach the database or durable DLQ before the deadline.
func (*Pipeline) Start ¶
Start spawns the worker pool. Safe to call once. Subsequent calls are no-ops; tests rely on this for reset semantics.
func (*Pipeline) Stats ¶
func (p *Pipeline) Stats() PipelineStats
Stats returns snapshot counters for tests and for telemetry that doesn't already use Prometheus instruments. The values are best-effort and not synchronized across atomics — sufficient for diagnostics.
Processed is the one exception and is load-bearing: it is incremented in process()'s outermost defer, so observing Processed >= N guarantees that N batches finished completely — writes attempted, callbacks fired, tenant slots and byte reservations released. That makes it the correct barrier for a caller (a test, a drain check) that wants to read state process() produced. Every other counter is a plain progress tally.
func (*Pipeline) Stop ¶
func (p *Pipeline) Stop()
Stop preserves the prior internal lifecycle surface. Production shutdown uses Shutdown so deadline and durability failures reach the process result.
func (*Pipeline) Submit ¶
func (p *Pipeline) Submit(b *Batch) (SubmitOutcome, error)
Submit enqueues a batch for asynchronous persistence. It returns a SubmitOutcome plus a nil error when the batch is accepted (or intentionally shed under soft backpressure), and ErrQueueFull when the queue is at hard capacity. Nil batches are no-ops.
Soft backpressure: when fullness >= SoftThreshold, healthy batches (Priority()==false) are dropped at the door and Submit returns (SubmitSoftDropped, nil) so the OTLP client sees a successful Export. Errors and slow traces always continue to the channel.
Hard backpressure: when the channel send fails (buffer at 100%) or the byte cap would be exceeded, Submit returns ErrQueueFull regardless of priority. The caller should translate this into a backpressure signal so the client retries with exponential backoff rather than tighter loops — except in AGGREGATE_MODE=aggregate, where the durable aggregate commit is already the ACK and the caller absorbs the rejection (see submitExemplars).
func (*Pipeline) TenantDropped ¶
TenantDropped reports the cumulative number of healthy submissions rejected because the submitting tenant was at the per-tenant cap. Distinct from RejectedFull (queue at hard capacity) and DroppedHealthy (soft-backpressure across the whole queue).
type PipelineConfig ¶
type PipelineConfig struct {
Capacity int // total queue depth across all signal types
Workers int // worker goroutines draining the queue
SoftThreshold float64 // fullness fraction above which healthy batches are dropped (0.0–1.0)
// MaxBytes caps the approximate bytes held by queued batches. Capacity
// alone cannot bound memory — a single Batch may carry arbitrarily
// large span/log payloads. At the cap, Submit rejects with ErrQueueFull
// even for priority batches: a 429 is recoverable by the client's retry
// loop, an OOM kill is not. <=0 falls back to the 512MB default.
MaxBytes int64
}
PipelineConfig holds the tunables for a Pipeline.
func DefaultPipelineConfig ¶
func DefaultPipelineConfig() PipelineConfig
DefaultPipelineConfig returns production-sized defaults.
type PipelineStats ¶
type PipelineStats struct {
Enqueued int64
Processed int64 // batches whose process() ran to completion (see Stats)
DroppedHealthy int64
RejectedFull int64
RejectedBytes int64 // batches rejected because the byte cap was exceeded
ProcessFailures int64
DLQEnqueued int64 // batches durably handed to the DLQ (persist failure or queue saturation)
DLQFailed int64 // batches the DLQ itself refused (data lost)
StoreFiltered int64 // logs dropped by STORE_MIN_SEVERITY at persist time
QueueDepth int
Capacity int
QueueBytes int64 // approx bytes currently reserved by in-flight batches
MaxBytes int64 // configured byte cap
}
PipelineStats is a snapshot of pipeline counters.
type Sampler ¶
type Sampler struct {
// contains filtered or unexported fields
}
Sampler decides whether a trace/span should be ingested. Always keeps: error traces, slow traces (duration > latencyThresholdMs), new services. Samples healthy traces at the configured rate using a per-service token bucket.
func NewSampler ¶
NewSampler creates a Sampler with the given parameters.
func (*Sampler) ShouldSample ¶
ShouldSample returns true if the trace should be ingested. isError: whether the trace/span has error status. durationMs: trace duration in milliseconds. serviceName: originating service.
type SignalType ¶
type SignalType uint8
SignalType identifies the OTLP signal a Batch carries. The label is exported on pipeline metrics so operators can attribute drops.
const ( SignalTraces SignalType = iota SignalLogs )
type SubmitOutcome ¶ added in v0.5.0
type SubmitOutcome uint8
SubmitOutcome is the success value of Submit. A hard rejection stays an error (ErrQueueFull); this type distinguishes the two ways Submit can succeed without the caller having to infer them from counters.
The distinction is load-bearing for the aggregate-shadow ordering state machine (#196 Q4): both outcomes are non-retry Export results, so both permit the shadow aggregate to be applied exactly once, but only SubmitEnqueued means the rows will actually reach the database.
const ( // SubmitEnqueued — the batch holds a queue slot and a worker will // persist it. Also returned for a nil or empty batch, which needs no // slot and loses nothing. SubmitEnqueued SubmitOutcome = iota // SubmitSoftDropped — the batch was intentionally shed at the door by // soft backpressure or the per-tenant admission cap. Already counted on // otelcontext_ingest_pipeline_dropped_total; the Export still succeeds. SubmitSoftDropped )
func (SubmitOutcome) String ¶ added in v0.5.0
func (o SubmitOutcome) String() string
String renders the outcome for logs and test failure messages.
type TraceServer ¶
type TraceServer struct {
coltracepb.UnimplementedTraceServiceServer
// contains filtered or unexported fields
}
func NewTraceServer ¶
func NewTraceServer(repo *storage.Repository, metrics *telemetry.Metrics, cfg *config.Config) *TraceServer
func (*TraceServer) Export ¶
func (s *TraceServer) Export(ctx context.Context, req *coltracepb.ExportTraceServiceRequest) (*coltracepb.ExportTraceServiceResponse, error)
Export handles incoming OTLP trace data.
func (*TraceServer) SetAggregateEngine ¶ added in v0.5.0
func (s *TraceServer) SetAggregateEngine(e *aggregate.Engine)
SetAggregateEngine enables aggregate accounting on this server. The reducer runs inside Export() ahead of the sampler and severity gates; passing nil (the AGGREGATE_MODE=legacy case) leaves the export path unchanged.
func (*TraceServer) SetExemplarPolicy ¶ added in v0.5.0
func (s *TraceServer) SetExemplarPolicy(p *ExemplarPolicy)
SetExemplarPolicy installs the bounded exemplar retention policy (#176). Wired only for AGGREGATE_MODE=aggregate, where it replaces the adaptive Sampler as the sole raw-retention gate. Passing nil (legacy / shadow) leaves the Sampler in charge and every byte of the export path unchanged.
func (*TraceServer) SetLogCallback ¶
func (s *TraceServer) SetLogCallback(cb func(storage.Log))
SetLogCallback sets the function to call when a new log is synthesized from a trace.
func (*TraceServer) SetPipeline ¶
func (s *TraceServer) SetPipeline(p *Pipeline)
SetPipeline enables the async ingest pipeline. When set, Export() returns to the caller as soon as the parsed batch is enqueued (or rejected), and persistence runs on the pipeline's worker pool. Pass nil to revert to the synchronous DB-write path.
func (*TraceServer) SetResourceRegistry ¶ added in v0.5.0
func (s *TraceServer) SetResourceRegistry(r *topology.Registry)
SetResourceRegistry wires the bounded resource registry (#279). Every resource batch registers its host identity ahead of the sampler; passing nil disables registration.
func (*TraceServer) SetSampler ¶
func (s *TraceServer) SetSampler(sm *Sampler)
SetSampler enables adaptive trace sampling. Pass nil to disable.
func (*TraceServer) SetSpanCallback ¶
func (s *TraceServer) SetSpanCallback(cb func(storage.Span))
SetSpanCallback sets the function to call when spans are persisted.
func (*TraceServer) SetTopologyObserver ¶ added in v0.3.1
func (s *TraceServer) SetTopologyObserver(cb func(tenant, traceID, spanID, parentSpanID, service string))
SetTopologyObserver wires a pre-sample hook invoked for every received span before the sampler runs, so the in-memory service map keeps cross-service flow direction even when sampling drops the spans that would have formed the edge. Pass nil to disable. Wired in main.go to graphrag.ObserveSpanTopology.