diff --git a/cmd/felis/api.go b/cmd/felis/api.go index 1c14ecb..ca7c3df 100644 --- a/cmd/felis/api.go +++ b/cmd/felis/api.go @@ -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 diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index 220e64e..cffc24b 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -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`, diff --git a/internal/api/metrics.go b/internal/api/metrics.go index b38562a..cac9849 100644 --- a/internal/api/metrics.go +++ b/internal/api/metrics.go @@ -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. diff --git a/internal/api/metrics_test.go b/internal/api/metrics_test.go index 534e4e2..6e33491 100644 --- a/internal/api/metrics_test.go +++ b/internal/api/metrics_test.go @@ -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://felis@127.0.0.1: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) + } + } +} diff --git a/internal/pgint/store_test.go b/internal/pgint/store_test.go new file mode 100644 index 0000000..27fd3b6 --- /dev/null +++ b/internal/pgint/store_test.go @@ -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) + } +} diff --git a/internal/store/sqldriver.go b/internal/store/sqldriver.go index 0e6ae98..d2c87ef 100644 --- a/internal/store/sqldriver.go +++ b/internal/store/sqldriver.go @@ -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) } diff --git a/internal/store/sqldriver_test.go b/internal/store/sqldriver_test.go new file mode 100644 index 0000000..de23c5e --- /dev/null +++ b/internal/store/sqldriver_test.go @@ -0,0 +1,52 @@ +package store + +import "testing" + +func TestConnConfigFillsSessionDefaults(t *testing.T) { + cfg, err := connConfig("postgres://felis:pw@127.0.0.1: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:pw@127.0.0.1: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:pw@127.0.0.1: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:pw@127.0.0.1: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) + } +}