ingest

package
v0.4.0-beta.2 Latest Latest
Warning

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

Go to latest
Published: Aug 24, 2026 License: MIT Imports: 38 Imported by: 0

Documentation

Index

Constants

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

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

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

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

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

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

func ParseSeverity(level string) int

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.

func (*Batch) Priority

func (b *Batch) Priority() bool

Priority reports whether the batch is protected from soft-backpressure drops. Used by Submit() to decide whether to enqueue at >= 90% fullness.

type BatchSink added in v0.5.0

type BatchSink interface {
	Enqueue(batch any) error
}

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

type ExemplarTraceStats struct {
	Truncated bool
	Retained  int
	Observed  int
}

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

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.

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

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.

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

func (p *Pipeline) OfferToDLQ(b *Batch) error

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

func (p *Pipeline) SetDLQ(sink BatchSink)

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

func (p *Pipeline) SetPerTenantCap(n int)

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

func (p *Pipeline) SetStoreMinSeverity(level int)

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

func (p *Pipeline) Start(ctx context.Context)

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 signals workers to exit and blocks until in-flight batches have been drained from the channel. Idempotent.

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

func (p *Pipeline) TenantDropped() int64

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

func NewSampler(rate float64, alwaysOnErrors bool, latencyThresholdMs float64) *Sampler

NewSampler creates a Sampler with the given parameters.

func (*Sampler) ShouldSample

func (s *Sampler) ShouldSample(serviceName string, isError bool, durationMs float64) bool

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.

func (*Sampler) Stats

func (s *Sampler) Stats() (int64, int64)

Stats returns (seen, dropped) counters for metrics.

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

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

Jump to

Keyboard shortcuts

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