fix(store): 数据库连接池限 25 条并设语句与空闲事务超时,迁移锁固定在取锁的连接上释放,API 导出连接池指标
This commit is contained in:
7 files changed
+348
-12
No files matched your search
@@ -276,6 +276,9 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int {
|
||||
// the same store the auth handlers write to, so a login and the next request
|
||||
// agree on what local auth knows.
|
||||
repo := api.NewPGRepo(drv.DB())
|
||||
if err := api.RegisterStorePool(drv.DB()); err != nil {
|
||||
fmt.Fprintf(stderr, "felis api: store pool metrics unavailable: %v\n", err)
|
||||
}
|
||||
|
||||
// The owner's on-demand backup levers come from [archive], the same keys the
|
||||
// backup Job and the reaper read. A malformed key leaves the defaults in
|
||||
|
||||
@@ -1317,6 +1317,14 @@ All four mandated metrics have real producers; scrape them when triaging:
|
||||
- The Go runtime and process series (`go_*`, `process_*`) of `felis-api`:
|
||||
goroutines, heap, open file descriptors. Goroutines that climb without
|
||||
falling back usually mean streams or uploads that never end.
|
||||
- `go_sql_*{db_name="felis"}` — the API's database pool. It is capped at 25
|
||||
connections, and every connection runs with `statement_timeout=15s` and
|
||||
`idle_in_transaction_session_timeout=60s` (a `[database] url` that sets
|
||||
either keeps its own). `go_sql_in_use_connections` sitting at
|
||||
`go_sql_max_open_connections` with `go_sql_wait_count_total` climbing means
|
||||
requests are queueing for a connection: look for a slow query or a lock
|
||||
(`SELECT pid, state, wait_event, query FROM pg_stat_activity`). A statement
|
||||
cut off by the limit logs `canceling statement due to statement timeout`.
|
||||
|
||||
The API also writes one access-log line per request to its log, in logfmt:
|
||||
`face`, `method`, `route`, `path`, `status`, `duration_ms`, `bytes`,
|
||||
@@ -1340,7 +1348,7 @@ annotated Service endpoints picks them up as is.
|
||||
`controller_runtime_reconcile_*` / `workqueue_*` series.
|
||||
- `felis-api` internal face `:8081/metrics` (Service `felis-api-internal`) —
|
||||
`felis_build_info{component="api"}`, the `felis_http_*` request series, the
|
||||
`go_*`/`process_*` runtime series,
|
||||
`go_*`/`process_*` runtime series, the `go_sql_*` pool series,
|
||||
`felis_image_build_failures_total`, and the sign-in series of §17
|
||||
(`felis_mail_total`, `felis_rate_limited_total`,
|
||||
`felis_auth_otp_lockouts_total`, `felis_auth_failures_total`,
|
||||
|
||||
+13
-3
@@ -1,6 +1,7 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"net/http"
|
||||
|
||||
"felis.lolicon.best/internal/metrics"
|
||||
@@ -15,10 +16,12 @@ import (
|
||||
// build reconcile loop lives in cmd/felis.reconcileBuilds), so this endpoint is
|
||||
// that counter's sole scrape path; the operator's :8080 carries the fleet
|
||||
// gauges instead.
|
||||
var apiMetricsHandler = newAPIMetricsHandler()
|
||||
var (
|
||||
apiRegistry = prometheus.NewRegistry()
|
||||
apiMetricsHandler = newAPIMetricsHandler(apiRegistry)
|
||||
)
|
||||
|
||||
func newAPIMetricsHandler() http.Handler {
|
||||
reg := prometheus.NewRegistry()
|
||||
func newAPIMetricsHandler(reg *prometheus.Registry) http.Handler {
|
||||
// A fresh registry cannot already hold a collector, so Register's
|
||||
// AlreadyRegistered tolerance arm never triggers here.
|
||||
_ = metrics.Register(reg)
|
||||
@@ -31,6 +34,13 @@ func newAPIMetricsHandler() http.Handler {
|
||||
return promhttp.HandlerFor(reg, promhttp.HandlerOpts{})
|
||||
}
|
||||
|
||||
// RegisterStorePool adds the store pool's go_sql_* series (db_name="felis"):
|
||||
// connections open, in use and idle against the cap, and how often and how
|
||||
// long requests waited for one, which is where a pool pinned at its cap shows.
|
||||
func RegisterStorePool(db *sql.DB) error {
|
||||
return apiRegistry.Register(collectors.NewDBStatsCollector(db, "felis"))
|
||||
}
|
||||
|
||||
// handleMetrics is mounted as a Public route on the internal face (the same
|
||||
// stance as the health probes): the internal listener is ClusterIP-only and a
|
||||
// Prometheus scrape carries no token. The external face never serves metrics.
|
||||
|
||||
@@ -1,11 +1,14 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"net/http"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"felis.lolicon.best/internal/metrics"
|
||||
|
||||
_ "github.com/jackc/pgx/v5/stdlib"
|
||||
)
|
||||
|
||||
// TestMetricsEndpoint locks the scrape surface: the internal face serves the
|
||||
@@ -42,3 +45,26 @@ func TestMetricsEndpoint(t *testing.T) {
|
||||
t.Fatalf("external GET /metrics = %d, want 404 (metrics stay internal)", w.Code)
|
||||
}
|
||||
}
|
||||
|
||||
// The store pool's counts reach the internal scrape once registered.
|
||||
func TestMetricsEndpointCarriesTheStorePool(t *testing.T) {
|
||||
pool, err := sql.Open("pgx", "postgres://[email protected]:1/felis")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer pool.Close()
|
||||
pool.SetMaxOpenConns(7)
|
||||
if err := RegisterStorePool(pool); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
w := do((&API{}).InternalHandler(), "GET", "/metrics", "", nil)
|
||||
for _, want := range []string{
|
||||
`go_sql_max_open_connections{db_name="felis"} 7`,
|
||||
`go_sql_wait_count_total{db_name="felis"} 0`,
|
||||
} {
|
||||
if !strings.Contains(w.Body.String(), want) {
|
||||
t.Errorf("metrics exposition missing %q", want)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,135 @@
|
||||
//go:build pgint
|
||||
|
||||
package pgint
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"felis.lolicon.best/internal/store"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgconn"
|
||||
)
|
||||
|
||||
// Every pooled connection starts with the server-side limits, so a runaway
|
||||
// statement or an abandoned transaction cannot hold a connection forever.
|
||||
func TestPoolSessionsCarryTheTimeouts(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
for setting, want := range map[string]string{
|
||||
"statement_timeout": "15s",
|
||||
"idle_in_transaction_session_timeout": "1min",
|
||||
} {
|
||||
var got string
|
||||
if err := db.QueryRowContext(ctx, "SHOW "+setting).Scan(&got); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got != want {
|
||||
t.Errorf("%s = %q, want %q", setting, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The limit is enforced by the server: a DSN that tightens it gets its
|
||||
// statements cancelled at that bound.
|
||||
func TestStatementTimeoutCancelsARunawayStatement(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
dsn := os.Getenv("FELIS_TEST_PG_URL")
|
||||
sep := "?"
|
||||
if strings.Contains(dsn, "?") {
|
||||
sep = "&"
|
||||
}
|
||||
drv, err := store.Open(ctx, dsn+sep+"statement_timeout=200")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer drv.Close()
|
||||
_, err = drv.DB().ExecContext(ctx, "SELECT pg_sleep(2)")
|
||||
var pgErr *pgconn.PgError
|
||||
if !errors.As(err, &pgErr) || pgErr.Code != "57014" { // query_canceled
|
||||
t.Fatalf("pg_sleep(2) under a 200ms limit: err = %v, want query_canceled", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Unlock releases the lock on the connection that took it, even when the pool
|
||||
// would hand that connection to someone else: another migrator can take it next.
|
||||
func TestMigrationLockIsReleasedOnItsOwnConnection(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
dsn := os.Getenv("FELIS_TEST_PG_URL")
|
||||
drv, err := store.Open(ctx, dsn)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer drv.Close()
|
||||
other, err := store.Open(ctx, dsn)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer other.Close()
|
||||
|
||||
if err := drv.Lock(ctx); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Hold the pool's idle connection, so a release that went through the pool
|
||||
// would land on a different session from the lock's.
|
||||
busy, err := drv.DB().Conn(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer busy.Close()
|
||||
|
||||
var held bool
|
||||
if err := other.DB().QueryRowContext(ctx, "SELECT pg_try_advisory_lock($1)", store.AdvisoryLockKey).Scan(&held); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if held {
|
||||
t.Fatal("a second migrator took the lock while the first held it")
|
||||
}
|
||||
|
||||
if err := drv.Unlock(ctx); err != nil {
|
||||
t.Fatalf("Unlock: %v", err)
|
||||
}
|
||||
conn, err := other.DB().Conn(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer conn.Close()
|
||||
if err := conn.QueryRowContext(ctx, "SELECT pg_try_advisory_lock($1)", store.AdvisoryLockKey).Scan(&held); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !held {
|
||||
t.Fatal("the lock is still held after Unlock")
|
||||
}
|
||||
if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", store.AdvisoryLockKey); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
// A migration body runs without the request-sized statement limit.
|
||||
func TestMigrationBodyHasNoStatementLimit(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
drv, err := store.Open(ctx, os.Getenv("FELIS_TEST_PG_URL"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer drv.Close()
|
||||
t.Cleanup(func() {
|
||||
_, _ = db.ExecContext(ctx, "DROP TABLE IF EXISTS timeout_probe")
|
||||
_, _ = db.ExecContext(ctx, "DELETE FROM schema_migrations WHERE version = 999999")
|
||||
})
|
||||
|
||||
probe := store.Migration{Version: 999999, Name: "timeout_probe",
|
||||
SQL: "CREATE TABLE timeout_probe AS SELECT current_setting('statement_timeout') AS v"}
|
||||
if err := drv.Apply(ctx, probe); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var v string
|
||||
if err := db.QueryRowContext(ctx, "SELECT v FROM timeout_probe").Scan(&v); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if v != "0" {
|
||||
t.Fatalf("statement_timeout inside a migration = %q, want 0", v)
|
||||
}
|
||||
}
|
||||
+110
-8
@@ -3,22 +3,83 @@ package store
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"database/sql/driver"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgconn"
|
||||
_ "github.com/jackc/pgx/v5/stdlib" // register the "pgx" database/sql driver
|
||||
"github.com/jackc/pgx/v5/stdlib"
|
||||
)
|
||||
|
||||
// PostgresDriver is the production Driver, backed by a database/sql pool using
|
||||
// the pgx stdlib driver.
|
||||
type PostgresDriver struct {
|
||||
db *sql.DB
|
||||
// lock is the one connection holding the migration advisory lock, from Lock
|
||||
// to Unlock. The lock belongs to a server session, so it must be released on
|
||||
// the connection that took it, never on whichever one the pool hands out.
|
||||
lock *sql.Conn
|
||||
}
|
||||
|
||||
// Pool bounds. Every API request that reads the store takes a connection, and
|
||||
// PostgreSQL's default max_connections of 100 is shared with the reaper and
|
||||
// backup Jobs, the offsite sync and a break-glass console. Uncapped, a flood of
|
||||
// public requests or a pile-up behind one lock would open connections until
|
||||
// those could not get one; capped, the excess waits in database/sql's queue
|
||||
// until its request context gives up.
|
||||
const (
|
||||
MaxOpenConns = 25
|
||||
maxIdleConns = 10
|
||||
connMaxLifetime = 30 * time.Minute
|
||||
connMaxIdleTime = 5 * time.Minute
|
||||
)
|
||||
|
||||
// SessionDefaults are the server-side limits every connection starts with. A
|
||||
// statement that runs away or waits on a lock is cancelled instead of holding
|
||||
// its connection forever, and a transaction a stuck caller left open is rolled
|
||||
// back with its locks. Nothing felis runs per request comes near these; a
|
||||
// migration lifts the statement limit for its own work (Lock, Apply). A DSN that
|
||||
// sets either one, as a query parameter or in options=-c, keeps its own value.
|
||||
var SessionDefaults = map[string]string{
|
||||
"statement_timeout": "15s",
|
||||
"idle_in_transaction_session_timeout": "60s",
|
||||
}
|
||||
|
||||
// connConfig parses dsn and fills in the SessionDefaults it does not set.
|
||||
func connConfig(dsn string) (*pgx.ConnConfig, error) {
|
||||
cfg, err := pgx.ParseConfig(dsn)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for k, v := range SessionDefaults {
|
||||
if _, set := cfg.RuntimeParams[k]; set || strings.Contains(cfg.RuntimeParams["options"], k) {
|
||||
continue
|
||||
}
|
||||
cfg.RuntimeParams[k] = v
|
||||
}
|
||||
return cfg, nil
|
||||
}
|
||||
|
||||
// newPool builds the bounded pool without dialing.
|
||||
func newPool(dsn string) (*sql.DB, error) {
|
||||
cfg, err := connConfig(dsn)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
db := stdlib.OpenDB(*cfg)
|
||||
db.SetMaxOpenConns(MaxOpenConns)
|
||||
db.SetMaxIdleConns(maxIdleConns)
|
||||
db.SetConnMaxLifetime(connMaxLifetime)
|
||||
db.SetConnMaxIdleTime(connMaxIdleTime)
|
||||
return db, nil
|
||||
}
|
||||
|
||||
// Open dials dsn and returns a PostgresDriver. The caller owns Close.
|
||||
func Open(ctx context.Context, dsn string) (*PostgresDriver, error) {
|
||||
db, err := sql.Open("pgx", dsn)
|
||||
db, err := newPool(dsn)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("open postgres: %w", err)
|
||||
}
|
||||
@@ -35,16 +96,53 @@ func (d *PostgresDriver) DB() *sql.DB { return d.db }
|
||||
// Close releases the pool.
|
||||
func (d *PostgresDriver) Close() error { return d.db.Close() }
|
||||
|
||||
// Lock takes the session-level advisory lock that serializes migrations.
|
||||
// Lock takes the session-level advisory lock that serializes migrations, on a
|
||||
// connection it keeps out of the pool until Unlock. Waiting behind another
|
||||
// migrator is not a runaway statement, so that session has no statement limit.
|
||||
func (d *PostgresDriver) Lock(ctx context.Context) error {
|
||||
_, err := d.db.ExecContext(ctx, "SELECT pg_advisory_lock($1)", AdvisoryLockKey)
|
||||
return err
|
||||
if d.lock != nil {
|
||||
return errors.New("migration lock already held")
|
||||
}
|
||||
conn, err := d.db.Conn(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := conn.ExecContext(ctx, "SET statement_timeout = 0"); err != nil {
|
||||
discard(conn)
|
||||
return err
|
||||
}
|
||||
if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_lock($1)", AdvisoryLockKey); err != nil {
|
||||
discard(conn)
|
||||
return err
|
||||
}
|
||||
d.lock = conn
|
||||
return nil
|
||||
}
|
||||
|
||||
// Unlock releases the advisory lock.
|
||||
// Unlock releases the advisory lock on the connection that took it, then
|
||||
// closes that connection rather than returning it to the pool, so neither a
|
||||
// lock that failed to release nor the lifted statement limit outlives the run.
|
||||
func (d *PostgresDriver) Unlock(ctx context.Context) error {
|
||||
_, err := d.db.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", AdvisoryLockKey)
|
||||
return err
|
||||
conn := d.lock
|
||||
if conn == nil {
|
||||
return errors.New("migration lock not held")
|
||||
}
|
||||
d.lock = nil
|
||||
defer discard(conn)
|
||||
var released bool
|
||||
if err := conn.QueryRowContext(ctx, "SELECT pg_advisory_unlock($1)", AdvisoryLockKey).Scan(&released); err != nil {
|
||||
return err
|
||||
}
|
||||
if !released {
|
||||
return errors.New("migration lock was not held by its connection")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// discard closes conn's server session instead of pooling it.
|
||||
func discard(conn *sql.Conn) {
|
||||
_ = conn.Raw(func(any) error { return driver.ErrBadConn })
|
||||
_ = conn.Close()
|
||||
}
|
||||
|
||||
// EnsureVersionTable creates the bookkeeping table if absent.
|
||||
@@ -90,6 +188,10 @@ func (d *PostgresDriver) Apply(ctx context.Context, m Migration) error {
|
||||
}
|
||||
defer tx.Rollback() //nolint:errcheck // rollback after a successful commit is a no-op
|
||||
|
||||
// A schema change on a grown table may rightly take longer than any request.
|
||||
if _, err := tx.ExecContext(ctx, "SET LOCAL statement_timeout = 0"); err != nil {
|
||||
return fmt.Errorf("lift statement timeout: %w", err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, m.SQL); err != nil {
|
||||
return fmt.Errorf("exec body: %w", err)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,52 @@
|
||||
package store
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestConnConfigFillsSessionDefaults(t *testing.T) {
|
||||
cfg, err := connConfig("postgres://felis:[email protected]:5432/felis?sslmode=disable")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for k, want := range map[string]string{
|
||||
"statement_timeout": "15s",
|
||||
"idle_in_transaction_session_timeout": "60s",
|
||||
} {
|
||||
if got := cfg.RuntimeParams[k]; got != want {
|
||||
t.Errorf("%s = %q, want %q", k, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestConnConfigKeepsTheDSNsOwnLimits(t *testing.T) {
|
||||
cfg, err := connConfig("postgres://felis:[email protected]:5432/felis?statement_timeout=2500")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := cfg.RuntimeParams["statement_timeout"]; got != "2500" {
|
||||
t.Errorf("statement_timeout = %q, want the DSN's 2500", got)
|
||||
}
|
||||
if got := cfg.RuntimeParams["idle_in_transaction_session_timeout"]; got != "60s" {
|
||||
t.Errorf("idle_in_transaction_session_timeout = %q, want the default 60s", got)
|
||||
}
|
||||
|
||||
// Set through options=-c, it is left out: the server applies startup
|
||||
// parameters after options, so a default sent beside it would win.
|
||||
cfg, err = connConfig("postgres://felis:[email protected]:5432/felis?options=-c%20statement_timeout%3D0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if v, set := cfg.RuntimeParams["statement_timeout"]; set {
|
||||
t.Errorf("statement_timeout = %q beside options=-c, want it unset", v)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPoolIsBounded(t *testing.T) {
|
||||
db, err := newPool("postgres://felis:[email protected]:1/felis")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
if got := db.Stats().MaxOpenConnections; got != MaxOpenConns || got == 0 {
|
||||
t.Fatalf("MaxOpenConnections = %d, want %d", got, MaxOpenConns)
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user