Documentation
¶
Overview ¶
Package sinks provides functionality to store monitored data in different ways.
At the moment we provide sink connectors for
- PostgreSQL and flavours,
- Prometheus,
- plain JSON files,
- and RPC servers.
To ensure the simultaneous storage of data in several storages, the `MultiWriter` class is implemented.
Feedback ¶
Sinks may optionally implement Feedbacker to report the epoch of the newest measurement they already hold for a source/metric pair, so that a stateful, resumable collector can continue from there instead of restarting from the current instant. The capability is negotiated at two levels: implementing the interface declares the sink kind capable, and CanFeedback answers for one specific pair.
- PostgresWriter answers from the measurement tables, bounded by retention.
- RPCWriter forwards the question to the remote server; servers that do not implement it are probed once and then left alone.
- MultiWriter reports the minimum epoch across capable sinks, so a resume never starves the sink that lags furthest behind.
- PrometheusWriter and JSONWriter deliberately do not implement it; see spec/design-sink-feedback.md for why.
Nothing in pgwatch queries feedback yet. Read the caller contract in that specification before wiring up a consumer.
Index ¶
- Variables
- func LoadTLSCredentials(CAFile string) (credentials.TransportCredentials, error)
- func NewPostgresSinkMigrator(ctx context.Context, connStr string) (db.Migrator, error)
- type CmdOpts
- type DbStorageSchemaType
- type ExistingPartitionInfo
- type Feedbacker
- type JSONWriter
- type MeasurementMessagePostgres
- type MetricsDefiner
- type MultiWriter
- func (mw *MultiWriter) AddWriter(w Writer)
- func (mw *MultiWriter) CanFeedback(sourceName, metricName string) bool
- func (mw *MultiWriter) Count() int
- func (mw *MultiWriter) DefineMetrics(metrics *metrics.Metrics) (err error)
- func (mw *MultiWriter) LastMeasurement(ctx context.Context, sourceName, metricName string) (int64, error)
- func (mw *MultiWriter) Migrate() (err error)
- func (mw *MultiWriter) NeedsMigration() (bool, error)
- func (mw *MultiWriter) SyncMetric(sourceName, metricName string, op SyncOp) (err error)
- func (mw *MultiWriter) Write(msg metrics.MeasurementEnvelope) (err error)
- type PostgresWriter
- func (pgw *PostgresWriter) AddDBUniqueMetricToListingTable(dbUnique, metric string) error
- func (pgw *PostgresWriter) CanFeedback(sourceName, metricName string) bool
- func (pgw *PostgresWriter) DeleteOldPartitions()
- func (pgw *PostgresWriter) EnsureBuiltinMetricDummies() (err error)
- func (pgw *PostgresWriter) EnsureMetricDummy(metric string) (err error)
- func (pgw *PostgresWriter) EnsureMetricTimePartsExist(metricPartBounds map[string]ExistingPartitionInfo) error
- func (pgw *PostgresWriter) EnsureMetricTimescale(pgPartBounds map[string]ExistingPartitionInfo) (err error)
- func (pgw *PostgresWriter) LastMeasurement(ctx context.Context, sourceName, metricName string) (int64, error)
- func (pgw *PostgresWriter) MaintainUniqueSources()
- func (pgw *PostgresWriter) Migrate() error
- func (pgw *PostgresWriter) NeedsMigration() (bool, error)
- func (pgw *PostgresWriter) ReadMetricSchemaType() (err error)
- func (pgw *PostgresWriter) SyncMetric(sourceName, metricName string, op SyncOp) error
- func (pgw *PostgresWriter) Write(msg metrics.MeasurementEnvelope) error
- type PromMetricCache
- type PrometheusWriter
- func (promw *PrometheusWriter) AddCacheEntry(dbUnique, metric string, msgArr metrics.MeasurementEnvelope)
- func (promw *PrometheusWriter) Collect(ch chan<- prometheus.Metric)
- func (promw *PrometheusWriter) DefineMetrics(metrics *metrics.Metrics) (err error)
- func (promw *PrometheusWriter) Describe(_ chan<- *prometheus.Desc)
- func (promw *PrometheusWriter) InitCacheEntry(dbUnique string)
- func (promw *PrometheusWriter) Println(v ...any)
- func (promw *PrometheusWriter) PurgeCacheEntry(dbUnique, metric string)
- func (promw *PrometheusWriter) SyncMetric(sourceName, metricName string, op SyncOp) error
- func (promw *PrometheusWriter) Write(msg metrics.MeasurementEnvelope) error
- func (promw *PrometheusWriter) WritePromMetrics(msg metrics.MeasurementEnvelope, ch chan<- prometheus.Metric) (written int, errorCount int)
- type RPCWriter
- func (rw *RPCWriter) CanFeedback(sourceName, metricName string) bool
- func (rw *RPCWriter) DefineMetrics(metrics *metrics.Metrics) error
- func (rw *RPCWriter) LastMeasurement(ctx context.Context, sourceName, metricName string) (int64, error)
- func (rw *RPCWriter) Ping() error
- func (rw *RPCWriter) SyncMetric(sourceName, metricName string, op SyncOp) error
- func (rw *RPCWriter) Write(msg metrics.MeasurementEnvelope) error
- type SyncOp
- type Writer
Constants ¶
This section is empty.
Variables ¶
var ErrFeedbackUnsupported = errors.New("sink does not support feedback for this source/metric")
ErrFeedbackUnsupported indicates that the sink cannot report a last-written epoch for the requested (sourceName, metricName) pair. It is not a failure: callers are expected to fall back to their default behaviour.
var ErrNeedsMigration = errors.New("sink database schema is outdated, please run migrations using `pgwatch config upgrade` command")
var ErrNoFeedbackData = errors.New("sink holds no measurements for this source/metric")
ErrNoFeedbackData indicates that the pair is supported but the sink holds no measurement for it yet.
Functions ¶
func LoadTLSCredentials ¶
func LoadTLSCredentials(CAFile string) (credentials.TransportCredentials, error)
Types ¶
type CmdOpts ¶
type CmdOpts struct {
Sinks []string `long:"sink" mapstructure:"sink" description:"URI where metrics will be stored, can be used multiple times" env:"PW_SINK"`
BatchingDelay time.Duration `` /* 170-byte string literal not displayed */
PartitionInterval string `` /* 202-byte string literal not displayed */
RetentionInterval string `` /* 161-byte string literal not displayed */
MaintenanceInterval string `` /* 273-byte string literal not displayed */
RealDbnameField string `` /* 151-byte string literal not displayed */
SystemIdentifierField string `` /* 169-byte string literal not displayed */
NoFeedback bool `` /* 205-byte string literal not displayed */
}
CmdOpts specifies the storage configuration to store metrics measurements
func (*CmdOpts) FeedbackEnabled ¶
FeedbackEnabled reports whether sinks may answer feedback queries. Feedback is on by default; go-flags booleans always default to false and can only be turned on, hence the negated NoFeedback flag.
type DbStorageSchemaType ¶
type DbStorageSchemaType int
const ( DbStorageSchemaPostgres DbStorageSchemaType = iota DbStorageSchemaTimescale )
type ExistingPartitionInfo ¶
type Feedbacker ¶
type Feedbacker interface {
// CanFeedback reports whether LastMeasurement can be answered for this
// pair. It must not perform I/O, must not block, and must be safe for
// concurrent use. A true result is advisory: LastMeasurement may still
// return ErrFeedbackUnsupported if state changed in between.
CanFeedback(sourceName, metricName string) bool
// LastMeasurement returns the epoch_ns (Unix nanoseconds) of the newest
// measurement the sink durably holds for the pair.
//
// Returns ErrFeedbackUnsupported when the pair cannot be answered, and
// ErrNoFeedbackData when the pair is supported but empty. Both are
// expected outcomes, not faults. The returned epoch is 0 whenever err is
// non-nil, and strictly positive whenever err is nil.
LastMeasurement(ctx context.Context, sourceName, metricName string) (int64, error)
}
Feedbacker is an optional interface that a Writer may implement to report back what it has already durably stored. It exists so that stateful, resumable collectors can continue from the last persisted measurement instead of restarting from the current instant.
Implementing Feedbacker declares the sink kind capable of feedback; CanFeedback declares whether one specific source/metric pair can be answered.
No pgwatch component calls these methods today. See spec/design-sink-feedback.md for the caller contract before wiring up a consumer.
type JSONWriter ¶
type JSONWriter struct {
// contains filtered or unexported fields
}
JSONWriter is a sink that writes metric measurements to a file in JSON format. It supports compression and rotation of output files. The default rotation is based on the file size (100Mb). JSONWriter is useful for debugging and testing purposes, as well as for integration with other systems, such as log aggregators, analytics systems, and data processing pipelines, ML models, etc.
func NewJSONWriter ¶
func NewJSONWriter(ctx context.Context, fname string) (*JSONWriter, error)
func (*JSONWriter) SyncMetric ¶
func (jw *JSONWriter) SyncMetric(_, _ string, _ SyncOp) error
func (*JSONWriter) Write ¶
func (jw *JSONWriter) Write(msg metrics.MeasurementEnvelope) error
type MetricsDefiner ¶
MetricDefiner is an interface for passing metric definitions to a sink.
type MultiWriter ¶
MultiWriter ensures the simultaneous storage of data in several storages.
func (*MultiWriter) AddWriter ¶
func (mw *MultiWriter) AddWriter(w Writer)
func (*MultiWriter) CanFeedback ¶
func (mw *MultiWriter) CanFeedback(sourceName, metricName string) bool
CanFeedback reports whether at least one contained writer can answer for the pair. Writers that do not implement Feedbacker are simply skipped.
func (*MultiWriter) Count ¶
func (mw *MultiWriter) Count() int
func (*MultiWriter) DefineMetrics ¶
func (mw *MultiWriter) DefineMetrics(metrics *metrics.Metrics) (err error)
func (*MultiWriter) LastMeasurement ¶
func (mw *MultiWriter) LastMeasurement(ctx context.Context, sourceName, metricName string) (int64, error)
LastMeasurement returns the oldest epoch reported by any feedback-capable writer. The minimum is deliberate: a consumer resuming from the maximum would leave the laggard sink permanently short, whereas resuming from the minimum only duplicates measurements the leading sink already has. A writer holding nothing at all is the extreme case of lagging, so it short-circuits the whole aggregate to ErrNoFeedbackData.
func (*MultiWriter) Migrate ¶
func (mw *MultiWriter) Migrate() (err error)
Migrate runs migrations on all writers that support it
func (*MultiWriter) NeedsMigration ¶
func (mw *MultiWriter) NeedsMigration() (bool, error)
NeedsMigration checks if any writer needs migration
func (*MultiWriter) SyncMetric ¶
func (mw *MultiWriter) SyncMetric(sourceName, metricName string, op SyncOp) (err error)
func (*MultiWriter) Write ¶
func (mw *MultiWriter) Write(msg metrics.MeasurementEnvelope) (err error)
type PostgresWriter ¶
type PostgresWriter struct {
// contains filtered or unexported fields
}
PostgresWriter is a sink that writes metric measurements to a Postgres database. At the moment, it supports both Postgres and TimescaleDB as a storage backend. However, one is able to use any Postgres-compatible database as a storage backend, e.g. PGEE, Citus, Greenplum, CockroachDB, etc.
func NewPostgresWriter ¶
func NewWriterFromPostgresConn ¶
func NewWriterFromPostgresConn(ctx context.Context, conn db.PgxPoolIface, opts *CmdOpts) (pgw *PostgresWriter, err error)
func (*PostgresWriter) AddDBUniqueMetricToListingTable ¶
func (pgw *PostgresWriter) AddDBUniqueMetricToListingTable(dbUnique, metric string) error
func (*PostgresWriter) CanFeedback ¶
func (pgw *PostgresWriter) CanFeedback(sourceName, metricName string) bool
CanFeedback reports whether LastMeasurement can be attempted for the pair. It is optimistic and lock-free by design: partitionMapMetric is populated only by the flush path, so it is empty at process start — exactly when a resuming collector asks. Whether the metric table actually exists is resolved authoritatively by LastMeasurement.
func (*PostgresWriter) DeleteOldPartitions ¶
func (pgw *PostgresWriter) DeleteOldPartitions()
DeleteOldPartitions is a background task that deletes old partitions from the measurements DB
func (*PostgresWriter) EnsureBuiltinMetricDummies ¶
func (pgw *PostgresWriter) EnsureBuiltinMetricDummies() (err error)
EnsureBuiltinMetricDummies creates empty tables for all built-in metrics if they don't exist
func (*PostgresWriter) EnsureMetricDummy ¶
func (pgw *PostgresWriter) EnsureMetricDummy(metric string) (err error)
EnsureMetricDummy creates an empty table for a metric measurements if it doesn't exist
func (*PostgresWriter) EnsureMetricTimePartsExist ¶
func (pgw *PostgresWriter) EnsureMetricTimePartsExist(metricPartBounds map[string]ExistingPartitionInfo) error
func (*PostgresWriter) EnsureMetricTimescale ¶
func (pgw *PostgresWriter) EnsureMetricTimescale(pgPartBounds map[string]ExistingPartitionInfo) (err error)
func (*PostgresWriter) LastMeasurement ¶
func (pgw *PostgresWriter) LastMeasurement(ctx context.Context, sourceName, metricName string) (int64, error)
LastMeasurement returns the epoch_ns of the newest measurement stored for the pair. The scan is bounded by the configured retention interval so it cannot degrade into a full scan of every partition.
func (*PostgresWriter) MaintainUniqueSources ¶
func (pgw *PostgresWriter) MaintainUniqueSources()
MaintainUniqueSources is a background task that maintains a mapping of unique sources in each metric table in admin.all_distinct_dbname_metrics. This is used to avoid listing the same source multiple times in Grafana dropdowns.
func (*PostgresWriter) Migrate ¶
func (pgw *PostgresWriter) Migrate() error
Migrate upgrades database with all migrations
func (*PostgresWriter) NeedsMigration ¶
func (pgw *PostgresWriter) NeedsMigration() (bool, error)
NeedsMigration checks if database needs migration
func (*PostgresWriter) ReadMetricSchemaType ¶
func (pgw *PostgresWriter) ReadMetricSchemaType() (err error)
func (*PostgresWriter) SyncMetric ¶
func (pgw *PostgresWriter) SyncMetric(sourceName, metricName string, op SyncOp) error
SyncMetric ensures that tables exist for newly added metrics and/or sources
func (*PostgresWriter) Write ¶
func (pgw *PostgresWriter) Write(msg metrics.MeasurementEnvelope) error
Write sends the measurements to the cache channel
type PromMetricCache ¶
type PromMetricCache = map[string]map[string]metrics.MeasurementEnvelope // [dbUnique][metric]lastly_fetched_data
type PrometheusWriter ¶
type PrometheusWriter struct {
sync.RWMutex
Namespace string
Cache PromMetricCache // [dbUnique][metric]lastly_fetched_data
// contains filtered or unexported fields
}
PrometheusWriter is a Prometheus exporter that implements the prometheus.Collector interface using the "unchecked collector" pattern (empty Describe method).
Design decisions based on Prometheus exporter guidelines (https://prometheus.io/docs/instrumenting/writing_exporters/#collectors):
Metrics are collected periodically by reaper and cached in-memory. On scrape, the collector reads a snapshot of the cache and emits fresh NewConstMetric values. The cache is NOT consumed on scrape, so parallel or back-to-back scrapes see the same data until the next Write() updates arrive.
This is an "unchecked collector": Describe() sends no descriptors, which tells the Prometheus registry to skip consistency checks. This is necessary because the set of metrics is dynamic (driven by monitored databases and their query results). Safety is ensured by deduplicating metric identities within each Collect() call.
Label keys are always sorted lexicographically before building descriptors and label value slices. This guarantees deterministic descriptor identity regardless of Go map iteration order.
func NewPrometheusWriter ¶
func NewPrometheusWriter(ctx context.Context, connstr string) (promw *PrometheusWriter, err error)
func (*PrometheusWriter) AddCacheEntry ¶
func (promw *PrometheusWriter) AddCacheEntry(dbUnique, metric string, msgArr metrics.MeasurementEnvelope)
func (*PrometheusWriter) Collect ¶
func (promw *PrometheusWriter) Collect(ch chan<- prometheus.Metric)
Collect implements prometheus.Collector. It reads a snapshot of the metric cache and emits const metrics. Parallel scrapes see the same data until background Write() calls update it
func (*PrometheusWriter) DefineMetrics ¶
func (promw *PrometheusWriter) DefineMetrics(metrics *metrics.Metrics) (err error)
DefineMetrics is called by reaper on startup and whenever metric definitions change
func (*PrometheusWriter) Describe ¶
func (promw *PrometheusWriter) Describe(_ chan<- *prometheus.Desc)
Describe is intentionally empty to make PrometheusWriter an "unchecked collector" per the prometheus.Collector contract
func (*PrometheusWriter) InitCacheEntry ¶
func (promw *PrometheusWriter) InitCacheEntry(dbUnique string)
func (*PrometheusWriter) Println ¶
func (promw *PrometheusWriter) Println(v ...any)
Println implements promhttp.Logger
func (*PrometheusWriter) PurgeCacheEntry ¶
func (promw *PrometheusWriter) PurgeCacheEntry(dbUnique, metric string)
func (*PrometheusWriter) SyncMetric ¶
func (promw *PrometheusWriter) SyncMetric(sourceName, metricName string, op SyncOp) error
SyncMetric is called by reaper when a metric or monitored source is removed or added, allowing the writer to purge or initialize cache entries as needed
func (*PrometheusWriter) Write ¶
func (promw *PrometheusWriter) Write(msg metrics.MeasurementEnvelope) error
Write is called by reaper whenever new measurement data arrives
func (*PrometheusWriter) WritePromMetrics ¶
func (promw *PrometheusWriter) WritePromMetrics(msg metrics.MeasurementEnvelope, ch chan<- prometheus.Metric) (written int, errorCount int)
WritePromMetrics converts a MeasurementEnvelope into Prometheus const metrics and sends them directly to ch. Returns the count of metrics written and errors encountered.
For prom-sourced envelopes (SourceKind == "prometheus") the original family name is used as the fully-qualified metric name with no pgwatch namespace prefix, and each measurement's own epoch_ns is used as the timestamp. For all other sources the existing namespace+family+field naming and a single envelope-level timestamp are applied.
type RPCWriter ¶
type RPCWriter struct {
// contains filtered or unexported fields
}
RPCWriter sends metric measurements to a remote server using gRPC. Remote servers should make use of the .proto file under api/pb/ to integrate with it. It's up to the implementer to define the behavior of the server. It can be a simple logger, external storage, alerting system, or an analytics system.
func NewRPCWriter ¶
func (*RPCWriter) CanFeedback ¶
CanFeedback reports whether a feedback query is worth attempting. It is optimistic until the remote server proves it does not implement the method.
func (*RPCWriter) DefineMetrics ¶
DefineMetrics sends metric definitions to the remote server
func (*RPCWriter) LastMeasurement ¶
func (rw *RPCWriter) LastMeasurement(ctx context.Context, sourceName, metricName string) (int64, error)
LastMeasurement asks the remote server for the newest measurement it holds for the pair. The caller's context drives cancellation; the credential metadata is carried over from the writer's own context.
func (*RPCWriter) SyncMetric ¶
SyncMetric synchronizes a metric and monitored source with the remote server
type SyncOp ¶
type SyncOp int32
SyncOp represents synchronization operations for metrics. These constants are used both in Go code and protobuf definitions.