metrics

package
v0.16.2 Latest Latest
Warning

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

Go to latest
Published: Sep 29, 2026 License: Apache-2.0 Imports: 24 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type CPUUsage

type CPUUsage struct {
	Log    *slog.Logger
	Writer *VictoriaMetricsWriter
	Reader *VictoriaMetricsReader
	// contains filtered or unexported fields
}

func NewCPUUsage added in v0.3.0

func NewCPUUsage(log *slog.Logger, writer *VictoriaMetricsWriter, reader *VictoriaMetricsReader) *CPUUsage

NewCPUUsage creates a new CPUUsage with the given dependencies. Writer and Reader can be nil for environments without metrics collection.

func (*CPUUsage) CPUUsageDayAgo

func (m *CPUUsage) CPUUsageDayAgo(entity string, day units.Days) (float64, error)

Returns the cpu usage in cores for the given entity for a specific day in the past

func (*CPUUsage) CPUUsageLastHour

func (m *CPUUsage) CPUUsageLastHour(entity string) ([]UsageAtTime, error)

func (*CPUUsage) CPUUsageOver

func (m *CPUUsage) CPUUsageOver(entity string, interval string) (float64, error)

Returns the cpu usage in cores for the given entity over a given interval

func (*CPUUsage) CPUUsageOverDay

func (m *CPUUsage) CPUUsageOverDay(entity string) (float64, error)

func (*CPUUsage) CPUUsageOverLastHour

func (m *CPUUsage) CPUUsageOverLastHour(entity string) (float64, error)

func (*CPUUsage) CurrentCPUUsage

func (m *CPUUsage) CurrentCPUUsage(entity string) (float64, error)

Returns the cpu usage in cores for the given entity in the last minute

func (*CPUUsage) RecordUsage

func (m *CPUUsage) RecordUsage(ctx context.Context, entity string, windowStart, windowEnd time.Time, cpuUsec units.Microseconds, attrs map[string]string) error

func (*CPUUsage) Setup

func (m *CPUUsage) Setup() error

type Data

type Data struct {
	ResultType string   `json:"resultType"`
	Result     []Result `json:"result"`
}

type ErrorBreakdown

type ErrorBreakdown struct {
	StatusCode int
	Count      int64
	Percentage float64
}

ErrorBreakdown represents error counts by status code

type Fanout added in v0.16.0

type Fanout struct {
	// contains filtered or unexported fields
}

Fanout delivers every batch to each attached sink.

Sinks can be attached after collectors have started emitting. The runtime's operational gauges (control-process memory, etcd health, node usage) start early in boot and write to the embedded VictoriaMetrics from their first tick, while the sink that ships them off the cluster only exists once the managed-metrics vmagent is up, several boot stages later. Attaching late is how that sink joins without the collectors being restarted or made aware of boot order.

func NewFanout added in v0.16.0

func NewFanout(sinks ...PointWriter) *Fanout

NewFanout creates a Fanout that starts with the given sinks.

func (*Fanout) Attach added in v0.16.0

func (f *Fanout) Attach(sink PointWriter)

Attach adds a sink. Batches written after Attach returns reach it; a batch already in flight may not. A nil sink is ignored.

func (*Fanout) Detach added in v0.16.0

func (f *Fanout) Detach(sink PointWriter)

Detach removes a sink. It takes the write lock, and WritePoints holds the read lock for the whole delivery, so once Detach returns no batch is still being handed to the removed sink; the caller can close it without a late write landing after its final flush. Detaching a sink that was never attached is a no-op.

func (*Fanout) WritePoints added in v0.16.0

func (f *Fanout) WritePoints(ctx context.Context, points []MetricPoint) error

WritePoints hands the batch to every sink. One sink failing does not stop the others from receiving the batch; all failures are returned joined.

The read lock is held across delivery, not just while reading the sink list. Sinks only buffer (the writers flush on their own goroutines), so the hold is brief, and it is what gives Detach its guarantee. The corollary is that a sink must not call Attach or Detach from inside its own WritePoints; those are lifecycle operations for whoever owns the sink, and a sink that tried to detach itself mid-delivery would be waiting on its own read lock.

type HTTPMetrics

