store

package
v1.0.0 Latest Latest
Warning

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

Go to latest
Published: Oct 7, 2026 License: Apache-2.0 Imports: 15 Imported by: 0

Documentation

Overview

Package store is Stampede's PostgreSQL (and TimescaleDB) persistence: connection pool, embedded migrations and typed queries generated by sqlc from internal/store/queries.

Index

Constants

View Source
const (
	Rollup10s = "10s"
	Rollup1m  = "1m"
)

Rollup resolutions of the per-second run metrics (migration 00011).

View Source
const LeaderLockKey = 7462

LeaderLockKey is the advisory lock held by the active server replica.

View Source
const MinMetricsRetention = 24 * time.Hour

MinMetricsRetention is the shortest retention accepted: the continuous aggregates refresh the last day, so per-second rows must outlive it.

Variables

This section is empty.

Functions

func IsNotFound

func IsNotFound(err error) bool

IsNotFound reports whether err means no rows matched.

func IsUniqueViolation

func IsUniqueViolation(err error) bool

IsUniqueViolation reports whether err is a unique constraint violation.

func RollupWidth

func RollupWidth(res string) time.Duration

RollupWidth is a resolution's bucket width.

Types

type Lock

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

Lock is a session-level advisory lock held on a dedicated connection.

func (*Lock) Alive

func (l *Lock) Alive(ctx context.Context) error

Alive reports whether the connection holding the lock still works. If it does not, the database has dropped the session and the lock with it.

func (*Lock) Release

func (l *Lock) Release()

Release frees the lock.

type RollupPoint

type RollupPoint struct {
	Bucket                       time.Time
	Requests, Failed             int64
	RPS, P50, P95, P99           float64
	VUs                          int
	Planned                      float64
	Dropped, Iterations, Samples int64
	SchedLag                     float64
}

RollupPoint is one bucket of a run's rolled-up metrics. P50 is the mean of the per-second medians; P95 and P99 are the worst per-second values.

type Store

type Store struct {
	Pool *pgxpool.Pool
	*db.Queries
}

Store wraps a connection pool and the generated queries.

func Open

func Open(ctx context.Context, url string) (*Store, error)

Open connects to PostgreSQL and checks the connection.

func (*Store) AdvisoryLock

func (s *Store) AdvisoryLock(ctx context.Context, key int64) (*Lock, error)

AdvisoryLock tries to take a session-level advisory lock on a dedicated connection, used for leader election between server replicas. It returns nil when another session holds the lock.

func (*Store) Close

func (s *Store) Close()

Close releases the pool.

func (*Store) ContinuousRollups

func (s *Store) ContinuousRollups(ctx context.Context) (bool, error)

ContinuousRollups reports whether the rollups are TimescaleDB continuous aggregates (kept after per-second rows expire) rather than plain views.

func (*Store) DeleteMetricsBefore

func (s *Store) DeleteMetricsBefore(ctx context.Context, t time.Time) (int64, error)

DeleteMetricsBefore deletes per-second run metrics older than t, in batches, and returns how many rows went.

func (*Store) InTx

func (s *Store) InTx(ctx context.Context, fn func(q *db.Queries) error) error

InTx runs fn in a transaction, committing when it returns nil.

func (*Store) InTxLocked

func (s *Store) InTxLocked(ctx context.Context, key int64, fn func(q *db.Queries) error) error

InTxLocked is InTx holding a transaction-scoped advisory lock on key, so concurrent callers with the same key run one at a time.

func (*Store) Listen

func (s *Store) Listen(ctx context.Context, channel string, fn func(payload string), onReconnect func())

Listen calls fn with every notification on channel until ctx ends. It holds a dedicated connection and reconnects after errors, calling onReconnect (if set) each time it listens again, since notifications sent while it was disconnected are lost.

func (*Store) Migrate

func (s *Store) Migrate(ctx context.Context, log *slog.Logger, dryRun bool) error

Migrate applies pending migrations. With dryRun it only reports them.

func (*Store) MigrationVersion

func (s *Store) MigrationVersion(ctx context.Context) (int64, error)

MigrationVersion reports the applied schema version.

func (*Store) Notify

func (s *Store) Notify(ctx context.Context, channel, payload string) error

Notify sends payload to everyone listening on channel (pg_notify).

func (*Store) RefreshRollups

func (s *Store) RefreshRollups(ctx context.Context, from, to time.Time) error

RefreshRollups brings continuous aggregates up to date for a time range (the refresh policy does this every minute; tests and backfills call it directly). It does nothing on plain PostgreSQL.

func (*Store) RunRollup

func (s *Store) RunRollup(ctx context.Context, run uuid.UUID, res string) ([]RollupPoint, error)

RunRollup returns a run's metrics in buckets of res (Rollup10s or Rollup1m), oldest first.

func (*Store) SetMetricsRetention

func (s *Store) SetMetricsRetention(ctx context.Context, d time.Duration) (managed bool, err error)

SetMetricsRetention keeps per-second run metrics for d (0 keeps them forever). With TimescaleDB a retention policy drops whole chunks and managed is true; on plain PostgreSQL the caller deletes old rows with DeleteMetricsBefore, periodically.

Directories

Path Synopsis
Package storetest starts a disposable TimescaleDB for tests.
Package storetest starts a disposable TimescaleDB for tests.

Jump to

Keyboard shortcuts

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