type HTTPMetrics struct {
	Log    *slog.Logger
	Writer *VictoriaMetricsWriter
	Reader *VictoriaMetricsReader
	// contains filtered or unexported fields
}

HTTPMetrics tracks HTTP request metrics for applications using VictoriaMetrics

func NewHTTPMetrics added in v0.3.0

func NewHTTPMetrics(log *slog.Logger, writer *VictoriaMetricsWriter, reader *VictoriaMetricsReader) *HTTPMetrics

NewHTTPMetrics creates a new HTTPMetrics with the given dependencies. Writer and Reader can be nil for environments without metrics collection.

func (*HTTPMetrics) Close

func (h *HTTPMetrics) Close() error

Close is a no-op for VictoriaMetrics (writer handles its own lifecycle)

func (*HTTPMetrics) ErrorsLastHour

func (h *HTTPMetrics) ErrorsLastHour(app string) ([]ErrorBreakdown, error)

ErrorsLastHour returns breakdown of errors by status code for the last hour

func (*HTTPMetrics) RPSLastMinute

func (h *HTTPMetrics) RPSLastMinute(app string) (float64, error)

RPSLastMinute returns requests per second for the last minute

func (*HTTPMetrics) RecordRequest

func (h *HTTPMetrics) RecordRequest(ctx context.Context, req HTTPRequest) error

RecordRequest records an HTTP request as metrics in VictoriaMetrics.

Only maxPathBoundCounters counters are kept for real paths. Anything beyond that is still counted, but reported with path and method "other", so an app serving unique URLs -- or a client walking random ones -- costs a fixed amount of memory here and a fixed number of series downstream.

func (*HTTPMetrics) Setup

func (h *HTTPMetrics) Setup() error

func (*HTTPMetrics) StatsLastHour

func (h *HTTPMetrics) StatsLastHour(app string) ([]RequestStats, error)

StatsLastHour returns request statistics for the last hour in 1-minute buckets

func (*HTTPMetrics) TopPaths

func (h *HTTPMetrics) TopPaths(app string, limit int) ([]PathStats, error)

TopPaths returns the busiest paths for an app over the last recentMinutes.

Served from the in-process counters rather than from VictoriaMetrics. The query it replaced asked for topk by path and then issued two more queries per path for the average duration and the error rate, all of them reading the `path` label -- the same label the cardinality cap collapses to "other". So the cap and this view could not both be right. Reading the counters directly removes the conflict, and removes 1+2N round trips per `miren app status`.

Only this node's ingress traffic is counted. That is the whole cluster today, since the ingress runs on the coordinator alone; it stops being true the day runners serve ingress, and this has to aggregate across them then.

type HTTPRequest

type HTTPRequest struct {
	Timestamp    time.Time
	App          string
	Method       string
	Path         string
	StatusCode   int
	DurationMs   int64
	ResponseSize int64
}

HTTPRequest represents a single HTTP request for metrics

type Labeled added in v0.16.0

type Labeled struct {
	Sink   PointWriter
	Labels map[string]string
}

Labeled stamps constant labels onto every point before passing it on.

This is how series pooled from many clusters into one store stay distinct. The collectors themselves carry no cluster identity, because inside a cluster's own store the cluster is implicit; the identity belongs to the shipping path, so it is applied by the sink that ships rather than by every emitter. A label the point already carries is kept, so a collector that knows its own runner is not overwritten by a process-wide default.

Label names are compared after the same sanitization the writer applies when formatting, so a collector's "miren.runner" and a constant "miren_runner" are recognized as the same label rather than emitted twice.

func (*Labeled) WritePoints added in v0.16.0

func (l *Labeled) WritePoints(ctx context.Context, points []MetricPoint) error

WritePoints copies each point's labels, merging in the constant labels, and writes the result to the wrapped sink. The caller's points are not mutated; collectors share one label map across every point in a batch.

type MemoryUsage

type MemoryUsage struct {
	Log    *slog.Logger
	Writer *VictoriaMetricsWriter
	Reader *VictoriaMetricsReader
	// contains filtered or unexported fields
}

func NewMemoryUsage added in v0.3.0

func NewMemoryUsage(log *slog.Logger, writer *VictoriaMetricsWriter, reader *VictoriaMetricsReader) *MemoryUsage

NewMemoryUsage creates a new MemoryUsage with the given dependencies. Writer and Reader can be nil for environments without metrics collection.

func (*MemoryUsage) RecordUsage

func (m *MemoryUsage) RecordUsage(
	ctx context.Context,
	entity string,
	ts time.Time,
	memory units.Bytes,
	attrs map[string]string,
) error

func (*MemoryUsage) Setup

func (m *MemoryUsage) Setup() error

func (*MemoryUsage) UsageLastHour

func (m *MemoryUsage) UsageLastHour(entity string) ([]MemoryUsageAtTime, error)

type MemoryUsageAtTime

type MemoryUsageAtTime struct {
	Timestamp time.Time
	Memory    units.Bytes
}

type MetricPoint

type MetricPoint struct {
	Name      string
	Labels    map[string]string
	Value     float64
	Timestamp time.Time
}

MetricPoint represents a single metric data point

type NodeUsage added in v0.15.0

type NodeUsage struct {
	Log    *slog.Logger
	Writer PointWriter

	// NodeID is the node entity ID these series describe, and RunnerID is the
	// runner identifier for the same host. Both are emitted: the entity ID is
	// what joins to the node record, the runner ID is what an operator typed.
	NodeID   string
	RunnerID string

	// DataPath is the filesystem whose capacity is reported. Empty skips the
	// storage series rather than reporting the root filesystem, which would be
	// a different disk on most real deployments.
	DataPath string
	// contains filtered or unexported fields
}

NodeUsage records what a whole host is doing, as opposed to what any single sandbox on it is doing.

Two questions need this. First, "how much of this machine is left?" -- a workload burning 200% of a core is unremarkable on a 64-core box and an emergency on a 2-core one, and nothing in the per-sandbox series says which box it is. Second, "is the load even an app?" -- containerd, buildkit, the registry, the log pipeline and miren itself run outside every sandbox cgroup, so they are invisible in the per-sandbox series. Subtracting the sum of sandbox usage from the host total is what surfaces them, and that subtraction needs a host total to start from.

Series are labeled by node so they join against the miren.node label the sandbox collectors now attach.

func NewNodeUsage added in v0.15.0

func NewNodeUsage(log *slog.Logger, writer PointWriter, nodeID, runnerID, dataPath string) *NodeUsage

NewNodeUsage creates a NodeUsage collector. Writer may be nil for environments without metrics collection, in which case Monitor is a no-op.

func (*NodeUsage) Monitor added in v0.15.0

func (n *NodeUsage) Monitor(ctx context.Context)

Monitor samples host stats every defaultNodeUsageInterval and pushes one batch of points per tick until ctx is cancelled. Mirrors the cadence and lifecycle of RuntimeMemory.Monitor.

type PathStats

type PathStats struct {
	Path          string
	Count         int64
	AvgDurationMs float64
	ErrorRate     float64
}

PathStats represents statistics for a specific path

type PointWriter added in v0.16.0

type PointWriter interface {
	WritePoints(ctx context.Context, points []MetricPoint) error
}

PointWriter is the sink the runtime's operational collectors emit to. The collectors only ever batch-write, so this is the whole surface they need; keeping it this narrow is what lets a collector be pointed at one store, a fan-out of several, or a labeling wrapper without knowing which.

type ProcessInfo added in v0.16.0

type ProcessInfo struct {
	Log    *slog.Logger
	Writer PointWriter

	// Entity is the value of the "entity" label on every emitted series.
	Entity string

	// StartTime is what process_start_time_seconds reports.
	StartTime time.Time

	// Version, Commit and Channel are the labels on miren_build_info. Version
	// and Commit are emitted as-is: a dev build reports "unknown", which is
	// worth seeing in the store, since an unknown build in a fleet is skew.
	// Channel is omitted when the build has none.
	Version string
	Commit  string
	Channel string
}

ProcessInfo publishes the identity of the miren control process: when it started and which build it is. Both are constants for the life of the process, and that is the point. A restart shows up as a step in process_start_time_seconds, an upgrade as a change in the commit label of miren_build_info, and either one explains a heap curve that resets or a goroutine count that gaps without anyone having to line the gap up against the commit log.

The series carry the same entity label as RuntimeMemory so the two can be joined, and pick up cluster and runner identity at the shipping layer like every other operational series.

func NewProcessInfo added in v0.16.0

func NewProcessInfo(log *slog.Logger, writer PointWriter) *ProcessInfo

NewProcessInfo creates a ProcessInfo collector describing this binary. Writer may be nil for environments without metrics collection, in which case Monitor is a no-op.

func (*ProcessInfo) Emit added in v0.16.0

func (p *ProcessInfo) Emit(ctx context.Context) error

Emit pushes one sample of each identity series. Both are constants, so an extra push between ticks is harmless; callers use it to prime a sink that attached after the last one. A nil Writer is a no-op.

func (*ProcessInfo) Monitor added in v0.16.0

func (p *ProcessInfo) Monitor(ctx context.Context)

Monitor pushes the identity series immediately and then once per defaultProcessInfoInterval until ctx is cancelled. The first push is not gated on the ticker so a process that restarts faster than the interval still records the new start time in the embedded store. The sink that ships off-cluster attaches later in boot and does not see that push; the boot stage that attaches it calls Emit so the central store gets the same coverage.

type QueryResult

type QueryResult struct {
	Status string `json:"status"`
	Data   Data   `json:"data"`
}

QueryResult represents the result of a MetricsQL query

type RequestStats

type RequestStats struct {
	Time          time.Time
	Count         int64
	AvgDurationMs float64
	P95DurationMs float64
	P99DurationMs float64
	ErrorRate     float64
}

RequestStats represents aggregated request statistics

type Result

type Result struct {
	Metric map[string]string `json:"metric"`
	Value  []any             `json:"value,omitempty"`  // For instant queries: [timestamp, value_string]
	Values [][]any           `json:"values,omitempty"` // For range queries: [[timestamp, value_string], ...]
}

type RuntimeMemory added in v0.11.1

type RuntimeMemory struct {
	Log    *slog.Logger
	Writer PointWriter

	// Entity is the value of the "entity" label on every emitted series.
	Entity string
}

RuntimeMemory collects Go runtime memory stats for the miren control process itself and pushes them to VictoriaMetrics in the same push-only style as the per-sandbox collectors (MemoryUsage/CPUUsage). The control process is not covered by any per-sandbox cgroup, so today there is no visibility into the coordinator's own heap/RSS growth — these series are that visibility.

All series carry entity="miren/control" so they sit alongside the per-app memory_usage_bytes{entity="app/..."} series and can be compared directly. The key diagnostic is the gap between process_resident_memory_bytes (the whole control process, including off-heap: cgo, mmap'd bbolt/etcd, the buildkit content store) and go_mem_heap_inuse_bytes (Go-managed heap only):

  • both balloon together -> Go heap; a pprof heap profile names the site.
  • RSS balloons, heap flat -> off-heap; pprof can't see it, look at bbolt/buildkit.

func NewRuntimeMemory added in v0.11.1

func NewRuntimeMemory(log *slog.Logger, writer PointWriter) *RuntimeMemory

NewRuntimeMemory creates a RuntimeMemory collector. Writer may be nil for environments without metrics collection, in which case Monitor is a no-op.

func (*RuntimeMemory) Monitor added in v0.11.1

func (r *RuntimeMemory) Monitor(ctx context.Context)

Monitor samples runtime memory every defaultRuntimeMemoryInterval and pushes one batch of points per tick until ctx is cancelled. Mirrors the cadence and lifecycle of (*sandbox.Metrics).Monitor.

type TimeSeriesPoint

type TimeSeriesPoint struct {
	Timestamp time.Time
	Value     float64
}

TimeSeriesPoint represents a single point in a time series

type UsageAtTime

type UsageAtTime struct {
	Timestamp time.Time
	Cores     float64
}

type VictoriaMetricsReader

type VictoriaMetricsReader struct {
	Log     *slog.Logger
	Address string
	Timeout time.Duration
	// contains filtered or unexported fields
}

VictoriaMetricsReader reads metrics from VictoriaMetrics using MetricsQL

func NewVictoriaMetricsReader

func NewVictoriaMetricsReader(log *slog.Logger, address string, timeout time.Duration) *VictoriaMetricsReader

NewVictoriaMetricsReader creates a new VictoriaMetrics reader

func (*VictoriaMetricsReader) GetAverage

func (r *VictoriaMetricsReader) GetAverage(ctx context.Context, metricName string, labels map[string]string, window string) (float64, error)

GetAverage calculates the average value over time

func (*VictoriaMetricsReader) GetLatestValue

func (r *VictoriaMetricsReader) GetLatestValue(ctx context.Context, metricName string, labels map[string]string) (float64, error)

GetLatestValue retrieves the latest value for a metric

func (*VictoriaMetricsReader) GetQuantile

func (r *VictoriaMetricsReader) GetQuantile(ctx context.Context, quantile float64, metricName string, labels map[string]string, window string) (float64, error)

GetQuantile calculates a quantile over time

func (*VictoriaMetricsReader) GetRate

func (r *VictoriaMetricsReader) GetRate(ctx context.Context, metricName string, labels map[string]string, window string) (float64, error)

GetRate calculates the rate of a counter metric over a time range

func (*VictoriaMetricsReader) GetTimeSeries

func (r *VictoriaMetricsReader) GetTimeSeries(ctx context.Context, metricName string, labels map[string]string, start, end time.Time, step string) ([]TimeSeriesPoint, error)

GetTimeSeries retrieves time series data for a metric over a time range

func (*VictoriaMetricsReader) InstantQuery

func (r *VictoriaMetricsReader) InstantQuery(ctx context.Context, query string, ts time.Time) (*QueryResult, error)

InstantQuery executes an instant MetricsQL query (point-in-time)

func (*VictoriaMetricsReader) RangeQuery

func (r *VictoriaMetricsReader) RangeQuery(ctx context.Context, query string, start, end time.Time, step string) (*QueryResult, error)

RangeQuery executes a range MetricsQL query (time series)

type VictoriaMetricsWriter

type VictoriaMetricsWriter struct {
	Log     *slog.Logger
	Address string // e.g., "localhost:8428"
	Timeout time.Duration
	// contains filtered or unexported fields
}

VictoriaMetricsWriter writes metrics to VictoriaMetrics using the import API

func NewVictoriaMetricsWriter

func NewVictoriaMetricsWriter(log *slog.Logger, address string, timeout time.Duration, opts ...WriterOption) *VictoriaMetricsWriter

NewVictoriaMetricsWriter creates a new VictoriaMetrics writer

func (*VictoriaMetricsWriter) Close

func (w *VictoriaMetricsWriter) Close() error

Close stops the background flush routine and flushes remaining data

func (*VictoriaMetricsWriter) Flush

func (w *VictoriaMetricsWriter) Flush()

Flush manually triggers a flush of buffered metrics

func (*VictoriaMetricsWriter) Start

func (w *VictoriaMetricsWriter) Start()

Start begins the background flush routine. It is idempotent and safe to call multiple times.

func (*VictoriaMetricsWriter) WritePoint

func (w *VictoriaMetricsWriter) WritePoint(ctx context.Context, point MetricPoint) error

WritePoint adds a metric point to the buffer

func (*VictoriaMetricsWriter) WritePoints

func (w *VictoriaMetricsWriter) WritePoints(ctx context.Context, points []MetricPoint) error

WritePoints adds multiple metric points to the buffer

type WriterOption added in v0.13.0

type WriterOption func(*VictoriaMetricsWriter)

WriterOption customizes a VictoriaMetricsWriter at construction.

func WithHTTPClient added in v0.13.0

func WithHTTPClient(client *http.Client) WriterOption

WithHTTPClient replaces the writer's default HTTP client.

Callers that cannot reach VictoriaMetrics with a plain client supply their own. A distributed runner ships metrics through the coordinator rather than dialing VictoriaMetrics directly, which means an HTTP/3 transport carrying a credential. Putting that in the client's RoundTripper keeps this writer unaware of how it is authenticated, so the credential has one home rather than being threaded through every send.

Jump to

Keyboard shortcuts

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