diff --git a/.github/workflows/bdd-tests.yml b/.github/workflows/bdd-tests.yml index f2298bc36..b650d646a 100644 --- a/.github/workflows/bdd-tests.yml +++ b/.github/workflows/bdd-tests.yml @@ -219,12 +219,13 @@ jobs: packages: read strategy: fail-fast: false - # GitHub-hosted runners share CPU under heavy matrix fan-out, and - # the timing-sensitive scenarios (sleep-based retain/lifetime - # waits, SCRAM passthrough reconnect) lose their margin when 20+ - # BDD jobs run in parallel. Cap concurrency so each suite gets a - # less contested runner. - max-parallel: 4 + # Six concurrent suites is the cap this matrix was tuned for. + # GitHub-hosted runners share CPU under heavy fan-out; with 20+ + # BDD jobs in parallel, the sleep-based lifecycle tests and the + # SCRAM passthrough reconnect tests lost timing margin. This cap + # keeps those suites stable while reducing total wall time versus + # the previous limit of four. + max-parallel: 6 matrix: suite: - { name: "Go", cargo: "test --test bdd -- --tags @go" } @@ -249,6 +250,7 @@ jobs: - { name: "Server TLS", cargo: "test --test bdd -- --tags @server-tls" } - { name: "TLS migration (vendored OpenSSL)", cargo: "test --features tls-migration --test bdd -- --tags @tls-migration" } - { name: "Startup parameters", cargo: "test --test bdd -- --tags @startup-parameters" } + - { name: "Web UI", cargo: "test --test bdd -- --tags @web-ui" } steps: - name: Checkout repository uses: actions/checkout@v4 diff --git a/.github/workflows/dashboard-validation.yml b/.github/workflows/dashboard-validation.yml index 0d95ad127..c7e90dadd 100644 --- a/.github/workflows/dashboard-validation.yml +++ b/.github/workflows/dashboard-validation.yml @@ -5,6 +5,7 @@ on: branches: [master] paths: - "grafana/**" + - "monitoring/prometheus-rules/**" - "scripts/dashboard-*" - "scripts/docker-smoke.sh" - "src/web/metrics/**" @@ -17,6 +18,7 @@ on: pull_request: paths: - "grafana/**" + - "monitoring/prometheus-rules/**" - "scripts/dashboard-*" - "scripts/docker-smoke.sh" - "src/web/metrics/**" diff --git a/Cargo.lock b/Cargo.lock index d8b726c3f..31ebfaa47 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2118,6 +2118,8 @@ dependencies = [ "bytes", "fallible-iterator", "postgres-protocol", + "serde", + "serde_json", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 73bdb6c75..efe811b98 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -68,7 +68,10 @@ iota = { version = "0.2.3" } pin-project-lite = "0.2.16" pam-client = { version = "0.5.0", optional = true } postgres = "0.19.10" -tokio-postgres = "0.7" +# `with-serde_json-1` lets `Row::try_get` decode `json`/`jsonb` +# columns directly into `serde_json::Value`, which allows auth_query +# startup_parameters to come from a native JSON column. +tokio-postgres = { version = "0.7", features = ["with-serde_json-1"] } postgres-native-tls = "0.5.1" flate2 = "1.0.28" sd-notify = "0.4" diff --git a/documentation/en/src/changelog.md b/documentation/en/src/changelog.md index 89eaa50e6..fc0711631 100644 --- a/documentation/en/src/changelog.md +++ b/documentation/en/src/changelog.md @@ -2,7 +2,18 @@ ### 3.9.1 -Web admin console refresh. +Web admin console refresh and a follow-up pass on `startup_parameters`. + +Upgrade notes for operators monitoring 3.9.0: + +- The pg_doorman-side budget rejection now returns `SQLSTATE 53400` + (`configuration_limit_exceeded`) instead of `54000`. Alert rules + and log filters keyed on `54000` need to switch. +- `PgDoormanStartupParameterPgRejection` is now `severity: warning` + (was `critical` in 3.9.0). Cascade-overflow stays `critical`. Review + the Alertmanager / on-call routing if you key on severity to page. + +#### Web admin console - Light theme by default. Three-position theme toggle (Light / System / Dark) in the sidebar footer; choice persists in localStorage. @@ -39,6 +50,29 @@ Backend: `web/access_log.rs` demotes authenticated 2xx reads to debug. Docs: `guides/web-ui.md` rewritten for the new pages and shortcuts. +#### startup_parameters follow-up + +- If the resolved `startup_parameters` set exceeds the startup packet + budget, backend startup now fails with `SQLSTATE 53400`. A + deterministic `general + pool` overflow is rejected at config load. +- The final `ParameterStatus` messages sent to the client no longer + overwrite operator-managed GUC names, so the client-visible values + match the backend checkout state. +- `auth_query` now rebuilds a dynamic pool after a successful MD5 + refetch, rejects the stale-overlay race in `create_dynamic_pool`, and + accepts native `json`/`jsonb` startup_parameter columns without a + `::text` cast. +- `/api/config` and `/api/pools` show literal startup_parameter values + only to `Admin`; SSO readers get the masked view. `/api/config` also + marks `general.host`, `general.port`, `web.host`, and `web.port` as + restart-required. +- Prometheus rules now cover PostgreSQL-side rejection, budget overflow, + malformed auth_query columns, dedicated-mode drops, and rejected SSO + credentials sent over insecure transport. +- Each pool now precomputes the merged startup map, budget decision, and + canonical operator-key set. Backend checkout reuses those cached + values instead of cloning and recalculating the map each time. + ### 3.9.0 Per-pool PostgreSQL startup parameters. pg_doorman can now add diff --git a/documentation/en/src/guides/web-ui.md b/documentation/en/src/guides/web-ui.md index 15b39ec21..5ad4fe2e3 100644 --- a/documentation/en/src/guides/web-ui.md +++ b/documentation/en/src/guides/web-ui.md @@ -130,7 +130,7 @@ proxy: | `sso_groups_claim` | Name of the JWT claim that carries the user's group memberships. Read together with `sso_admin_groups`. | `"groups"` | | `sso_admin_groups` | Group names that promote an SSO user to `Admin`. Empty keeps every SSO login at the read-only `Sso` role. | `[]` | | `sso_require_https` | Reject Bearer/cookie/query SSO credentials presented over plain HTTP. The listener treats a request as secure only when the TCP peer is in `trusted_proxies` and `X-Forwarded-Proto: https` is forwarded. Defaults to off so SSO keeps working through a TLS-terminating proxy that reaches pg_doorman over a private HTTP leg. | `false` | -| `trusted_proxies` | CIDR ranges trusted to set `X-Forwarded-For` / `Forwarded` / `X-Forwarded-Proto`. Empty trusts only the listener's own peer. See [Access log](#access-log). | `[]` | +| `trusted_proxies` | CIDR ranges trusted to set `X-Forwarded-For` / `Forwarded` / `X-Forwarded-Proto`. With an empty list, pg_doorman ignores forwarded headers and uses the TCP peer address. If `sso_require_https = true` is behind a TLS-terminating proxy, add that proxy CIDR so `X-Forwarded-Proto: https` is accepted. See [Access log](#access-log). | `[]` | ### Promoting SSO users to Admin via group claim diff --git a/documentation/en/src/tutorials/startup-parameters.md b/documentation/en/src/tutorials/startup-parameters.md index 1a3f0709b..dbf2c819b 100644 --- a/documentation/en/src/tutorials/startup-parameters.md +++ b/documentation/en/src/tutorials/startup-parameters.md @@ -52,15 +52,27 @@ FROM pg_authid WHERE rolname = $1; ``` -The column must serialize as `text`. If the SQL returns `json` or -`jsonb`, add an explicit `::text` cast. pg_doorman reads the column -as `text` and logs a warning for that fetched row when the type does -not match. +The column may be `text`, `json`, or `jsonb`; pg_doorman dispatches by +the column type without requiring a cast. The content must be a JSON +object whose values are strings. Other PostgreSQL types (or a custom +domain on top of `jsonb`) log a warning and the per-user overlay is +ignored. Dedicated `auth_query` mode (`server_user` set) ignores the per-user column and logs once per (pool, username): one shared backend serves many users, so a per-user override cannot apply. +Changes to a per-user `startup_parameters` row apply to **new** backend +connections, but only after pg_doorman re-reads the row. The +`auth_query` cache holds positive entries for `auth_query.cache_ttl` +(default one hour) and on a refresh detects the overlay change and +drops the dynamic pool so the next login rebuilds it against the new +values. Until the cache entry expires, reconnecting clients still see +the old overlay. To force an immediate rollout: lower `cache_ttl` and +reload the config, restart pg_doorman, or wait for the TTL to elapse. +Backends that are already checked out keep the values captured when +their pool was created. + ## What pg_doorman does with the values pg_doorman adds the resolved parameter set to the PostgreSQL @@ -92,21 +104,24 @@ At config load: - Keys must match PG GUC naming `^[A-Za-z_][A-Za-z0-9_.]*$`. Namespaced names like `auto_explain.log_min_duration` are accepted; arbitrary punctuation is not. -- Reserved keys (`user`, `database`, `replication`, `options`, and - anything starting with `_pq_.`) are refused. pg_doorman manages - them itself or PG treats them specially in the StartupMessage. +- Reserved keys (`user`, `database`, `replication`, `options`, `role`, + `session_authorization`, and anything starting with `_pq_.`) are + refused. pg_doorman manages them itself or PG treats them specially in + the StartupMessage. - Values must not contain null bytes. - Each level (general or per-pool) must fit within the startup-parameter budget: `MAX_STARTUP_PACKET_LENGTH` (10 000 bytes) minus 512 bytes reserved for pg_doorman-managed keys. -Before each backend spawn pg_doorman checks the resolved parameter set -against the same cap. Two layers that fit on their own can overflow once -`auth_query` adds a third layer. If only the `auth_query` layer pushes -the set over the cap, pg_doorman drops that layer and keeps the -general/pool baseline. If the baseline itself or the full startup packet -does not fit, pg_doorman skips all configured parameters for that -spawn and logs the byte counts. +Before each backend start, pg_doorman checks the resolved parameter set +against the same cap. Layers that fit individually can exceed the limit +after merging: `general + pool` can already be too large, and an +`auth_query` row can push a valid baseline over the limit. Any overflow +now returns a PostgreSQL-style error (`SQLSTATE 53400`) to the client +instead of sending a partial or empty `StartupMessage`. The warning log +records the byte counts, and +`pg_doorman_startup_parameters_dropped_total` increments for each +rejected backend start. ## What happens when PG rejects a parameter diff --git a/documentation/ru/src/authentication/auth-query.md b/documentation/ru/src/authentication/auth-query.md index 338d112c1..a4820d55b 100644 --- a/documentation/ru/src/authentication/auth-query.md +++ b/documentation/ru/src/authentication/auth-query.md @@ -30,10 +30,11 @@ pools: Запрос должен возвращать колонку с именем `passwd` или `password`, содержащую хеш MD5 или SCRAM. Дополнительные колонки игнорируются, кроме -необязательной `startup_parameters`. В passthrough-режиме pg_doorman -читает её как JSON-объект в `text` с пользовательскими параметрами -запуска PostgreSQL. Dedicated-режим игнорирует её и пишет -предупреждение. +необязательной `startup_parameters`. В сквозном режиме pg_doorman +читает эту колонку как `text`, `json` или `jsonb`; значение должно быть +JSON-объектом с параметрами запуска PostgreSQL для конкретного +пользователя. В выделенном режиме эта колонка игнорируется, а в лог +пишется предупреждение. `user` и `password` — это учётные данные, под которыми pg_doorman выполняет lookup-запрос. У них должно быть право читать колонку с учётными данными. Либо выдайте доступ к специально созданному представлению (рекомендуется), либо используйте пользователя из группы `pg_read_server_files`. diff --git a/documentation/ru/src/guides/web-ui.md b/documentation/ru/src/guides/web-ui.md index 909767841..e83c06d6a 100644 --- a/documentation/ru/src/guides/web-ui.md +++ b/documentation/ru/src/guides/web-ui.md @@ -130,7 +130,7 @@ SSO опциональный. По умолчанию (`[web].sso_enabled = fals | `sso_groups_claim` | Имя JWT-claim, в котором лежат группы пользователя. Читается вместе с `sso_admin_groups`. | `"groups"` | | `sso_admin_groups` | Группы, которые поднимают SSO-пользователя до `Admin`. Пустой список оставляет каждый SSO-логин на роли `Sso` только для чтения. | `[]` | | `sso_require_https` | Отклонять Bearer/cookie/query SSO-credentials, пришедшие по plain HTTP. Запрос считается защищённым только если TCP-peer входит в `trusted_proxies` и прокси прислал `X-Forwarded-Proto: https`. По умолчанию выключено, чтобы SSO продолжал работать в схеме «TLS терминирует прокси → pg_doorman слушает HTTP во внутренней сети». | `false` | -| `trusted_proxies` | CIDR доверенных обратных прокси (используется для `X-Forwarded-For` / `Forwarded` / `X-Forwarded-Proto`). Пустой список — доверять только непосредственному TCP-peer. См. [Журнал доступа](#журнал-доступа). | `[]` | +| `trusted_proxies` | CIDR доверенных обратных прокси для `X-Forwarded-For`, `Forwarded` и `X-Forwarded-Proto`. При пустом списке pg_doorman игнорирует эти заголовки и берёт адрес прямого TCP-пира. Если `sso_require_https = true` работает за прокси, который завершает TLS, добавьте CIDR этого прокси, чтобы доверять `X-Forwarded-Proto: https`. См. [Журнал доступа](#журнал-доступа). | `[]` | ### Поднятие SSO-пользователя до Admin через claim с группами diff --git a/documentation/ru/src/reference/general.md b/documentation/ru/src/reference/general.md index badb9ba8f..926faa35c 100644 --- a/documentation/ru/src/reference/general.md +++ b/documentation/ru/src/reference/general.md @@ -684,9 +684,11 @@ hostnossl all all 192.168.1.0/24 trust протокольные ключи (`user`, `database`, `replication`, `options`, `_pq_.*`), имена GUC, нулевые байты и размер этого уровня. Перед каждым запуском бэкенда объединённый набор параметров снова проверяется по -лимиту `MAX_STARTUP_PACKET_LENGTH` PostgreSQL; если он не помещается, -pg_doorman пропускает операторские параметры для этого запуска и пишет -предупреждение. +лимиту `MAX_STARTUP_PACKET_LENGTH` PostgreSQL (10 000 байт). Каскад и +проверка пакета отклоняют запуск бэкенда с SQLSTATE 53400 +(`configuration_limit_exceeded`); если каскад выходит за бюджет только +из-за overlay из `auth_query`, pg_doorman сбрасывает этот overlay и +продолжает с baseline. Если PostgreSQL отвергает параметр при запуске бэкенда, pg_doorman возвращает клиенту `ErrorResponse` PostgreSQL без изменений: повторной diff --git a/documentation/ru/src/tutorials/startup-parameters.md b/documentation/ru/src/tutorials/startup-parameters.md index 7c1fdd0ae..27fb0d969 100644 --- a/documentation/ru/src/tutorials/startup-parameters.md +++ b/documentation/ru/src/tutorials/startup-parameters.md @@ -42,10 +42,11 @@ work_mem = "64MB" настроек PostgreSQL по умолчанию. Уже открытые бэкенды не меняются: новые значения вступают в силу по мере ротации соединений. -В passthrough-режиме `auth_query`, когда `server_user` не задан, запрос +В сквозном режиме `auth_query`, когда `server_user` не задан, запрос аутентификации может вернуть необязательную колонку `startup_parameters` -типа `text` с JSON-объектом. Значения из этой колонки переопределяют -`general` и настройки пула только для конкретного пользователя. +типа `text`, `json` или `jsonb` с JSON-объектом. Значения из этой +колонки переопределяют `general` и настройки пула только для +конкретного пользователя. ```sql SELECT @@ -58,15 +59,28 @@ FROM pg_authid WHERE rolname = $1; ``` -Колонка должна возвращаться как `text`. Если SQL отдаёт `json` или -`jsonb`, добавьте явное приведение типа `::text`. pg_doorman читает её -именно как `text` и пишет предупреждение для полученной строки, если -тип не совпал. +pg_doorman выбирает декодер по типу колонки, поэтому приведение `::text` +не требуется. Содержимое должно быть JSON-объектом, а значения — +строками. Для других типов PostgreSQL, в том числе доменов поверх +`jsonb`, pg_doorman пишет предупреждение и игнорирует пользовательские +параметры из этой строки. -Dedicated-режим `auth_query`, когда `server_user` задан, игнорирует эту +Выделенный режим `auth_query`, когда `server_user` задан, игнорирует эту колонку и один раз пишет предупреждение на пару `(пул, пользователь)`. В этом режиме один серверный пул обслуживает разных пользователей, -поэтому per-user значения применить нельзя. +поэтому пользовательские значения применить нельзя. + +Изменения в пользовательской строке `startup_parameters` применяются +только к **новым** серверным подключениям и только после того, как +pg_doorman перечитает строку. Кеш `auth_query` хранит положительные +записи в течение `auth_query.cache_ttl` (по умолчанию час); при +обновлении кеша pg_doorman замечает изменение overlay и сбрасывает +динамический пул, чтобы следующий логин пересоздал его с новыми +значениями. Пока запись кеша не истекла, новые подключения клиентов +получают прежний overlay. Чтобы выкатить изменения немедленно: уменьшить +`cache_ttl` и сделать reload конфигурации, перезапустить pg_doorman или +дождаться истечения TTL. Бэкенды, которые уже выданы клиентам, +продолжают работать со значениями, сохранёнными при создании пула. ## Что pg_doorman делает со значениями @@ -99,31 +113,34 @@ checkout=> SET plan_cache_mode = 'auto'; RESET ALL; SHOW plan_cache_mode; `^[A-Za-z_][A-Za-z0-9_.]*$`. Составные имена вроде `auto_explain.log_min_duration` допустимы; произвольная пунктуация нет. -- Зарезервированные ключи (`user`, `database`, `replication`, `options` - и всё, что начинается с `_pq_.`) отклоняются. pg_doorman управляет - ими сам, либо PostgreSQL обрабатывает их в `StartupMessage` особым - образом. +- Зарезервированные ключи (`user`, `database`, `replication`, `options`, + `role`, `session_authorization` и всё, что начинается с `_pq_.`) + отклоняются. pg_doorman управляет ими сам, либо PostgreSQL + обрабатывает их в `StartupMessage` особым образом. - Значения не должны содержать нулевой байт. - Каждый уровень (`general` или `pool`) должен помещаться в лимит для операторских параметров: `MAX_STARTUP_PACKET_LENGTH` (10000 байт) минус 512 байт, зарезервированных под служебные ключи pg_doorman. -Перед запуском каждого бэкенда pg_doorman заново проверяет объединённый -набор параметров по тому же лимиту. Два уровня, которые помещались по -отдельности, могут вместе выйти за лимит, особенно когда `auth_query` -добавляет третий слой. Если лимит превышает только слой `auth_query`, -pg_doorman отбрасывает этот слой и сохраняет baseline из `general` и -пула. Если не помещается сам baseline или полный startup-пакет, -pg_doorman пропускает все операторские параметры для этого -запуска и пишет размеры в лог. +Перед запуском каждого бэкенда pg_doorman снова проверяет объединённый +набор параметров по тому же лимиту. Слои, которые помещаются по +отдельности, могут превысить лимит после объединения: `general + pool` +может быть слишком большим сам по себе, а строка `auth_query` может +переполнить уже допустимый базовый набор. Любое превышение теперь +возвращается клиенту как PostgreSQL-ошибка с `SQLSTATE 53400`; пустой +или урезанный `StartupMessage` не отправляется. В предупреждении +записываются размеры в байтах, а +`pg_doorman_startup_parameters_dropped_total` увеличивается на каждой +отклонённой попытке запуска бэкенда. ## Что происходит, если PG отвергает параметр Если PostgreSQL отвергает заданный оператором параметр при запуске бэкенда, pg_doorman возвращает клиенту `ErrorResponse` PostgreSQL без -изменений. Клиент видит тот же sqlstate (`22023`, `42704`, `42501`, -`55P02` или любой другой код из стартового семейства) и то же сообщение, -что увидел бы при прямом подключении к PostgreSQL. +изменений. Клиент видит тот же `SQLSTATE` (`22023`, `42704`, `42501`, +`55P02` или любой другой код, который PostgreSQL вернул при отклонении +`StartupMessage`) и то же сообщение, что увидел бы при прямом +подключении к PostgreSQL. pg_doorman не пробует повторить подключение без отклонённого параметра и не отключает этот ключ автоматически для пула. Следующее подключение @@ -153,8 +170,14 @@ admin> SHOW STARTUP_PARAMETERS; параметра, заданного оператором. Имя параметра и пользователя пишутся в строку лога уровня `warn`; в лейблы они не включены, чтобы динамические `auth_query`-пулы не раздували количество серий. - -Разумная отправная точка для алерта: если +- `pg_doorman_startup_parameters_dropped_total{pool, reason}` считает + случаи, когда pg_doorman отклонил `startup_parameters` до отправки + `StartupMessage`: превышение лимита, неподдерживаемый тип или + неверный JSON из `auth_query`, недопустимые ключи или значения, а + также пользовательские значения, проигнорированные в выделенном + режиме. + +Практичное условие для алерта: если `pg_doorman_backend_startup_parameter_errors_total` растёт по одному и тому же пулу несколько минут подряд, новые подключения к этому пулу падают на одном и том же GUC. Конфигурацию нужно исправить до возврата @@ -177,8 +200,8 @@ admin> SHOW STARTUP_PARAMETERS; - [Общие настройки](../reference/general.md): `startup_parameters`. - [Настройки пула](../reference/pool.md): `pools..startup_parameters`. -- [auth_query](../authentication/auth-query.md): passthrough- и - dedicated-режимы, чтение колонки `startup_parameters`. +- [auth_query](../authentication/auth-query.md): сквозной и выделенный + режимы, чтение колонки `startup_parameters`. - [Команды администратора](../observability/admin-commands.md): `SHOW STARTUP_PARAMETERS`. - [Метрики Prometheus](../reference/prometheus.md): полный список. diff --git a/documentation/ru/src/tutorials/troubleshooting.md b/documentation/ru/src/tutorials/troubleshooting.md index ab73289fc..8d90e33b4 100644 --- a/documentation/ru/src/tutorials/troubleshooting.md +++ b/documentation/ru/src/tutorials/troubleshooting.md @@ -18,7 +18,7 @@ SELECT usename, passwd FROM pg_shadow WHERE usename = 'your_user'; ### Когда username пула отличается от роли на backend -Когда обращённый к клиенту `username` в PgDoorman не совпадает с реальной ролью PostgreSQL, passthrough работать не может: у pg_doorman нет пароля для backend-роли. Дайте явные credentials: +Когда `username`, под которым клиент подключается к pg_doorman, не совпадает с реальной ролью PostgreSQL, passthrough работать не может: у pg_doorman нет пароля для backend-роли. Укажите явные credentials: ```yaml users: diff --git a/frontend/dist/.source-hash b/frontend/dist/.source-hash index 1457695a1..499933def 100644 --- a/frontend/dist/.source-hash +++ b/frontend/dist/.source-hash @@ -1 +1 @@ -10634db6ec2fc51c3c7ce7485d71f1939758e0e39350e7f46c88a6742f08486c +c8bfd5e0b4dac8f2ea1c9087b2da622babd97e9e6147fbbd82e0ae2cbd5d15f7 diff --git a/frontend/dist/assets/index-n2o-PXii.js.gz b/frontend/dist/assets/index-Dg2wtyD8.js.gz similarity index 67% rename from frontend/dist/assets/index-n2o-PXii.js.gz rename to frontend/dist/assets/index-Dg2wtyD8.js.gz index 1c992a0ea..ebbd8338a 100644 Binary files a/frontend/dist/assets/index-n2o-PXii.js.gz and b/frontend/dist/assets/index-Dg2wtyD8.js.gz differ diff --git a/frontend/dist/index.html.gz b/frontend/dist/index.html.gz index 8653a516d..1800d71e9 100644 Binary files a/frontend/dist/index.html.gz and b/frontend/dist/index.html.gz differ diff --git a/frontend/src/components/AuthGate.tsx b/frontend/src/components/AuthGate.tsx index c56520cf1..26b82635c 100644 --- a/frontend/src/components/AuthGate.tsx +++ b/frontend/src/components/AuthGate.tsx @@ -430,8 +430,8 @@ function BasicBlock({ > That user/password was rejected. Recheck{" "} [general].admin_username and{" "} - [general].admin_password in{" "} - pg_doorman.toml. + [general].admin_password in the + active config file.

)}
diff --git a/grafana/pg_doorman.json b/grafana/pg_doorman.json index 7f77e562d..981f765a5 100644 --- a/grafana/pg_doorman.json +++ b/grafana/pg_doorman.json @@ -3086,7 +3086,7 @@ } ], "title": "PG-Side Rejections by SQLSTATE", - "description": "Per-pool rate of backend startups PG rejected because of an operator-supplied parameter. Split by SQLSTATE: 22023 invalid_value, 42704 undefined_object, 42501 insufficient_privilege, 55P02 cant_change_runtime_param. Non-zero for the same pool over a few minutes means every connect through that pool fails on the same operator GUC \u2014 fix general/pool/auth_query. Filters on $user/$database do not apply \u2014 the counter has only the `pool` label.", + "description": "Rate of backend startups PostgreSQL rejected because of a configured startup parameter. Split by SQLSTATE: 22023 invalid_value, 42704 undefined_object, 42501 insufficient_privilege, 55P02 cant_change_runtime_param. Growth for the same pool over several minutes means new connects through that pool fail on the same GUC; fix general/pool/auth_query. Filters on $user/$database do not apply: the counter has only the `pool` label.", "datasource": { "uid": "prometheus" }, @@ -3178,7 +3178,7 @@ } ], "title": "Pre-Wire Drops by Reason", - "description": "Operator-supplied entries pg_doorman dropped BEFORE the StartupMessage went on the wire \u2014 the failure mode the PG-side counter above cannot see. Reasons: cascade_budget_exceeded (merged map past 9 488 bytes), packet_cap_exceeded (full packet past PG MAX_STARTUP_PACKET_LENGTH 10 000 bytes), auth_query_oversize (per-user JSON column past operator budget), auth_query_overlay_oversize (overlay pushes cascade over budget but baseline alone fits \u2014 pg_doorman ships the baseline), auth_query_bad_type / auth_query_invalid_json / auth_query_invalid_shape (column type, JSON parse, or non-object payload), auth_query_invalid_entry (one or more JSON entries failed validation), dedicated_mode (per-user GUC ignored because the pool shares one backend across users). Non-zero on any reason needs operator attention. Filters on $user/$database do not apply \u2014 the counter has only `pool` and `reason` labels.", + "description": "Startup parameter drop events before pg_doorman sends StartupMessage. Reasons: cascade_budget_exceeded (resolved set above 9 488 bytes), packet_cap_exceeded (full packet above PG MAX_STARTUP_PACKET_LENGTH 10 000 bytes), auth_query_oversize (per-user JSON column above the startup-parameter budget), auth_query_overlay_oversize (auth_query overlay overflows but the general/pool baseline still fits), auth_query_bad_type / auth_query_invalid_json / auth_query_invalid_shape (column type, JSON parse, or non-object payload), auth_query_invalid_entry (one or more JSON entries failed validation), dedicated_mode (per-user GUC ignored because the pool shares one backend across users). Any increase should be investigated. Filters on $user/$database do not apply: the counter has only `pool` and `reason` labels.", "datasource": { "uid": "prometheus" }, diff --git a/monitoring/prometheus-rules/startup-parameters.yaml b/monitoring/prometheus-rules/startup-parameters.yaml new file mode 100644 index 000000000..05ae9543d --- /dev/null +++ b/monitoring/prometheus-rules/startup-parameters.yaml @@ -0,0 +1,126 @@ +groups: + - name: startup_parameters_alerts + interval: 30s + rules: + - alert: PgDoormanStartupParameterPgRejection + expr: | + sum by (job, instance, pool, sqlstate) ( + increase(pg_doorman_backend_startup_parameter_errors_total[5m]) + ) > 0 + for: 1m + labels: + severity: warning + service: pg_doorman + annotations: + summary: "pg_doorman: PG rejected configured startup_parameter for pool {{ $labels.pool }} (SQLSTATE {{ $labels.sqlstate }})" + runbook: | + PostgreSQL rejected at least one backend startup for this + pool because of a startup_parameter sent in StartupMessage. + The metric is labelled only by pool and SQLSTATE, so the + source may be one auth_query row or the pool-wide baseline. + Use the warning log line to find the failing key and + username, then fix the value in general/pool config or in + the auth_query row. SQLSTATE usually identifies the cause: + 22023 invalid_value, 42704 undefined_object, 42501 + insufficient_privilege, or 55P02 cannot_alter_runtime_param. + + - alert: PgDoormanStartupParameterCascadeOverflow + expr: | + sum by (job, instance, pool, reason) ( + increase(pg_doorman_startup_parameters_dropped_total{reason=~"cascade_budget_exceeded|packet_cap_exceeded"}[5m]) + ) > 0 + for: 1m + labels: + severity: critical + service: pg_doorman + annotations: + summary: "pg_doorman: startup_parameters cascade does not fit StartupMessage budget for pool {{ $labels.pool }}" + runbook: | + The merged general/pool/auth_query map does not fit the + PostgreSQL StartupMessage budget. pg_doorman rejects backend + startup for this pool until the configured map is smaller. + Trim general.startup_parameters, + pools..startup_parameters, or the per-user auth_query + row. The reason label shows which gate fired: + cascade_budget_exceeded or packet_cap_exceeded. + + - alert: PgDoormanStartupParameterOverlayOversize + expr: | + sum by (job, instance, pool) ( + increase(pg_doorman_startup_parameters_dropped_total{reason="auth_query_overlay_oversize"}[5m]) + ) > 0 + for: 1m + labels: + severity: warning + service: pg_doorman + annotations: + summary: "pg_doorman: auth_query startup_parameters exceed budget for pool {{ $labels.pool }}" + runbook: | + A per-user auth_query row makes the resolved startup map + exceed the StartupMessage budget. pg_doorman returns + SQLSTATE 53400 and rejects backend startup for the affected + dynamic pool. Shrink the per-user startup_parameters row, + the pool baseline, or the general baseline. + + - alert: PgDoormanStartupParameterOverlayRaw + expr: | + sum by (job, instance, pool) ( + increase(pg_doorman_startup_parameters_dropped_total{reason="auth_query_oversize"}[5m]) + ) > 0 + for: 1m + labels: + severity: warning + service: pg_doorman + annotations: + summary: "pg_doorman: auth_query startup_parameters row exceeds parse-time budget for pool {{ $labels.pool }}" + runbook: | + The auth_query row for a user returned a + startup_parameters payload too large to be a valid + StartupMessage overlay. pg_doorman ignores the per-user + values for that row, so the affected user authenticates + without those GUCs. Inspect the auth_query view for the + affected user and shrink the row. + + - alert: PgDoormanStartupParameterAuthQueryInvalid + expr: | + sum by (job, instance, pool, reason) ( + increase(pg_doorman_startup_parameters_dropped_total{reason=~"auth_query_bad_type|auth_query_invalid_json|auth_query_invalid_shape|auth_query_invalid_entry"}[5m]) + ) > 0 + for: 5m + labels: + severity: warning + service: pg_doorman + annotations: + summary: "pg_doorman: auth_query startup_parameters column malformed for pool {{ $labels.pool }} (reason {{ $labels.reason }})" + runbook: | + auth_query returned a startup_parameters value pg_doorman + could not use: unsupported PG type, malformed JSON, + non-object top level, or invalid key/value entry. Affected + users authenticate, but backend startup proceeds without + their per-user GUCs. Fix the auth_query SELECT or the + underlying row. The warning log line has the exact failure. + + - alert: PgDoormanStartupParameterDedicatedDrops + # `dedicated_mode` ticks only when auth_query actually fetches + # a row, not on cache hits. With the default `cache_ttl = 1h` + # the counter fires once per user per hour at most, so the + # window has to span at least one full TTL or the alert will + # silently miss a steady misconfiguration. `for: 5m` keeps the + # signal from flapping on a single noisy refetch. + expr: | + sum by (job, instance, pool) ( + increase(pg_doorman_startup_parameters_dropped_total{reason="dedicated_mode"}[1h]) + ) > 0 + for: 5m + labels: + severity: warning + service: pg_doorman + annotations: + summary: "pg_doorman: auth_query dedicated mode ignored per-user startup_parameters for pool {{ $labels.pool }}" + runbook: | + Dedicated auth_query mode uses a shared server_user pool + and cannot apply per-user startup_parameters. If the + operator filled that column, either switch the pool to + passthrough mode or remove the per-user values. A sustained + non-zero rate means the auth_query SELECT keeps returning + values that this pool mode cannot use. diff --git a/monitoring/prometheus-rules/web-sso.yaml b/monitoring/prometheus-rules/web-sso.yaml index a7cf33201..18525e145 100644 --- a/monitoring/prometheus-rules/web-sso.yaml +++ b/monitoring/prometheus-rules/web-sso.yaml @@ -122,3 +122,25 @@ groups: minutes. Common causes: SSO operators trying admin paths (403), expired tokens reaching API after silent refresh failed, broken proxy adding malformed Authorization headers. + + - alert: PgDoormanWebSsoInsecureTransport + expr: | + sum by (job, instance) ( + clamp_min( + rate(pg_doorman_web_sso_validation_errors_total{reason="insecure_transport"}[5m]), + 0 + ) + ) > 0 + for: 5m + labels: + severity: warning + service: pg_doorman + annotations: + summary: "pg_doorman: SSO requests rejected over insecure transport on {{ $labels.instance }}" + runbook: | + sso_require_https is enabled, and pg_doorman received an + SSO credential on a request it classified as plain HTTP. + Usually either the TLS-terminating proxy is missing from + trusted_proxies, so X-Forwarded-Proto is ignored, or the + client is really using HTTP. Check trusted_proxies, the + proxy's X-Forwarded-Proto header, and the client URL. diff --git a/pg_doorman.toml b/pg_doorman.toml index 574813c5f..34ce4732a 100644 --- a/pg_doorman.toml +++ b/pg_doorman.toml @@ -416,8 +416,9 @@ hba = [] # validates reserved keys, GUC names, null bytes, and this level's # size. Before backend startup, pg_doorman checks the resolved # parameter set again; if it does not fit PG's -# MAX_STARTUP_PACKET_LENGTH (10000 bytes), pg_doorman skips -# configured GUCs for that startup and logs a warning. +# MAX_STARTUP_PACKET_LENGTH (10000 bytes), pg_doorman rejects +# backend startup with SQLSTATE 53400 (configuration_limit_exceeded) +# and writes the event to the warning log. # Example: startup_parameters = { plan_cache_mode = "force_custom_plan" } # Default: {} (empty) # startup_parameters = { plan_cache_mode = "force_custom_plan", work_mem = "64MB" } diff --git a/pg_doorman.yaml b/pg_doorman.yaml index f8f3082e1..abb4b7253 100644 --- a/pg_doorman.yaml +++ b/pg_doorman.yaml @@ -456,8 +456,9 @@ general: # validates reserved keys, GUC names, null bytes, and this level's # size. Before backend startup, pg_doorman checks the resolved # parameter set again; if it does not fit PG's - # MAX_STARTUP_PACKET_LENGTH (10000 bytes), pg_doorman skips - # configured GUCs for that startup and logs a warning. + # MAX_STARTUP_PACKET_LENGTH (10000 bytes), pg_doorman rejects + # backend startup with SQLSTATE 53400 (configuration_limit_exceeded) + # and writes the event to the warning log. # Example: startup_parameters = { plan_cache_mode = "force_custom_plan" } # Default: {} (empty) # startup_parameters: diff --git a/src/admin/show.rs b/src/admin/show.rs index f19a91478..c85dd084c 100644 --- a/src/admin/show.rs +++ b/src/admin/show.rs @@ -432,8 +432,12 @@ where { let config = &get_config(); let config: HashMap = config.into(); - // Configs that cannot be changed without restarting. - let immutables = ["host", "port", "connect_timeout"]; + // Configs that cannot be changed without restarting. The keys here + // are the bare names that `From<&Config> for HashMap` emits — the + // Web `/api/config` view uses flattened paths (`general.host`, + // `web.host`, …) and its own matcher in `web::routes::collect::config`. + // `connect_timeout` is reloadable on SIGHUP; do not list it. + let immutables = ["host", "port"]; // Columns let columns = vec![ ("key", DataType::Text), diff --git a/src/app/generate/fields.yaml b/src/app/generate/fields.yaml index 01ad67927..5b8871d11 100644 --- a/src/app/generate/fields.yaml +++ b/src/app/generate/fields.yaml @@ -1158,8 +1158,9 @@ fields: validates reserved keys, GUC names, null bytes, and this level's size. Before backend startup, pg_doorman checks the resolved parameter set again; if it does not fit PG's - MAX_STARTUP_PACKET_LENGTH (10000 bytes), pg_doorman skips - configured GUCs for that startup and logs a warning. + MAX_STARTUP_PACKET_LENGTH (10000 bytes), pg_doorman rejects + backend startup with SQLSTATE 53400 (configuration_limit_exceeded) + and writes the event to the warning log. Example: startup_parameters = { plan_cache_mode = "force_custom_plan" } ru: | Базовые GUC PostgreSQL, которые pg_doorman добавляет в @@ -1170,15 +1171,16 @@ fields: зарезервированные ключи, имена GUC, нулевые байты и размер этого уровня. Объединённый набор параметров снова проверяется при запуске бэкенда; если он не помещается в лимит PG - MAX_STARTUP_PACKET_LENGTH (10000 байт), pg_doorman пропускает - операторские GUC для этого запуска и пишет предупреждение. + MAX_STARTUP_PACKET_LENGTH (10000 байт), pg_doorman отклоняет + запуск бэкенда с SQLSTATE 53400 (configuration_limit_exceeded) + и пишет предупреждение в лог. Пример: startup_parameters = { plan_cache_mode = "force_custom_plan" } doc: | Map of PostgreSQL configuration parameter names to string values. pg_doorman writes them into each new backend `StartupMessage`; PostgreSQL stores them as the session reset defaults, so client `RESET ALL` / `DISCARD ALL` returns to these values. Cascade order: `general.startup_parameters`, then `pools..startup_parameters`, then the optional `startup_parameters` JSON column returned by passthrough `auth_query`. Later layers win per key. Dedicated-mode `auth_query` pools ignore the per-user column because one shared backend serves multiple roles. - Validation at config load rejects reserved protocol keys (`user`, `database`, `replication`, `options`, anything starting with `_pq_.`), invalid GUC names, null bytes, and per-level maps that exceed the startup-parameter budget. Before each backend startup, pg_doorman checks the resolved parameter set against PG's `MAX_STARTUP_PACKET_LENGTH` (10 000 bytes); if it does not fit, pg_doorman drops the auth_query overlay when the baseline still fits, otherwise it drops all configured keys for that startup and logs the event. + Validation at config load rejects reserved protocol keys (`user`, `database`, `replication`, `options`, anything starting with `_pq_.`), invalid GUC names, null bytes, and per-level maps that exceed the startup-parameter budget. Before each backend startup, pg_doorman checks the resolved parameter set against PG's `MAX_STARTUP_PACKET_LENGTH` (10 000 bytes). The cascade and packet-cap gates reject backend startup with SQLSTATE 53400 (`configuration_limit_exceeded`); when only the auth_query overlay tips the cascade over budget, pg_doorman drops that overlay and continues with the baseline. If PostgreSQL rejects a parameter at backend startup, pg_doorman returns PostgreSQL's `ErrorResponse` to the client unchanged. There is no retry with the key removed, and pg_doorman does not automatically disable that key for the pool. The cumulative count is exported as `pg_doorman_backend_startup_parameter_errors_total{pool, sqlstate}`; the parameter name and username are written to the corresponding warning log line. @@ -1641,7 +1643,7 @@ fields: doc: | SQL query to fetch credentials. It must return a column named `passwd` or `password` containing the MD5 or SCRAM hash. If the query returns exactly one column, it is used regardless of name. - Extra columns are ignored except for the optional `startup_parameters` text column. In passthrough mode, pg_doorman reads that column as a JSON object with per-user PostgreSQL startup parameters. Dedicated mode ignores it and logs a warning. Use `$1` as the placeholder for the username parameter. + Extra columns are ignored except for the optional `startup_parameters` column. The column may be `text`, `json`, or `jsonb`; pg_doorman dispatches by the column type and the content must be a JSON object whose values are strings. Custom domains over `jsonb` are not accepted without an explicit cast. In passthrough mode, the map applies as per-user startup parameters. Dedicated mode ignores it and logs a warning. Use `$1` as the placeholder for the username parameter. Example: `"SELECT passwd FROM pg_shadow WHERE usename = $1"` diff --git a/src/auth/auth_query.rs b/src/auth/auth_query.rs index 9a0eebd87..034fd9a4e 100644 --- a/src/auth/auth_query.rs +++ b/src/auth/auth_query.rs @@ -17,6 +17,7 @@ use log::{debug, error, info, warn}; use crate::utils::format_elapsed; use tokio::sync::mpsc; use tokio::sync::Mutex as TokioMutex; +use tokio_postgres::types::{FromSql, Type}; use tokio_postgres::{Client, NoTls}; use crate::config::{AuthQueryConfig, Duration}; @@ -29,6 +30,46 @@ use crate::stats::auth_query::AuthQueryStats; /// memory exhaustion from very long usernames. const MAX_USERNAME_LEN: usize = 63; +/// Marker string the custom `LimitedJson` decoder embeds in its error +/// message when a `json`/`jsonb` row exceeds `MAX_OPERATOR_BUDGET`. The +/// outer `try_get` error has to stay generic (`Box`), so the +/// auth_query reader matches on this substring to decide between +/// `auth_query_oversize` (operator footgun) and `auth_query_bad_type` +/// (decoder failure on unexpected wire bytes). +const LIMITED_JSON_OVERSIZE_TAG: &str = "auth_query startup_parameters oversize"; + +/// Custom `FromSql` wrapper for `json`/`jsonb` columns that enforces +/// `MAX_OPERATOR_BUDGET` on the raw wire bytes BEFORE `serde_json` walks +/// the value tree. Without this, a malicious or accidentally large +/// `jsonb` row on the auth_query path would force pg_doorman to +/// allocate and parse the full tree on every cache miss, even though +/// the result is later rejected by the size gate inside +/// `parse_startup_parameters_value`. The text decoder already has this +/// pre-parse gate; this wrapper closes the asymmetry for jsonb. +struct LimitedJson(serde_json::Value); + +impl<'a> FromSql<'a> for LimitedJson { + fn from_sql( + ty: &Type, + raw: &'a [u8], + ) -> Result> { + let budget = crate::config::startup_parameters::MAX_OPERATOR_BUDGET; + if raw.len() > budget { + return Err(format!( + "{LIMITED_JSON_OVERSIZE_TAG}: raw column is {} bytes, exceeding operator budget {budget}", + raw.len() + ) + .into()); + } + let value = ::from_sql(ty, raw)?; + Ok(LimitedJson(value)) + } + + fn accepts(ty: &Type) -> bool { + matches!(*ty, Type::JSON | Type::JSONB) + } +} + // --------------------------------------------------------------------------- // PasswordFetcher trait (allows mocking AuthQueryExecutor in unit tests) // --------------------------------------------------------------------------- @@ -386,11 +427,14 @@ impl AuthQueryExecutor { } } - /// Read the optional `startup_parameters` text column from the auth_query - /// row and parse it as a JSON object. A missing column yields an empty - /// map; a present column whose type does not coerce to `Option` - /// logs a warning and yields an empty map. Actual JSON parsing and - /// per-entry validation happen in `parse_startup_parameters_text`. + /// Read the optional `startup_parameters` column from the auth_query + /// row. A missing column yields an empty map. Column type drives the + /// decoder so a `jsonb` row goes straight through `serde_json::Value` + /// without a `text`-decode failure or a serialize/re-parse + /// round-trip; an unsupported column type logs a warning, ticks the + /// drop counter, and yields an empty map. Actual JSON shape / + /// per-entry validation is shared via + /// `parse_startup_parameters_value`. fn extract_startup_parameters( row: &tokio_postgres::Row, username: &str, @@ -403,30 +447,66 @@ impl AuthQueryExecutor { let Some(column) = column else { return std::collections::HashMap::new(); }; - let raw: Option = match row.try_get::<_, Option>("startup_parameters") { - Ok(v) => v, - Err(e) => { - warn!( - "[{username}@{pool_name}] auth_query startup_parameters column has type \ - `{ty}` but pg_doorman reads it as `text`: {e}. If the SELECT returns \ - json or jsonb, add `::text` (for example: \ - `jsonb_build_object(...)::text AS startup_parameters`); per-user \ - parameters are ignored for this row.", - ty = column.type_().name() - ); - crate::web::metrics::STARTUP_PARAMETERS_DROPPED_TOTAL - .with_label_values(&[pool_name, "auth_query_bad_type"]) - .inc(); - return std::collections::HashMap::new(); + let col_type = column.type_(); + if matches!(*col_type, Type::JSON | Type::JSONB) { + // The custom `LimitedJson` wrapper enforces the + // `MAX_OPERATOR_BUDGET` cap on the raw wire bytes from PG + // before `serde_json` walks the tree. Without this gate the + // text path (which checks `text.len()` before + // `from_str`) and the json/jsonb path diverge: a malicious + // or accidentally oversize `jsonb` row would still force + // pg_doorman to materialise the full `Value` tree before + // discarding it. + match row.try_get::<_, Option>("startup_parameters") { + Ok(Some(LimitedJson(value))) => { + Self::parse_startup_parameters_value(value, username, pool_name) + } + Ok(None) => std::collections::HashMap::new(), + Err(json_err) => { + if json_err.to_string().contains(LIMITED_JSON_OVERSIZE_TAG) { + warn!( + "[{username}@{pool_name}] auth_query startup_parameters: {json_err}; \ + parameters ignored" + ); + crate::web::metrics::STARTUP_PARAMETERS_DROPPED_TOTAL + .with_label_values(&[pool_name, "auth_query_oversize"]) + .inc(); + } else { + warn!( + "[{username}@{pool_name}] auth_query startup_parameters column has type \ + `{ty}` but pg_doorman could not decode it as json: {json_err}. \ + Per-user parameters are ignored for this row.", + ty = col_type.name() + ); + crate::web::metrics::STARTUP_PARAMETERS_DROPPED_TOTAL + .with_label_values(&[pool_name, "auth_query_bad_type"]) + .inc(); + } + std::collections::HashMap::new() + } } - }; - Self::parse_startup_parameters_text(raw.as_deref(), username, pool_name) + } else { + match row.try_get::<_, Option>("startup_parameters") { + Ok(raw) => Self::parse_startup_parameters_text(raw.as_deref(), username, pool_name), + Err(text_err) => { + warn!( + "[{username}@{pool_name}] auth_query startup_parameters column has type \ + `{ty}` but pg_doorman reads it as `text`, `json`, or `jsonb`: {text_err}. \ + Per-user parameters are ignored for this row.", + ty = col_type.name() + ); + crate::web::metrics::STARTUP_PARAMETERS_DROPPED_TOTAL + .with_label_values(&[pool_name, "auth_query_bad_type"]) + .inc(); + std::collections::HashMap::new() + } + } + } } - /// Parse the optional `startup_parameters` JSON object returned by - /// auth_query. Valid string entries become per-user GUCs. Invalid keys, - /// non-string values, malformed JSON, and non-object JSON are logged and - /// ignored; authentication still continues. + /// Parse the optional `startup_parameters` JSON object received as + /// `text` from auth_query. Wraps `parse_startup_parameters_value` + /// after the size gate and `serde_json::from_str`. fn parse_startup_parameters_text( text: Option<&str>, username: &str, @@ -466,7 +546,19 @@ impl AuthQueryExecutor { return std::collections::HashMap::new(); } }; - let serde_json::Value::Object(obj) = parsed else { + Self::parse_startup_parameters_value(parsed, username, pool_name) + } + + /// Per-entry validation shared between the `text` and `json`/`jsonb` + /// auth_query decoders. The json/jsonb path used to serialise the + /// `serde_json::Value` back into a string and re-parse it here; this + /// helper avoids the round-trip. + fn parse_startup_parameters_value( + value: serde_json::Value, + username: &str, + pool_name: &str, + ) -> std::collections::HashMap { + let serde_json::Value::Object(obj) = value else { warn!( "[{username}@{pool_name}] auth_query startup_parameters: top-level value is not a \ JSON object; ignored" @@ -489,7 +581,13 @@ impl AuthQueryExecutor { had_invalid_entry = true; continue; } - out.insert(k, s); + // Canonicalise tracked GUC names so the per-user + // overlay merges with the general/pool cascade by + // canonical key. Without this, an auth_query row + // returning `timezone` would not override a pool + // `TimeZone` baseline and both would survive. + let canonical = crate::server::parameters::canonicalize_param_name(k); + out.insert(canonical, s); } other => { let kind = match other { @@ -517,6 +615,28 @@ impl AuthQueryExecutor { .with_label_values(&[pool_name, "auth_query_invalid_entry"]) .inc(); } + // The text decoder enforces `MAX_OPERATOR_BUDGET` on the raw + // column before parsing. The json/jsonb decoder has no raw + // bytes to size, so apply the same gate to the wire-shape of + // the resulting map. Without this a large `jsonb` row would be + // cached and then rejected at backend startup with a less + // helpful 53400; here we drop it as `auth_query_oversize` like + // the text path. + let max_bytes = crate::config::startup_parameters::MAX_OPERATOR_BUDGET; + let serialized = out + .iter() + .map(|(k, v)| k.len() + 1 + v.len() + 1) + .sum::(); + if serialized > max_bytes { + warn!( + "[{username}@{pool_name}] auth_query startup_parameters: decoded map is \ + {serialized} bytes, exceeding operator budget {max_bytes}; parameters ignored" + ); + crate::web::metrics::STARTUP_PARAMETERS_DROPPED_TOTAL + .with_label_values(&[pool_name, "auth_query_oversize"]) + .inc(); + return std::collections::HashMap::new(); + } out } } @@ -525,6 +645,47 @@ impl AuthQueryExecutor { // CacheEntry // --------------------------------------------------------------------------- +/// Immutable snapshot of a user's auth_query overlay: the wire map plus +/// the precomputed hash the dynamic-pool fast path consumes. The two +/// values move together by contract — direct field writes risk drift +/// (e.g. swapping the map without touching the hash), so the struct +/// only exposes accessors and constructors. Cheap to clone because the +/// inner `Arc` is shared. +#[derive(Clone, Debug)] +pub struct StartupOverlay { + map: Arc>, + hash: u64, +} + +impl StartupOverlay { + pub fn empty() -> Self { + Self { + map: Arc::new(std::collections::HashMap::new()), + hash: crate::pool::empty_overlay_hash(), + } + } + + pub fn from_map(map: std::collections::HashMap) -> Self { + let hash = crate::pool::per_user_overlay_hash(map.iter()); + Self { + map: Arc::new(map), + hash, + } + } + + pub fn map(&self) -> &Arc> { + &self.map + } + + pub fn hash(&self) -> u64 { + self.hash + } + + pub fn is_empty(&self) -> bool { + self.map.is_empty() + } +} + /// Single cache entry for a username's credentials. #[derive(Clone, Debug)] pub struct CacheEntry { @@ -542,14 +703,12 @@ pub struct CacheEntry { /// None for MD5 users or before first SCRAM auth. pub client_key: Option>, /// Per-user startup parameters returned by the optional auth_query - /// `startup_parameters` JSON column. Empty when the column is absent, - /// empty/NULL, or filtered out in dedicated auth_query mode. - /// - /// Wrapped in `Arc` so cache hits do not clone the underlying map. - /// For a user with a wide row (a dozen extension GUCs) every - /// `cache.get_or_fetch` previously paid an `O(map)` clone; the - /// `Arc::clone` here is two atomic increments instead. - pub startup_parameters: Arc>, + /// `startup_parameters` JSON column, paired with their overlay hash. + /// The map is empty (and hash is `empty_overlay_hash()`) when the + /// column is absent, empty/NULL, or filtered out in dedicated + /// auth_query mode. `StartupOverlay` keeps map and hash from + /// drifting — mutate via [`Self::set_startup_overlay`]. + pub startup_overlay: StartupOverlay, } impl CacheEntry { @@ -560,7 +719,7 @@ impl CacheEntry { is_negative: false, last_refetch_at: None, client_key: None, - startup_parameters: Arc::new(std::collections::HashMap::new()), + startup_overlay: StartupOverlay::empty(), } } @@ -571,10 +730,19 @@ impl CacheEntry { is_negative: true, last_refetch_at: None, client_key: None, - startup_parameters: Arc::new(std::collections::HashMap::new()), + startup_overlay: StartupOverlay::empty(), } } + /// Replace the overlay with one built from the given map; the new + /// hash is computed inside the constructor. + pub fn set_startup_overlay( + &mut self, + startup_parameters: std::collections::HashMap, + ) { + self.startup_overlay = StartupOverlay::from_map(startup_parameters); + } + fn is_expired(&self, cache_ttl: &Duration, cache_failure_ttl: &Duration) -> bool { let ttl_ms = if self.is_negative { cache_failure_ttl.as_millis() @@ -655,7 +823,7 @@ impl AuthQueryCache { /// apply them. Drop the parsed map before it reaches downstream code /// and warn once per (pool, username) so the operator notices. fn dedicated_mode_filter(&self, entry: &mut CacheEntry, username: &str) { - if !self.is_dedicated || entry.startup_parameters.is_empty() { + if !self.is_dedicated || entry.startup_overlay.is_empty() { return; } // One increment per drop event (a single fetched row whose @@ -677,7 +845,7 @@ impl AuthQueryCache { pool = self.pool_name ); } - entry.startup_parameters = Arc::new(std::collections::HashMap::new()); + entry.set_startup_overlay(std::collections::HashMap::new()); } /// When a fresh auth_query fetch produces a per-user @@ -779,7 +947,7 @@ impl AuthQueryCache { Ok(Some((password_hash, startup_params))) => { self.inc(|s| &s.cache_misses); let mut entry = CacheEntry::positive(password_hash); - entry.startup_parameters = Arc::new(startup_params); + entry.set_startup_overlay(startup_params); self.dedicated_mode_filter(&mut entry, username); // Publish the fresh entry first so any concurrent // create_dynamic_pool peeks the new overlay, then drop the @@ -789,7 +957,7 @@ impl AuthQueryCache { // create_dynamic_pool would rebuild against that stale // map and immediately drift again. self.entries.insert(username.to_string(), entry.clone()); - self.drop_dynamic_pool_if_overlay_drifted(username, &entry.startup_parameters); + self.drop_dynamic_pool_if_overlay_drifted(username, entry.startup_overlay.map()); Ok(Some(entry)) } Ok(None) => { @@ -854,12 +1022,12 @@ impl AuthQueryCache { match self.executor.fetch_credentials(username).await { Ok(Some((password_hash, startup_params))) => { let mut entry = CacheEntry::positive(password_hash); - entry.startup_parameters = Arc::new(startup_params); + entry.set_startup_overlay(startup_params); entry.last_refetch_at = Some(Instant::now()); self.dedicated_mode_filter(&mut entry, username); // Insert before drop — see comment in get_or_fetch. self.entries.insert(username.to_string(), entry.clone()); - self.drop_dynamic_pool_if_overlay_drifted(username, &entry.startup_parameters); + self.drop_dynamic_pool_if_overlay_drifted(username, entry.startup_overlay.map()); Ok(Some(entry)) } Ok(None) => { @@ -919,7 +1087,7 @@ impl AuthQueryCache { if entry.is_expired(&self.cache_ttl, &self.cache_failure_ttl) { return None; } - Some(f(&entry.startup_parameters)) + Some(f(entry.startup_overlay.map())) } /// Number of cached entries (for metrics/admin). @@ -1376,11 +1544,11 @@ mod tests { #[test] fn parse_startup_parameters_oversize_text_returns_empty() { - // HIGH #9 regression guard: pathological auth_query row should not - // make serde_json walk megabytes of JSON. The raw text cap matches - // `MAX_OPERATOR_BUDGET`, so anything past that returns empty before - // we even start parsing. Drop the same value into a giant string - // so the byte length crosses the cap independently of JSON shape. + // A pathological auth_query row should not make serde_json walk + // megabytes of JSON. The raw text cap matches `MAX_OPERATOR_BUDGET`, + // so anything past that returns empty before parsing starts. Use a + // giant string so the byte length crosses the cap independently of + // JSON shape. let cap = crate::config::startup_parameters::MAX_OPERATOR_BUDGET; let bytes = "a".repeat(cap + 1); let r = AuthQueryExecutor::parse_startup_parameters_text(Some(&bytes), "u", "p"); @@ -1432,7 +1600,7 @@ mod tests { // because the backend identity is shared in dedicated mode. let entry = cache.get_or_fetch("alice").await.unwrap().unwrap(); assert!( - entry.startup_parameters.is_empty(), + entry.startup_overlay.is_empty(), "params must be cleared in dedicated mode" ); @@ -1441,7 +1609,7 @@ mod tests { // tracker still holds exactly one entry after a second miss-and-fill. cache.invalidate("alice"); let entry = cache.get_or_fetch("alice").await.unwrap().unwrap(); - assert!(entry.startup_parameters.is_empty()); + assert!(entry.startup_overlay.is_empty()); assert_eq!(cache.dedicated_warnings.len(), 1); // clear() resets the warning tracker so a config reload re-arms it. @@ -1458,7 +1626,11 @@ mod tests { let cache = make_cache(fetcher.clone(), &config); let entry = cache.get_or_fetch("alice").await.unwrap().unwrap(); assert_eq!( - entry.startup_parameters.get("work_mem").map(String::as_str), + entry + .startup_overlay + .map() + .get("work_mem") + .map(String::as_str), Some("64MB") ); } @@ -1501,11 +1673,10 @@ mod tests { #[tokio::test] async fn peek_startup_parameters_returns_none_for_expired_entry() { - // HIGH #7 regression guard: a positive cache entry that has lived - // past `cache_ttl` must not pin a stale per-user startup parameter - // onto a backend the replenishment loop spawns later. Mirrors - // `test_cache_ttl_expiration` but exercises the peek path the - // backend-spawn hot path uses. + // A positive cache entry that has lived past `cache_ttl` must not + // pin a stale per-user startup parameter onto a backend spawned by + // the replenishment loop. Mirrors `test_cache_ttl_expiration` but + // exercises the peek path used by backend startup. let fetcher = Arc::new(MockFetcher::new()); fetcher.add_user_with_params("alice", "md5abc123", &[("work_mem", "64MB")]); let mut config = test_config(); diff --git a/src/auth/mod.rs b/src/auth/mod.rs index 8bd6fd9c7..0e59c6e6c 100644 --- a/src/auth/mod.rs +++ b/src/auth/mod.rs @@ -41,6 +41,28 @@ use crate::pool::{ }; use crate::server::ServerParameters; +/// Canonicalised set of GUC names the operator put under +/// `general.startup_parameters` / `pool.startup_parameters` / +/// `auth_query` for the (db, user) pair that authenticated this +/// connection. The same `Arc` lives on `ConnectionPool` so cloning +/// stays zero-copy. Client startup uses this to drop `ParameterStatus` +/// entries the client sent for keys the backend session already has +/// pinned by `startup_parameters`. +pub type OperatorManagedKeys = Arc>; + +/// Outcome of [`authenticate`]: everything the client startup path needs +/// to finish the StartupMessage exchange. `operator_managed_keys` is +/// captured from the same `ConnectionPool` snapshot that produced +/// `server_parameters`, so the client startup filter cannot drift +/// against a concurrent RELOAD or auth_query overlay refetch the way a +/// second global `POOLS` lookup would. +pub struct AuthOutcome { + pub transaction_mode: bool, + pub server_parameters: ServerParameters, + pub prepared_statements_enabled: bool, + pub operator_managed_keys: Option, +} + /// Authenticate a user based on the provided parameters pub async fn authenticate( read: &mut S, @@ -49,7 +71,7 @@ pub async fn authenticate( client_identifier: &mut ClientIdentifier, pool_name: &str, username_from_parameters: &str, -) -> Result<(bool, ServerParameters, bool), Error> +) -> Result where S: AsyncReadExt + Unpin, T: AsyncWriteExt + Unpin, @@ -57,7 +79,7 @@ where let mut prepared_statements_enabled = false; // Authenticate admin user. - let (transaction_mode, server_parameters) = if admin { + let (transaction_mode, server_parameters, operator_managed_keys) = if admin { if client_identifier.hba_md5 == CheckResult::Trust || client_identifier.hba_scram == CheckResult::Trust { @@ -65,7 +87,12 @@ where "HBA trust: admin user={username_from_parameters}, addr={}", client_identifier.addr ); - return Ok((false, ServerParameters::admin(), false)); + return Ok(AuthOutcome { + transaction_mode: false, + server_parameters: ServerParameters::admin(), + prepared_statements_enabled: false, + operator_managed_keys: None, + }); } if client_identifier.hba_md5 == CheckResult::Deny || client_identifier.hba_scram == CheckResult::Deny @@ -77,7 +104,8 @@ where wrong_password(write, username_from_parameters).await?; return Err(error); } - authenticate_admin(read, write, username_from_parameters).await? + let (tx, sp) = authenticate_admin(read, write, username_from_parameters).await?; + (tx, sp, None) } // Authenticate normal user. else { @@ -92,11 +120,12 @@ where .await? }; - Ok(( + Ok(AuthOutcome { transaction_mode, server_parameters, prepared_statements_enabled, - )) + operator_managed_keys, + }) } /// Authenticate an admin user with MD5 @@ -191,7 +220,7 @@ async fn authenticate_normal_user( pool_name: &str, username_from_parameters: &str, prepared_statements_enabled: &mut bool, -) -> Result<(bool, ServerParameters), Error> +) -> Result<(bool, ServerParameters, Option), Error> where S: AsyncReadExt + Unpin, T: AsyncWriteExt + Unpin, @@ -355,7 +384,17 @@ where } }; - Ok((transaction_mode, server_parameters)) + // Capture operator-managed startup-parameter keys from the same + // pool snapshot that produced `server_parameters`. The client + // startup path used to read this set with a second `POOLS` global + // lookup, which could observe a RELOAD between authentication and + // the lookup and send `ParameterStatus` values for keys the + // backend session already has set via operator-managed + // `startup_parameters`. Snapshotting on this side guarantees the + // two views stay in step. + let operator_managed_keys = Some(pool.database.server_pool().operator_managed_startup_keys()); + + Ok((transaction_mode, server_parameters, operator_managed_keys)) } /// Authenticate a user with PAM @@ -633,7 +672,7 @@ async fn try_auth_query( pool_name: &str, username: &str, prepared_statements_enabled: &mut bool, -) -> Result<(bool, ServerParameters), Error> +) -> Result<(bool, ServerParameters, Option), Error> where S: AsyncReadExt + Unpin, T: AsyncWriteExt + Unpin, @@ -681,8 +720,11 @@ where } }; - // 3. Fetch password hash from cache or PG - let cache_entry = match cache.get_or_fetch(username).await { + // 3. Fetch password hash from cache or PG. `mut` because a + // successful MD5 refetch below swaps in the fresh entry so the + // backend pool gets the rotated password hash and the rotated + // per-user startup_parameters, not the stale snapshot. + let mut cache_entry = match cache.get_or_fetch(username).await { Ok(Some(entry)) => entry, Ok(None) => { // User not found @@ -705,10 +747,8 @@ where } }; - let pool_password = &cache_entry.password_hash; - // 4. HBA check - let hba_decision = eval_hba_for_pool_password(pool_password, client_identifier); + let hba_decision = eval_hba_for_pool_password(&cache_entry.password_hash, client_identifier); if hba_decision == CheckResult::Deny { error_response_terminal( write, @@ -730,26 +770,48 @@ where if hba_decision == CheckResult::Trust { // HBA trust — skip password check - } else if pool_password.starts_with(MD5_PASSWORD_PREFIX) { + } else if cache_entry.password_hash.starts_with(MD5_PASSWORD_PREFIX) { // MD5 challenge-response let salt = md5_challenge(write).await?; let password_response = read_password(read).await?; - let expected = md5_hash_second_pass(pool_password.strip_prefix("md5").unwrap(), &salt); + let expected = md5_hash_second_pass( + cache_entry.password_hash.strip_prefix("md5").unwrap(), + &salt, + ); if expected != password_response { // Password mismatch — try re-fetch (password may have changed in PG) let mut auth_ok = false; + let mut refreshed: Option = None; if let Ok(Some(new_entry)) = cache.refetch_on_failure(username).await { - if new_entry.password_hash != *pool_password - && new_entry.password_hash.starts_with(MD5_PASSWORD_PREFIX) - { - let new_expected = md5_hash_second_pass( - new_entry.password_hash.strip_prefix("md5").unwrap(), - &salt, - ); - if new_expected == password_response { - auth_ok = true; - info!("[{username}@{pool_name}] auth_query: re-fetched password matched"); + if new_entry.password_hash != cache_entry.password_hash { + if new_entry.password_hash.starts_with(MD5_PASSWORD_PREFIX) { + let new_expected = md5_hash_second_pass( + new_entry.password_hash.strip_prefix("md5").unwrap(), + &salt, + ); + if new_expected == password_response { + auth_ok = true; + info!( + "[{username}@{pool_name}] auth_query: re-fetched password matched" + ); + refreshed = Some(new_entry); + } + } else { + // The refetched verifier is no longer MD5 — the + // operator switched `password_encryption` mid-flight + // (typically MD5 → SCRAM). The current MD5 proof + // cannot validate against a SCRAM verifier; reject + // this attempt and invalidate the cache so the next + // reconnect hits `cache.get_or_fetch` and takes the + // SCRAM branch immediately rather than waiting for + // `cache_ttl`. + warn!( + "[{username}@{pool_name}] auth_query: refetched verifier changed type ({stored} → {fresh}); cache invalidated, client must reconnect with the new mechanism", + stored = if cache_entry.password_hash.starts_with(MD5_PASSWORD_PREFIX) { "md5" } else { "scram" }, + fresh = if new_entry.password_hash.starts_with(MD5_PASSWORD_PREFIX) { "md5" } else { "scram" } + ); + cache.invalidate(username); } } } @@ -763,10 +825,17 @@ where "MD5 authentication failed for auth_query user: {username}" ))); } + // Swap in the refetched snapshot so backend_auth and the + // dynamic-pool overlay below are built from the rotated + // credentials, not the stale ones that just failed the + // first challenge. + if let Some(new_entry) = refreshed { + cache_entry = new_entry; + } } - } else if pool_password.starts_with(SCRAM_SHA_256) { + } else if cache_entry.password_hash.starts_with(SCRAM_SHA_256) { // SCRAM-SHA-256 challenge-response - let server_secret = match parse_server_secret(pool_password) { + let server_secret = match parse_server_secret(&cache_entry.password_hash) { Ok(s) => s, Err(err) => { error!( @@ -927,12 +996,20 @@ where shared_pool_id ); - Ok((transaction_mode, server_parameters)) + let operator_managed_keys = + Some(pool.database.server_pool().operator_managed_startup_keys()); + Ok((transaction_mode, server_parameters, operator_managed_keys)) } None => { // === Passthrough mode: each dynamic user gets their own pool === - let backend_auth = if pool_password.starts_with(MD5_PASSWORD_PREFIX) { - Some(BackendAuthMethod::Md5PassTheHash(pool_password.clone())) + // After an MD5 refetch matched the rotated password, + // `cache_entry` already points at the new snapshot, so + // `password_hash` and `startup_parameters` below reflect the + // credentials PG will accept on the backend side. + let backend_auth = if cache_entry.password_hash.starts_with(MD5_PASSWORD_PREFIX) { + Some(BackendAuthMethod::Md5PassTheHash( + cache_entry.password_hash.clone(), + )) } else { auth_client_key.map(BackendAuthMethod::ScramPassthrough) }; @@ -941,14 +1018,19 @@ where // authenticated this user. That keeps dynamic-pool creation // tied to this login instead of reading the global cache // again while TTL expiry or a concurrent refetch is changing it. - let fetched_overlay = Arc::clone(&cache_entry.startup_parameters); - let mut pool = create_dynamic_pool(pool_name, username, backend_auth, fetched_overlay) - .map_err(|err| { - error!( - "[{username}@{pool_name}] auth_query: failed to create dynamic pool: {err}" - ); - err - })?; + let fetched_overlay = Arc::clone(cache_entry.startup_overlay.map()); + let fetched_overlay_hash = cache_entry.startup_overlay.hash(); + let mut pool = create_dynamic_pool( + pool_name, + username, + backend_auth, + fetched_overlay, + fetched_overlay_hash, + ) + .map_err(|err| { + error!("[{username}@{pool_name}] auth_query: failed to create dynamic pool: {err}"); + err + })?; // Do NOT change client_identifier.username — stay as the dynamic user // so that Client.username matches the pool's user for get_pool() lookups. @@ -967,6 +1049,17 @@ where } = &err { error!("[{username}@{pool_name}] auth_query passthrough: PG rejected operator-supplied startup parameter: {pg_message}"); + // Invalidate before dropping the pool. A + // concurrent reconnect that races us between + // these two calls would otherwise read the + // still-cached bad overlay, rebuild the same + // dynamic pool against it, and trigger the + // same rejection. Invalidating first guarantees + // any racing get_or_fetch refetches before + // reconstructing the pool. + cache.invalidate(username); + let identifier = crate::pool::PoolIdentifier::new(pool_name, username); + crate::pool::drop_dynamic_pool(&identifier); error_response(write, pg_message, sqlstate).await?; return Err(err); } @@ -984,7 +1077,9 @@ where aq_state.stats.auth_success.fetch_add(1, Ordering::Relaxed); info!("[{username}@{pool_name}] auth_query: authenticated (passthrough mode)"); - Ok((transaction_mode, server_parameters)) + let operator_managed_keys = + Some(pool.database.server_pool().operator_managed_startup_keys()); + Ok((transaction_mode, server_parameters, operator_managed_keys)) } } } diff --git a/src/client/startup.rs b/src/client/startup.rs index 1f5317a36..562963ffc 100644 --- a/src/client/startup.rs +++ b/src/client/startup.rs @@ -367,7 +367,7 @@ where let secret_key: i32 = rand::random(); // Authenticate user - let (transaction_mode, mut server_parameters, prepared_statements_enabled) = authenticate( + let auth_outcome = authenticate( &mut read, &mut write, admin, @@ -376,9 +376,34 @@ where username_from_parameters, ) .await?; - - // Update the parameters to merge what the application sent and what's originally on the server - server_parameters.set_from_hashmap(¶meters, false); + let transaction_mode = auth_outcome.transaction_mode; + let mut server_parameters = auth_outcome.server_parameters; + let prepared_statements_enabled = auth_outcome.prepared_statements_enabled; + + // Merge the startup parameters sent by the client with the + // server defaults. Operator-managed startup_parameters must + // win over the client packet; otherwise ParameterStatus + // reports the client value while the backend keeps the + // operator default, a protocol-visible mismatch. The + // backend's sync_parameters already applies the same filter + // on checkout, so this aligns the two client views. The key + // set comes from the same pool snapshot that produced + // `server_parameters`, so a concurrent RELOAD or auth_query + // overlay refetch between authentication and this filter + // cannot make the two views diverge. + match auth_outcome.operator_managed_keys { + Some(keys) if !keys.is_empty() => { + for (key, value) in ¶meters { + let canonical = crate::server::parameters::canonicalize_param_name(key.clone()); + if !keys.contains(&canonical) { + server_parameters.set_param(key.clone(), value.clone(), false); + } + } + } + _ => { + server_parameters.set_from_hashmap(¶meters, false); + } + } let mut buf = BytesMut::new(); { let mut auth_ok = BytesMut::with_capacity(9); diff --git a/src/config/mod.rs b/src/config/mod.rs index 567e49cae..d3f3e0911 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -428,6 +428,85 @@ impl Config { &self.general.startup_parameters, "general.startup_parameters", )?; + // Reject deterministic `general + pool` overflows at config load. + // For each configured user, mirror the runtime full-packet size + // check so `pg_doorman -t` fails even when the parameter body fits + // but `user`/`database`/`application_name` would push the full + // StartupMessage over `MAX_STARTUP_PACKET_LENGTH`. The checks + // here only cover size: reserved-key and shape validation has + // already run per level, and auth_query overlays are still + // checked at backend startup because they come from PostgreSQL. + for (pool_name, pool_config) in &self.pools { + // Same canonical cascade build the runtime does in + // `ServerPool::new`. Without the canonicalisation here, a + // pool that overrides `timezone` with `TimeZone` would + // serialise two rows during validation and disagree with + // the runtime byte count. + let merged = startup_parameters::cascade_canonical_keys(&[ + &self.general.startup_parameters, + &pool_config.startup_parameters, + ]); + let merged_size = startup_parameters::serialized_bytes(&merged); + if merged_size > startup_parameters::MAX_OPERATOR_BUDGET { + return Err(Error::BadConfig(format!( + "merged general + pools.{pool_name}.startup_parameters: serialized \ + size {merged_size} bytes exceeds operator budget {} (PG \ + StartupMessage cap is {} bytes; reduce general or pool startup_parameters)", + startup_parameters::MAX_OPERATOR_BUDGET, + startup_parameters::MAX_STARTUP_PACKET_SIZE, + ))); + } + let server_database = pool_config + .server_database + .as_deref() + .unwrap_or(pool_name.as_str()); + // Runtime resolves the StartupMessage application_name as + // pool override → `"pg_doorman"`. Mirror that default so + // `pg_doorman -t` doesn't accept a config whose only safe + // case is the empty-string assumption. + let application_name = pool_config + .application_name + .as_deref() + .unwrap_or("pg_doorman"); + let validate_user_identity = |display_kind: &str, + display_user: &str, + server_username: &str| + -> Result<(), Error> { + let (packet_bytes, _body_bytes) = startup_parameters::packet_and_body_bytes( + server_username, + server_database, + application_name, + &merged, + ); + if packet_bytes > startup_parameters::MAX_STARTUP_PACKET_SIZE { + return Err(Error::BadConfig(format!( + "merged general + pools.{pool_name}.startup_parameters: full StartupMessage \ + for {display_kind} '{display_user}' is {packet_bytes} bytes, exceeding \ + the PG cap of {} bytes (user/database/application_name overhead \ + included); reduce general or pool startup_parameters", + startup_parameters::MAX_STARTUP_PACKET_SIZE, + ))); + } + Ok(()) + }; + for user in &pool_config.users { + let server_username = user + .server_username + .as_deref() + .unwrap_or(user.username.as_str()); + validate_user_identity("user", &user.username, server_username)?; + } + // Dedicated auth_query mode opens one shared backend + // connection identified by `auth_query.server_user`; that + // identity must fit the packet just like a static user. + // Use a distinct display kind so operators don't waste time + // hunting for the name in `pool_config.users`. + if let Some(aq) = pool_config.auth_query.as_ref() { + if let Some(shared_user) = aq.server_user.as_deref() { + validate_user_identity("auth_query server_user", shared_user, shared_user)?; + } + } + } if self.general.tls_rate_limit_per_second < 100 && self.general.tls_rate_limit_per_second != 0 diff --git a/src/config/startup_parameters.rs b/src/config/startup_parameters.rs index f8ea832fc..73b5c1952 100644 --- a/src/config/startup_parameters.rs +++ b/src/config/startup_parameters.rs @@ -65,6 +65,19 @@ pub fn validate(map: &BTreeMap, scope: &str) -> Result<(), Error validate_key(k, scope)?; validate_value(k, v, scope)?; } + // PG GUC names are case-insensitive. Two keys that collapse to the + // same canonical form would silently lose one entry to BTreeMap + // iteration order in the cascade merge — reject the configuration + // instead and tell the operator which two keys clash. + let mut seen: std::collections::HashMap = std::collections::HashMap::new(); + for k in map.keys() { + let canonical = crate::server::parameters::canonicalize_param_name(k.clone()); + if let Some(existing) = seen.insert(canonical.clone(), k.clone()) { + return Err(Error::BadConfig(format!( + "{scope}: '{existing}' and '{k}' both refer to PG GUC '{canonical}'; keep one" + ))); + } + } validate_total_size(map, scope) } @@ -88,7 +101,13 @@ fn validate_key(key: &str, scope: &str) -> Result<(), Error> { "{scope}: '{key}' is reserved and managed by pg_doorman" ))); } - if key.starts_with(RESERVED_PREFIX) { + // Case-insensitive prefix check: PG GUC lookup is case-insensitive + // and `canonicalize_param_name` lowercases non-tracked keys before + // they reach the wire, so `_PQ_.foo` would otherwise pass validation + // here and emerge on the wire as `_pq_.foo`. + if key.len() >= RESERVED_PREFIX.len() + && key.as_bytes()[..RESERVED_PREFIX.len()].eq_ignore_ascii_case(RESERVED_PREFIX.as_bytes()) + { return Err(Error::BadConfig(format!( "{scope}: '{key}' uses the reserved '_pq_.' prefix" ))); @@ -129,6 +148,25 @@ pub fn serialized_bytes(map: &BTreeMap) -> usize { map.iter().map(|(k, v)| k.len() + 1 + v.len() + 1).sum() } +/// Merge `general`/`pool`/`auth_query` startup_parameter layers with +/// PostgreSQL case-insensitive GUC semantics. Each layer's keys are +/// canonicalised before insertion, so a pool `TimeZone` correctly wins +/// over a general `timezone` regardless of the raw casing the operator +/// wrote. Layers later in the slice override earlier ones, mirroring +/// the cascade order documented in the tutorial. +pub fn cascade_canonical_keys(layers: &[&BTreeMap]) -> BTreeMap { + let mut out = BTreeMap::new(); + for layer in layers { + for (k, v) in layer.iter() { + out.insert( + crate::server::parameters::canonicalize_param_name(k.clone()), + v.clone(), + ); + } + } + out +} + /// Exact byte length of the full StartupMessage pg_doorman will put on the /// wire for one backend spawn, *including* the 4-byte length prefix. The /// layout mirrors `crate::messages::protocol::startup`: @@ -275,6 +313,30 @@ mod tests { assert!(matches!(err, Error::BadConfig(_))); } + #[test] + fn pq_prefix_rejected_case_insensitive() { + // PG GUC lookup is case-insensitive, and `canonicalize_param_name` + // lowercases non-tracked keys before they reach the wire. Without + // a case-insensitive check the reserved protocol-extension + // namespace can be smuggled in as `_PQ_.foo` or `_Pq_.foo` and + // emerges on the wire as `_pq_.foo`. + for key in ["_PQ_.fancy_ext", "_Pq_.fancy_ext", "_pQ_.FANCY"] { + let err = validate(&m(&[(key, "x")]), "scope").unwrap_err(); + assert!( + matches!(err, Error::BadConfig(ref msg) if msg.contains("reserved")), + "key {key} must be rejected as reserved" + ); + } + } + + #[test] + fn validate_entry_rejects_pq_prefix_case_insensitive() { + // auth_query JSON values flow through `validate_entry`, so the + // reserved-prefix guard must cover that path identically. + let err = validate_entry("_PQ_.fancy_ext", "x", "scope").unwrap_err(); + assert!(matches!(err, Error::BadConfig(ref msg) if msg.contains("reserved"))); + } + #[test] fn empty_key_rejected() { let err = validate(&m(&[("", "x")]), "scope").unwrap_err(); diff --git a/src/config/tests.rs b/src/config/tests.rs index 8924cffbc..331f6217a 100644 --- a/src/config/tests.rs +++ b/src/config/tests.rs @@ -1687,9 +1687,9 @@ fn check_hba_legacy_empty_allows_unix() { #[test] fn check_hba_legacy_list_bypassed_for_unix() { - // Reproduces the "silent privilege expansion" case from review: the - // operator restricts TCP access with a CIDR whitelist, but Unix clients - // must still be allowed because the legacy list has no transport concept. + // The legacy CIDR allowlist applies only to TCP clients. Unix socket + // clients must still be allowed because the legacy list has no + // transport concept. let mut general = General::default(); general.hba = vec!["10.0.0.0/8".parse().unwrap()]; @@ -2074,3 +2074,35 @@ async fn reject_reserved_in_pool_startup_parameters() { other => panic!("expected BadConfig, got {other:?}"), } } + +#[tokio::test] +async fn reject_merged_general_pool_startup_parameters_overflow() { + // Two layers can fit on their own and still overflow once merged. + // Config validation should catch that at `pg_doorman -t`, before the + // first client tries to connect. + let mut cfg = Config::default(); + cfg.general.tls_rate_limit_per_second = 0; + let filler = "x".repeat(4800); + cfg.general + .startup_parameters + .insert("aaa_big".to_string(), filler.clone()); + let mut pool = Pool::default(); + pool.startup_parameters + .insert("bbb_big".to_string(), filler); + pool.users.push(User { + username: "u".to_string(), + password: "p".to_string(), + pool_size: 1, + ..User::default() + }); + cfg.pools.insert("p".to_string(), pool); + let err = cfg.validate().await.unwrap_err(); + match err { + Error::BadConfig(msg) => assert!( + msg.contains("merged general + pools.p.startup_parameters") + && msg.contains("exceeds operator budget"), + "unexpected message: {msg}" + ), + other => panic!("expected BadConfig, got {other:?}"), + } +} diff --git a/src/pool/dynamic.rs b/src/pool/dynamic.rs index 1d95dfc23..2131ae0da 100644 --- a/src/pool/dynamic.rs +++ b/src/pool/dynamic.rs @@ -36,17 +36,45 @@ pub fn create_dynamic_pool( username: &str, backend_auth: Option, fetched_overlay: Arc>, + fetched_overlay_hash: u64, ) -> Result { - // Fast path: pool already exists + // Fast path: pool already exists. The cache-side refetch path + // already drops the live pool when an auth_query refetch changes + // the overlay (see `drop_dynamic_pool_if_overlay_drifted`), but a + // concurrent login can still arrive after the cache published the + // fresh entry yet before the drop fires, or with a fetched_overlay + // newer than what the live pool was frozen with. Check the overlay + // hash here too so that login rebuilds the pool against the + // current snapshot instead of inheriting a stale one. The hash is + // precomputed on `CacheEntry`, so the fast path skips the sort + + // SipHash on every login. if let Some(existing) = get_pool(pool_name, username) { - // Update backend_auth (credentials may have changed) - if let (Some(ref ba_lock), Some(new_ba)) = (&existing.address.backend_auth, &backend_auth) { - debug!( - "[{username}@{pool_name}] auth_query: dynamic pool already exists, updating backend_auth" + let identifier = super::PoolIdentifier::new(pool_name, username); + let live_hash = existing.per_user_startup_overlay_hash; + let is_dyn = super::is_dynamic_pool(&identifier); + if !should_rebuild_for_overlay_drift(live_hash, fetched_overlay_hash, is_dyn) { + // Hash matches, or the live pool is static and the empty + // baseline does not match an auth_query overlay — either + // way the existing pool wins. Refresh `backend_auth` only + // on hash match: a password rotation between cache + // fetches still applies, but a static pool is left alone. + if live_hash == fetched_overlay_hash { + if let (Some(ref ba_lock), Some(new_ba)) = + (&existing.address.backend_auth, &backend_auth) + { + debug!( + "[{username}@{pool_name}] auth_query: dynamic pool already exists, updating backend_auth" + ); + *ba_lock.write() = new_ba.clone(); + } + } + return Ok(existing); + } + if super::drop_dynamic_pool(&identifier) { + info!( + "[{username}@{pool_name}] auth_query: per-user startup_parameters overlay drift on login — dynamic pool dropped, rebuilding" ); - *ba_lock.write() = new_ba.clone(); } - return Ok(existing); } let config = get_config(); @@ -130,14 +158,12 @@ pub fn create_dynamic_pool( // snapshot. Dynamic auth_query pools follow the same lifecycle as // static pools: rebuilt on RELOAD when the underlying base changes // (see `general_startup_parameters_changed` in pool/mod.rs). - let base_startup_parameters = { - let mut merged: std::collections::BTreeMap = - config.general.startup_parameters.clone(); - for (k, v) in &pool_config.startup_parameters { - merged.insert(k.clone(), v.clone()); - } - std::sync::Arc::new(merged) - }; + let base_startup_parameters = std::sync::Arc::new( + crate::config::startup_parameters::cascade_canonical_keys(&[ + &config.general.startup_parameters, + &pool_config.startup_parameters, + ]), + ); // Convert the caller's HashMap snapshot into the BTreeMap shape // ServerPool stores. The snapshot comes from the auth_query row used @@ -185,11 +211,12 @@ pub fn create_dynamic_pool( per_user_startup_overlay.clone(), ); - // Snapshot the overlay hash before the Arc moves into ServerPool. // The auth_query cache compares the new fetched per-user map against // this value after every refetch; a mismatch drops the dynamic pool - // so the next connect rebuilds with the new reset_val. - let overlay_hash = super::per_user_overlay_hash(per_user_startup_overlay.iter()); + // so the next connect rebuilds with the new reset_val. The caller + // already has the hash precomputed on the `CacheEntry`, so we reuse + // it instead of re-running per_user_overlay_hash on the same map. + let overlay_hash = fetched_overlay_hash; let queue_strategy = match config.general.server_round_robin { true => QueueMode::Fifo, @@ -248,15 +275,36 @@ pub fn create_dynamic_pool( let current = POOLS.load(); let mut new_pools = (**current).clone(); - // Re-check after clone (another thread may have created it) + // Re-check after clone (another thread may have created it). The + // fast path at the top of this function already validates the + // overlay hash; do the same here so the slow path doesn't reuse a + // pool another login built with a stale `startup_parameters` + // snapshot. Without this compare, two concurrent logins after an + // auth_query row update can race: one wins the slow path with the + // new overlay, the other finds the loser's `existing` and inherits + // the stale `reset_val` until TTL or RELOAD. if let Some(existing) = new_pools.get(&identifier) { - if let (Some(ref ba_lock), Some(ref new_ba)) = ( - &existing.address.backend_auth, - &conn_pool.address.backend_auth, - ) { - *ba_lock.write() = new_ba.read().clone(); + let live_hash = existing.per_user_startup_overlay_hash; + let is_dyn = super::is_dynamic_pool(&identifier); + if !should_rebuild_for_overlay_drift(live_hash, overlay_hash, is_dyn) { + // Same reasoning as the fast path: refresh backend_auth + // only when the live pool is a hash-matching dynamic. A + // static pool registered concurrently with the in-flight + // dynamic-pool build is preserved unchanged. + if live_hash == overlay_hash { + if let (Some(ref ba_lock), Some(ref new_ba)) = ( + &existing.address.backend_auth, + &conn_pool.address.backend_auth, + ) { + *ba_lock.write() = new_ba.read().clone(); + } + } + return Ok(existing.clone()); } - return Ok(existing.clone()); + info!( + "[{username}@{pool_name}] auth_query: per-user startup_parameters overlay drift on slow-path race — replacing concurrently-built pool" + ); + new_pools.remove(&identifier); } let auth_method = match &conn_pool.address.backend_auth { @@ -301,3 +349,46 @@ pub fn create_dynamic_pool( Ok(conn_pool) } + +/// Decide whether `create_dynamic_pool` should replace an existing +/// `(pool, user)` entry in `POOLS`. Hash drift alone is not enough — +/// a static pool registered for the same identifier (during a config +/// reload race with an in-flight auth_query login) keeps the empty +/// overlay hash, and replacing it would silently swap the operator's +/// configured backend auth/startup-parameters for the auth_query +/// passthrough version. Rebuild only when the live pool is dynamic. +fn should_rebuild_for_overlay_drift(live_hash: u64, fetched_hash: u64, is_dynamic: bool) -> bool { + live_hash != fetched_hash && is_dynamic +} + +#[cfg(test)] +mod tests { + use super::should_rebuild_for_overlay_drift; + + #[test] + fn overlay_drift_reuses_on_hash_match() { + let h = 0x1234_5678_9abc_def0_u64; + assert!(!should_rebuild_for_overlay_drift(h, h, true)); + assert!(!should_rebuild_for_overlay_drift(h, h, false)); + } + + #[test] + fn overlay_drift_rebuilds_dynamic_on_hash_mismatch() { + assert!(should_rebuild_for_overlay_drift(0xAAAA, 0xBBBB, true)); + } + + #[test] + fn overlay_drift_preserves_static_on_hash_mismatch() { + // A static pool registered during reload races with an + // in-flight auth_query login that fetched a non-empty overlay. + // The live pool's hash is `empty_overlay_hash()`; the fetched + // hash is non-empty. Static-overrides-dynamic must hold, so + // the existing pool wins and is not replaced. + let empty = crate::pool::empty_overlay_hash(); + let fetched = 0xBEEF_0000_0000_0001_u64; + assert_ne!(empty, fetched); + assert!(!should_rebuild_for_overlay_drift( + empty, fetched, /*is_dynamic=*/ false + )); + } +} diff --git a/src/pool/mod.rs b/src/pool/mod.rs index c96f52b6b..eff81d95c 100644 --- a/src/pool/mod.rs +++ b/src/pool/mod.rs @@ -535,14 +535,12 @@ impl ConnectionPool { // second reload between this iteration and constructor // execution would write a different baseline to the pool // than the one the reuse hash captured. - let base_startup_parameters = { - let mut merged: std::collections::BTreeMap = - config.general.startup_parameters.clone(); - for (k, v) in &pool_config.startup_parameters { - merged.insert(k.clone(), v.clone()); - } - Arc::new(merged) - }; + let base_startup_parameters = Arc::new( + crate::config::startup_parameters::cascade_canonical_keys(&[ + &config.general.startup_parameters, + &pool_config.startup_parameters, + ]), + ); let manager = ServerPool::new( address.clone(), @@ -753,14 +751,12 @@ impl ConnectionPool { let fallback_state = build_fallback_state(pool_name, pool_config, &config.general); - let base_startup_parameters = { - let mut merged: std::collections::BTreeMap = - config.general.startup_parameters.clone(); - for (k, v) in &pool_config.startup_parameters { - merged.insert(k.clone(), v.clone()); - } - Arc::new(merged) - }; + let base_startup_parameters = Arc::new( + crate::config::startup_parameters::cascade_canonical_keys(&[ + &config.general.startup_parameters, + &pool_config.startup_parameters, + ]), + ); let manager = ServerPool::new( address.clone(), diff --git a/src/pool/server_pool.rs b/src/pool/server_pool.rs index 5f8ae3d9d..f5e7377eb 100644 --- a/src/pool/server_pool.rs +++ b/src/pool/server_pool.rs @@ -4,7 +4,7 @@ //! server connections. It handles connect timeouts, lifetime checks, alive //! checks, pause/resume, and reconnect epoch management. -use std::collections::BTreeMap; +use std::collections::{BTreeMap, HashSet}; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::Arc; use std::time::Duration; @@ -59,6 +59,107 @@ impl BudgetReason { } } +/// Precompute the wire-ready `startup_parameters` map for a pool and +/// the budget classification reported on each backend start. Called +/// once from `ServerPool::new`; saves the startup path from +/// cloning the BTreeMap, re-running `packet_and_body_bytes`, and +/// re-evaluating the overlay-drop branch on every checkout. +/// +/// Cases follow `classify_startup_parameters`: +/// 1. Empty overlay + budget OK → reuse the base `Arc`. +/// 2. Non-empty overlay + budget OK → allocate the merged map once. +/// 3. Merged over budget but baseline fits → reuse the base `Arc` and +/// log the drop once. The backend start path rejects this decision. +/// 4. Nothing fits → empty map, log the drop. +fn build_resolved_startup( + base: &Arc>, + overlay: &Arc>, + server_username: &str, + database: &str, + application_name: &str, + pool_name: &str, + client_username: &str, +) -> (Arc>, BudgetDecision) { + // PostgreSQL GUC names are case-insensitive, but `general` and + // `pool` are deserialised as opaque maps with exact-case keys. A + // pool that sets `TimeZone = "UTC"` over a + // `general.timezone = "Europe/Berlin"` would otherwise ship both + // rows in `StartupMessage` and let backend-side merge order decide + // which wins. Canonicalise every key here so the merged view, the + // baseline-only fallback, and the admin/API read model all agree + // on the GUC name PostgreSQL will use. + let canonical_base: BTreeMap = base + .iter() + .map(|(k, v)| { + ( + crate::server::parameters::canonicalize_param_name(k.clone()), + v.clone(), + ) + }) + .collect(); + let merged: BTreeMap = if overlay.is_empty() { + canonical_base.clone() + } else { + let mut m = canonical_base.clone(); + for (k, v) in overlay.iter() { + m.insert( + crate::server::parameters::canonicalize_param_name(k.clone()), + v.clone(), + ); + } + m + }; + let (packet_bytes, body_bytes) = + sp::packet_and_body_bytes(server_username, database, application_name, &merged); + let over_budget = body_bytes > sp::MAX_OPERATOR_BUDGET; + let over_packet = packet_bytes > sp::MAX_STARTUP_PACKET_SIZE; + if !over_budget && !over_packet { + return (Arc::new(merged), BudgetDecision::FullCascade); + } + + // Merged cascade is over a limit. Try the canonical baseline-only. + if !overlay.is_empty() { + let (baseline_packet, baseline_body) = + sp::packet_and_body_bytes(server_username, database, application_name, &canonical_base); + if baseline_body <= sp::MAX_OPERATOR_BUDGET + && baseline_packet <= sp::MAX_STARTUP_PACKET_SIZE + { + warn!( + "[{client_username}@{pool_name}] auth_query per-user startup_parameters pushes the cascade \ + over the operator budget (merged {body_bytes} bytes, packet {packet_bytes} bytes); \ + backend spawns will be rejected until the overlay or baseline shrinks" + ); + return ( + Arc::new(canonical_base), + BudgetDecision::OverlayDroppedBaselineKept { + body_bytes, + packet_bytes, + }, + ); + } + } + + let reason = if over_packet { + BudgetReason::PacketCapExceeded + } else { + BudgetReason::CascadeBudgetExceeded + }; + warn!( + "[{client_username}@{pool_name}] effective startup_parameters serialize to {body_bytes} bytes \ + (packet {packet_bytes} bytes), exceeding operator budget {} / PG cap {}; \ + backend spawns will be rejected until general/pool startup_parameters shrink", + sp::MAX_OPERATOR_BUDGET, sp::MAX_STARTUP_PACKET_SIZE, + ); + ( + Arc::new(BTreeMap::new()), + BudgetDecision::EmptyDueToBudget { + reason, + body_bytes, + packet_bytes, + }, + ) +} + /// Wrapper for the connection pool. pub struct ServerPool { /// Server address. @@ -121,16 +222,6 @@ pub struct ServerPool { /// Notify to wake up clients blocked on PAUSE. resume_notify: Notify, - /// `general.startup_parameters` merged with this pool's - /// `pool.startup_parameters` — the baseline that every backend spawn - /// from this pool ships in `StartupMessage`, before the optional - /// per-user auth_query overlay is applied. Cached once at pool - /// construction: the reload path rebuilds the pool whenever either - /// layer's hash changes (see `ConnectionPool::from_config`), so this - /// view is immutable for the lifetime of the pool object. Shared as - /// `Arc` so backend spawns can borrow it without per-call cloning. - base_startup_parameters: Arc>, - /// Per-user auth_query overlay captured at pool construction. Dynamic /// passthrough pools populate this from a fresh `cache.get_or_fetch` /// snapshot taken right after auth, so every backend spawn from this @@ -141,6 +232,30 @@ pub struct ServerPool { /// drain logic in `pool/mod.rs`; this field is immutable for the /// lifetime of the pool object. per_user_startup_overlay: Arc>, + + /// Canonicalized GUC names this pool injects through + /// `startup_parameters` (general + pool + auth_query overlay). The + /// backend startup path (`Server::startup`) uses this set so that + /// `sync_parameters` on checkout does not let a client value overwrite + /// the operator default. The client startup path filters the + /// `StartupMessage` parameters against the same set before sending + /// `ParameterStatus`, so a driver sees the same value PG will use for + /// the session. Built once per pool and shared as `Arc`. + operator_managed_startup_keys: Arc>, + + /// Wire-ready merged startup map, precomputed at pool creation. The + /// spawn path used to recompute this on every backend create — clone + /// the baseline, merge the auth_query overlay, recheck the budget. + /// `base_startup_parameters` and `per_user_startup_overlay` are + /// immutable for the pool's lifetime, so the merge is too: every call + /// can hand out an `Arc` to the same map instead of cloning a + /// 50-entry BTreeMap and re-running `packet_and_body_bytes`. + resolved_startup_map: Arc>, + + /// Budget classification for `resolved_startup_map`. Precomputed + /// together with the map so the spawn path skips the re-validation + /// walk on every backend create. + resolved_startup_decision: BudgetDecision, } impl std::fmt::Debug for ServerPool { @@ -189,6 +304,26 @@ impl ServerPool { base_startup_parameters: Arc>, per_user_startup_overlay: Arc>, ) -> ServerPool { + let operator_managed_startup_keys = Arc::new( + base_startup_parameters + .keys() + .chain(per_user_startup_overlay.keys()) + .map(|k| crate::server::parameters::canonicalize_param_name(k.clone())) + .collect::>(), + ); + let server_username = user + .server_username + .as_deref() + .unwrap_or(user.username.as_str()); + let (resolved_startup_map, resolved_startup_decision) = build_resolved_startup( + &base_startup_parameters, + &per_user_startup_overlay, + server_username, + database, + &application_name, + address.pool_name.as_str(), + user.username.as_str(), + ); ServerPool { address, user: user.clone(), @@ -209,11 +344,18 @@ impl ServerPool { resume_notify: Notify::new(), session_mode, fallback_state, - base_startup_parameters, per_user_startup_overlay, + operator_managed_startup_keys, + resolved_startup_map, + resolved_startup_decision, } } + /// See `operator_managed_startup_keys` field. + pub fn operator_managed_startup_keys(&self) -> Arc> { + self.operator_managed_startup_keys.clone() + } + /// Attempts to create a new connection. /// Uses a semaphore to limit concurrent connection creation instead of serializing with mutex. pub async fn create(&self) -> Result { @@ -272,6 +414,16 @@ impl ServerPool { self.address.host, self.address.port, ); + // Resolve before any `ServerStats` is registered. The budget + // preflight can return `ServerStartupParameterRejection`; if we + // had already published the stats entry via `stats.register`, + // the `?` exit would leak that entry in `sv_login` forever + // because the disconnect path runs only after the spawn + // attempt produces a result. The resolved map and the TLS + // retry share the same Arc, so the plain attempt and the + // sslmode=allow retry still see one parameter set. + let startup_parameters = self.resolved_startup_parameters()?; + let stats = Arc::new(ServerStats::new( self.address.clone(), crate::utils::clock::now(), @@ -279,12 +431,6 @@ impl ServerPool { stats.register(stats.clone()); - // Resolve once for this spawn attempt: the plain attempt and the - // optional sslmode=allow TLS retry must see the same parameter set, - // otherwise a config RELOAD landing between the two would silently - // ship different StartupMessages for the same client request. - let startup_parameters = self.resolved_startup_parameters(); - let result = startup_with_timeout( self.connect_timeout, &self.address.host, @@ -301,6 +447,7 @@ impl ServerPool { self.application_name.clone(), self.session_mode, &startup_parameters, + self.operator_managed_startup_keys.clone(), ), ) .await; @@ -358,6 +505,7 @@ impl ServerPool { self.application_name.clone(), self.session_mode, &startup_parameters, + self.operator_managed_startup_keys.clone(), ), ) .await; @@ -459,8 +607,20 @@ impl ServerPool { // — both are side effects of the spawn path // (resolved_startup_parameters). SHOW polling and /api/pools // page refreshes are safe to call repeatedly. - let (wire_cow, _decision) = self.classify_startup_parameters(); + let (wire_cow, decision) = self.classify_startup_parameters(); let wire = wire_cow.as_ref(); + // If an auth_query overlay would push the cascade over budget, + // every overlay key is dropped due to budget, even when the + // baseline carries the same key with another value. Both + // over-budget decisions reject backend startup, so the admin/API + // view must show that no configured key reaches PostgreSQL. + // Reporting those keys as `Applied` or `Stale` would point the + // operator at RELOAD/refetch, which cannot fix a budget failure. + let cascade_rejected = matches!( + decision, + BudgetDecision::OverlayDroppedBaselineKept { .. } + | BudgetDecision::EmptyDueToBudget { .. } + ); let mut out: std::collections::BTreeMap< String, ( @@ -471,10 +631,14 @@ impl ServerPool { > = configured .into_iter() .map(|(k, (v, src))| { - let state = match wire.get(&k) { - Some(wire_v) if wire_v == &v => ApplicationState::Applied, - Some(_) => ApplicationState::Stale, - None => ApplicationState::DroppedDueToBudget, + let state = if cascade_rejected { + ApplicationState::DroppedDueToBudget + } else { + match wire.get(&k) { + Some(wire_v) if wire_v == &v => ApplicationState::Applied, + Some(_) => ApplicationState::Stale, + None => ApplicationState::DroppedDueToBudget, + } }; (k, (v, src, state)) }) @@ -486,6 +650,11 @@ impl ServerPool { // but my backends are still getting force_custom_plan" — that is // exactly the stale snapshot RELOAD has not yet recycled (or that // the auth_query cache has not yet refetched). + let wire_state = if cascade_rejected { + ApplicationState::DroppedDueToBudget + } else { + ApplicationState::Stale + }; for (k, wire_v) in wire { if !out.contains_key(k) { let frozen_source = if self.per_user_startup_overlay.contains_key(k) { @@ -493,96 +662,28 @@ impl ServerPool { } else { super::startup_resolver::ParameterSource::Pool }; - out.insert( - k.clone(), - (wire_v.clone(), frozen_source, ApplicationState::Stale), - ); + out.insert(k.clone(), (wire_v.clone(), frozen_source, wire_state)); } } out } - /// Resolve the operator-supplied startup_parameters map that this pool /// Pure classifier shared between the spawn path and the read-only /// admin/API views. Returns the wire-ready map and the budget /// decision the runtime would make for this spawn, **without** - /// touching `STARTUP_PARAMETERS_DROPPED_TOTAL` or emitting warn - /// logs. The spawn-side `resolved_startup_parameters` wraps this - /// with the counter + log; admin/API callers call this directly so - /// `SHOW STARTUP_PARAMETERS` polling cannot inflate the drop counter - /// or spam the warn log. + /// touching `STARTUP_PARAMETERS_DROPPED_TOTAL`. Both fields are + /// precomputed at pool construction (`build_resolved_startup`), so + /// every call hands back a borrow of the cached `Arc` instead of + /// cloning and re-validating the merge. fn classify_startup_parameters( &self, ) -> ( std::borrow::Cow<'_, BTreeMap>, BudgetDecision, ) { - let merged: std::borrow::Cow<'_, BTreeMap> = - if self.per_user_startup_overlay.is_empty() { - std::borrow::Cow::Borrowed(&*self.base_startup_parameters) - } else { - let mut owned = (*self.base_startup_parameters).clone(); - for (k, v) in self.per_user_startup_overlay.iter() { - owned.insert(k.clone(), v.clone()); - } - std::borrow::Cow::Owned(owned) - }; - - let username_for_wire = self - .user - .server_username - .as_deref() - .unwrap_or(self.user.username.as_str()); - - let (packet_bytes, body_bytes) = sp::packet_and_body_bytes( - username_for_wire, - &self.database, - &self.application_name, - &merged, - ); - - let over_budget = body_bytes > sp::MAX_OPERATOR_BUDGET; - let over_packet = packet_bytes > sp::MAX_STARTUP_PACKET_SIZE; - if !over_budget && !over_packet { - return (merged, BudgetDecision::FullCascade); - } - - // The merged cascade is over a limit. If we have an overlay, try - // dropping just the overlay and keep the baseline. - let has_overlay = !self.per_user_startup_overlay.is_empty(); - if has_overlay { - let baseline = &*self.base_startup_parameters; - let (baseline_packet, baseline_body) = sp::packet_and_body_bytes( - username_for_wire, - &self.database, - &self.application_name, - baseline, - ); - if baseline_body <= sp::MAX_OPERATOR_BUDGET - && baseline_packet <= sp::MAX_STARTUP_PACKET_SIZE - { - return ( - std::borrow::Cow::Borrowed(baseline), - BudgetDecision::OverlayDroppedBaselineKept { - body_bytes, - packet_bytes, - }, - ); - } - } - - let reason = if over_packet { - BudgetReason::PacketCapExceeded - } else { - BudgetReason::CascadeBudgetExceeded - }; ( - std::borrow::Cow::Owned(BTreeMap::new()), - BudgetDecision::EmptyDueToBudget { - reason, - body_bytes, - packet_bytes, - }, + std::borrow::Cow::Borrowed(&*self.resolved_startup_map), + self.resolved_startup_decision, ) } @@ -593,51 +694,75 @@ impl ServerPool { /// values. /// /// The merged map is checked against the operator budget and the full - /// PostgreSQL startup-packet limit. On overflow, pg_doorman logs and sends - /// no operator-supplied parameters for this backend startup. - fn resolved_startup_parameters(&self) -> std::borrow::Cow<'_, BTreeMap> { + /// PostgreSQL startup-packet limit. Any overflow, whether caused by + /// the auth_query overlay, the general/pool baseline, or the full + /// StartupMessage, becomes a PG-style + /// `ServerStartupParameterRejection`. That gives the client a clear + /// error instead of silently dropping configured startup parameters. + fn resolved_startup_parameters( + &self, + ) -> Result>, Error> { let (map, decision) = self.classify_startup_parameters(); match decision { - BudgetDecision::FullCascade => map, + BudgetDecision::FullCascade => Ok(map), BudgetDecision::OverlayDroppedBaselineKept { body_bytes, packet_bytes, } => { - warn!( - "[{}@{}] auth_query per-user startup_parameters pushes the cascade \ - over the operator budget (merged {} bytes, packet {} bytes); \ - dropping the per-user overlay and keeping the general/pool \ - baseline for this backend spawn", - self.user.username, self.address.pool_name, body_bytes, packet_bytes, - ); crate::web::metrics::STARTUP_PARAMETERS_DROPPED_TOTAL .with_label_values(&[ self.address.pool_name.as_str(), "auth_query_overlay_oversize", ]) .inc(); - map + Err(Error::ServerStartupParameterRejection { + sqlstate: "53400".to_string(), + message: format!( + "auth_query startup_parameters for pool '{}' would exceed the \ + operator budget (merged body {} bytes, full packet {} bytes; \ + operator budget {} bytes, PG StartupMessage cap {} bytes). \ + Reduce the per-user startup_parameters row, the pool baseline, \ + or the general baseline.", + self.address.pool_name, + body_bytes, + packet_bytes, + sp::MAX_OPERATOR_BUDGET, + sp::MAX_STARTUP_PACKET_SIZE, + ), + server_identifier: crate::app::errors::ServerIdentifier::new( + self.user.username.clone(), + &self.database, + &self.address.pool_name, + ), + }) } BudgetDecision::EmptyDueToBudget { reason, body_bytes, packet_bytes, } => { - warn!( - "[{}@{}] effective startup_parameters serialize to {} bytes \ - (packet {} bytes), exceeding operator budget {} / PG cap {}; \ - all operator-supplied parameters dropped for this backend spawn", - self.user.username, - self.address.pool_name, - body_bytes, - packet_bytes, - sp::MAX_OPERATOR_BUDGET, - sp::MAX_STARTUP_PACKET_SIZE, - ); crate::web::metrics::STARTUP_PARAMETERS_DROPPED_TOTAL .with_label_values(&[self.address.pool_name.as_str(), reason.as_str()]) .inc(); - map + Err(Error::ServerStartupParameterRejection { + sqlstate: "53400".to_string(), + message: format!( + "startup_parameters cascade for pool '{}' does not fit the operator \ + budget (body {} bytes, full packet {} bytes; operator budget {} bytes, \ + PG StartupMessage cap {} bytes). Reduce general or pool \ + startup_parameters.", + self.address.pool_name, + body_bytes, + packet_bytes, + sp::MAX_OPERATOR_BUDGET, + sp::MAX_STARTUP_PACKET_SIZE, + ), + server_identifier: crate::app::errors::ServerIdentifier::new( + self.user.username.clone(), + &self.database, + &self.address.pool_name, + ), + }) } } } @@ -690,25 +815,34 @@ impl ServerPool { }; let (result, source) = self.run_fallback_round(fallback).await; - match (result, source) { - (Ok(conn), _) => Ok(conn), - (Err(err), super::fallback::TargetSource::WhitelistCache) => { - // Cached host was stale; wipe it and try with full discovery - // exactly once more. If discovery fails too, return that - // failure without a third attempt. - info!( - "[{}@{}] fallback: whitelist round failed ({err}), retrying with fresh discovery", - self.address.username, self.address.pool_name, - ); - fallback.clear_whitelist(); - let (retry_result, _) = self.run_fallback_round(fallback).await; - retry_result.map_err(|e2| { - Error::ConnectError(format!( - "fallback exhausted (whitelist round: {err}; discovery round: {e2})" - )) - }) - } - (Err(err), super::fallback::TargetSource::Discovery) => Err(err), + match result { + Ok(conn) => Ok(conn), + // The startup_parameters preflight is host-independent. If it + // rejected the cascade, every candidate would fail with the + // same SQLSTATE until the operator fixes the config. Return + // the rejection unchanged so the client sees the PG-style + // error, and skip the whitelist clear so a healthy whitelist + // is not wiped because of a configuration bug. + Err(err @ Error::ServerStartupParameterRejection { .. }) => Err(err), + Err(err) => match source { + super::fallback::TargetSource::WhitelistCache => { + // Cached host was stale; wipe it and try with full + // discovery exactly once more. If discovery fails + // too, return that failure without a third attempt. + info!( + "[{}@{}] fallback: whitelist round failed ({err}), retrying with fresh discovery", + self.address.username, self.address.pool_name, + ); + fallback.clear_whitelist(); + let (retry_result, _) = self.run_fallback_round(fallback).await; + retry_result.map_err(|e2| { + Error::ConnectError(format!( + "fallback exhausted (whitelist round: {err}; discovery round: {e2})" + )) + }) + } + super::fallback::TargetSource::Discovery => Err(err), + }, } } @@ -770,8 +904,12 @@ impl ServerPool { // for one client checkout — and the merge result is host- // independent, so the per-candidate work was pure waste. // `try_fallback_target` borrows this map for both the plain - // attempt and the optional sslmode=allow TLS retry. - let startup_parameters_round = self.resolved_startup_parameters(); + // attempt and the optional sslmode=allow TLS retry. An + // over-budget cascade fails the whole fallback round here. + let startup_parameters_round = match self.resolved_startup_parameters() { + Ok(map) => map, + Err(err) => return (Err(err), source), + }; // Whitelist-cache hit: single target, race-of-one is just a startup. if matches!(source, super::fallback::TargetSource::WhitelistCache) { @@ -883,14 +1021,16 @@ impl ServerPool { "[{}@{}] fallback: all fallback candidates rejected ({summary_str})", self.address.username, self.address.pool_name, ); - // If every candidate failed solely on operator-supplied startup - // parameter rejection, surface PG's actual sqlstate/message so the - // client gets the real error instead of a generic 53300. Healthy - // hosts are not blacklisted (mark_unhealthy skips this category), - // so the same misconfiguration will keep failing until the - // operator fixes the config. - if summary.all_startup_parameter_rejection() { - if let Some(err) = summary.into_last_err() { + // If at least one fallback candidate reached PG's StartupMessage + // stage and got rejected for an operator-supplied startup + // parameter, surface PG's actual sqlstate/message. The + // misconfiguration is host-independent — every other reachable + // candidate would fail the same way once we got past transport + // issues — and healthy hosts are not blacklisted on this reason, + // so wrapping it in 53300 would just rewrite the actionable + // error. + if summary.has_startup_parameter_rejection() { + if let Some(err) = summary.into_startup_rejection() { return (Err(err), source); } } @@ -1023,6 +1163,7 @@ impl ServerPool { self.application_name.clone(), self.session_mode, startup_parameters, + self.operator_managed_startup_keys.clone(), ), ) .await; @@ -1074,6 +1215,7 @@ impl ServerPool { self.application_name.clone(), self.session_mode, startup_parameters, + self.operator_managed_startup_keys.clone(), ), ) .await; @@ -1249,19 +1391,61 @@ fn is_backend_unreachable(err: &Error) -> bool { #[derive(Default)] struct FailureSummary { last_err: Option, + /// The first `ServerStartupParameterRejection` observed in the + /// round. PG rejection of an operator startup_parameter is a + /// host-independent config error: even one PG candidate that + /// reached the StartupMessage stage and rejected the cascade + /// means the same SQLSTATE/message will repeat on every other + /// reachable PG. Mixed-failure rounds with one PG rejection plus + /// transport failures on the rest must still surface the PG + /// error so the operator sees the actionable cause. + typed_startup_rejection: Option, counts: std::collections::HashMap, } impl FailureSummary { fn record(&mut self, err: Error, reason: super::fallback::FailureReason) { *self.counts.entry(reason).or_insert(0) += 1; + if matches!( + reason, + super::fallback::FailureReason::StartupParameterRejection + ) && self.typed_startup_rejection.is_none() + { + // Clone the rejection — the original moves into `last_err` + // so the aggregate `format()` and surrounding code keep + // working unchanged. + if let Error::ServerStartupParameterRejection { + sqlstate, + message, + server_identifier, + } = &err + { + self.typed_startup_rejection = Some(Error::ServerStartupParameterRejection { + sqlstate: sqlstate.clone(), + message: message.clone(), + server_identifier: server_identifier.clone(), + }); + } + } self.last_err = Some(err); } - /// True when the recorded failures are non-empty and contain only - /// `StartupParameterRejection`. Used to decide whether to surface the - /// original PG error to the client instead of the aggregate - /// "all candidates rejected" wrapper. + /// Was at least one fallback candidate rejected by PG specifically + /// for an operator-supplied startup_parameter? Mixed-failure rounds + /// use this to keep the typed PG error from being wrapped in the + /// generic transport-aggregate message. + fn has_startup_parameter_rejection(&self) -> bool { + self.typed_startup_rejection.is_some() + } + + fn into_startup_rejection(self) -> Option { + self.typed_startup_rejection + } + + /// Kept for the legacy "every failure was rejection" tests. The + /// production path uses `has_startup_parameter_rejection` and only + /// needs the typed copy retained separately. + #[cfg(test)] fn all_startup_parameter_rejection(&self) -> bool { !self.counts.is_empty() && self @@ -1270,6 +1454,7 @@ impl FailureSummary { .all(|r| matches!(r, super::fallback::FailureReason::StartupParameterRejection)) } + #[cfg(test)] fn into_last_err(self) -> Option { self.last_err } @@ -1436,8 +1621,34 @@ mod tests { ); assert!( !s.all_startup_parameter_rejection(), - "a single non-rejection cause must veto the shortcut" + "a single non-rejection cause must veto the all-rejection shortcut" ); + // But the typed rejection must still be retained so a mixed + // round of transport failures + one PG rejection surfaces the + // PG SQLSTATE/message instead of the generic aggregate. + assert!( + s.has_startup_parameter_rejection(), + "a single PG rejection must be retained even alongside transport failures" + ); + } + + #[test] + fn into_startup_rejection_returns_first_typed_rejection() { + let mut s = FailureSummary::default(); + s.record( + timeout_err(), + super::super::fallback::FailureReason::Timeout, + ); + s.record( + rejection_err(), + super::super::fallback::FailureReason::StartupParameterRejection, + ); + match s.into_startup_rejection() { + Some(Error::ServerStartupParameterRejection { sqlstate, .. }) => { + assert_eq!(sqlstate, "22023") + } + other => panic!("expected the retained rejection, got {other:?}"), + } } #[test] diff --git a/src/pool/startup_resolver.rs b/src/pool/startup_resolver.rs index 895da0408..9f57cfbfe 100644 --- a/src/pool/startup_resolver.rs +++ b/src/pool/startup_resolver.rs @@ -89,16 +89,25 @@ pub fn resolve_with_sources( pool: &BTreeMap, auth_query_params: Option<&HashMap>, ) -> BTreeMap { + // Canonicalise keys at insert so the read model matches the + // wire-ready map produced by `cascade_canonical_keys` in the + // runtime. Without this, a `timezone` baseline and a `TimeZone` + // override surface as two entries in `SHOW STARTUP_PARAMETERS` / + // `/api/pools`, and the wire compare in + // `effective_startup_parameters_with_sources` flags one variant as + // `dropped_due_to_budget`/`stale`. Later layers win because their + // canonical key replaces earlier inserts at the same key. + let canon = crate::server::parameters::canonicalize_param_name; let mut out: BTreeMap = BTreeMap::new(); for (k, v) in general { - out.insert(k.clone(), (v.clone(), ParameterSource::General)); + out.insert(canon(k.clone()), (v.clone(), ParameterSource::General)); } for (k, v) in pool { - out.insert(k.clone(), (v.clone(), ParameterSource::Pool)); + out.insert(canon(k.clone()), (v.clone(), ParameterSource::Pool)); } if let Some(extra) = auth_query_params { for (k, v) in extra { - out.insert(k.clone(), (v.clone(), ParameterSource::AuthQuery)); + out.insert(canon(k.clone()), (v.clone(), ParameterSource::AuthQuery)); } } out @@ -187,6 +196,59 @@ mod tests { ); } + #[test] + fn resolve_with_sources_canonicalises_tracked_keys() { + // PG GUC lookup is case-insensitive; the runtime cascade + // canonicalises names before they reach the wire (e.g. pool + // `TimeZone` collapses with general `timezone`). The read model + // backing `SHOW STARTUP_PARAMETERS` and `/api/pools` must do the + // same — otherwise the same logical GUC appears twice in the + // SHOW output and the wire-map compare flags one variant as + // `dropped_due_to_budget`/`stale`. + let g = b(&[("timezone", "UTC")]); + let p = b(&[("TimeZone", "Europe/Moscow")]); + let r = resolve_with_sources(&g, &p, None); + assert!( + !r.contains_key("timezone"), + "raw lowercase key must not survive canonicalisation" + ); + assert_eq!( + r.get("TimeZone"), + Some(&("Europe/Moscow".to_string(), ParameterSource::Pool)) + ); + } + + #[test] + fn resolve_with_sources_canonicalises_non_tracked_keys() { + // For non-tracked GUCs canonicalisation is `to_ascii_lowercase`, + // mirroring `canonicalize_param_name`. Without it `Work_Mem` + // (pool) and `work_mem` (general) would both appear in the SHOW + // output and only one would be present on the wire. + let g = b(&[("work_mem", "64MB")]); + let p = b(&[("Work_Mem", "256MB")]); + let r = resolve_with_sources(&g, &p, None); + assert!(!r.contains_key("Work_Mem")); + assert_eq!( + r.get("work_mem"), + Some(&("256MB".to_string(), ParameterSource::Pool)) + ); + } + + #[test] + fn resolve_with_sources_canonicalises_auth_query_keys() { + // auth_query JSON rows can return any casing for tracked GUCs. + // The read model must canonicalise them too so the auth_query + // overlay overrides the same canonical key the runtime ships. + let p = b(&[("TimeZone", "UTC")]); + let a = h(&[("timezone", "Europe/Moscow")]); + let r = resolve_with_sources(&BTreeMap::new(), &p, Some(&a)); + assert_eq!( + r.get("TimeZone"), + Some(&("Europe/Moscow".to_string(), ParameterSource::AuthQuery)) + ); + assert!(!r.contains_key("timezone")); + } + #[test] fn application_name_can_cascade_too() { // operator-wins on application_name extends through the cascade: diff --git a/src/server/parameters.rs b/src/server/parameters.rs index b5f98318e..edc20b1bf 100644 --- a/src/server/parameters.rs +++ b/src/server/parameters.rs @@ -8,25 +8,35 @@ static TRACKED_PARAMETERS: Lazy> = Lazy::new(|| { let mut set = HashSet::new(); set.insert("client_encoding".to_string()); set.insert("DateStyle".to_string()); + set.insert("IntervalStyle".to_string()); set.insert("TimeZone".to_string()); set.insert("standard_conforming_strings".to_string()); set.insert("application_name".to_string()); set }); -/// Canonicalise a PostgreSQL session parameter name so that startup-time -/// lowercase forms (`timezone`, `datestyle`) match the -/// `ParameterStatus` casing PG sends back (`TimeZone`, `DateStyle`). -/// Used both by `ServerParameters::set_param` (where it has lived since -/// day one) and by `Server::startup` when it captures the operator- -/// managed key set: without the canonical form, `sync_parameters` -/// filters by exact-string match and a client startup value reported -/// as `TimeZone` would overwrite an operator value set as `timezone`. +/// Canonicalise a PostgreSQL session parameter name. PG GUC lookups are +/// case-insensitive, so pg_doorman needs one normalised form per name +/// for every internal compare-by-key path (operator_managed key set, +/// cascade merge, dynamic-pool overlay hash, admin/Web read model). The +/// rule: +/// +/// * Tracked parameters (`TRACKED_PARAMETERS`) return their fixed +/// spelling — the same casing PG reports back in `ParameterStatus`. +/// This keeps `sync_parameters` aligned with what the client expects +/// to see at the wire. +/// * Every other GUC is folded to ASCII lower case. PG itself accepts +/// any casing, but the cascade and admin views need a stable form so +/// `general.work_mem` plus `pool.Work_Mem` collapse to one entry +/// instead of shipping both rows in `StartupMessage`. pub fn canonicalize_param_name(key: String) -> String { - if key == "timezone" { - "TimeZone".to_string() - } else if key == "datestyle" { - "DateStyle".to_string() + for tracked in TRACKED_PARAMETERS.iter() { + if key.eq_ignore_ascii_case(tracked) { + return tracked.clone(); + } + } + if key.chars().any(|c| c.is_ascii_uppercase()) { + key.to_ascii_lowercase() } else { key } @@ -143,3 +153,63 @@ impl From<&ServerParameters> for BytesMut { bytes } } + +#[cfg(test)] +mod tests { + use super::canonicalize_param_name; + + #[test] + fn canonicalize_timezone_matches_any_case() { + assert_eq!(canonicalize_param_name("timezone".to_string()), "TimeZone"); + assert_eq!(canonicalize_param_name("TIMEZONE".to_string()), "TimeZone"); + assert_eq!(canonicalize_param_name("TimeZone".to_string()), "TimeZone"); + assert_eq!(canonicalize_param_name("TimezONE".to_string()), "TimeZone"); + } + + #[test] + fn canonicalize_datestyle_matches_any_case() { + assert_eq!( + canonicalize_param_name("datestyle".to_string()), + "DateStyle" + ); + assert_eq!( + canonicalize_param_name("DATESTYLE".to_string()), + "DateStyle" + ); + assert_eq!( + canonicalize_param_name("DateStyle".to_string()), + "DateStyle" + ); + } + + #[test] + fn canonicalize_intervalstyle_matches_any_case() { + assert_eq!( + canonicalize_param_name("intervalstyle".to_string()), + "IntervalStyle" + ); + assert_eq!( + canonicalize_param_name("INTERVALSTYLE".to_string()), + "IntervalStyle" + ); + } + + #[test] + fn canonicalize_lowercases_non_tracked_keys() { + // PG GUC lookup is case-insensitive, so untracked names must + // collapse to one canonical form too, otherwise a `work_mem` + // baseline and a `Work_Mem` pool override become two rows on + // the wire instead of one cascaded value. + assert_eq!(canonicalize_param_name("work_mem".to_string()), "work_mem"); + assert_eq!(canonicalize_param_name("Work_Mem".to_string()), "work_mem"); + assert_eq!(canonicalize_param_name("WORK_MEM".to_string()), "work_mem"); + assert_eq!( + canonicalize_param_name("statement_timeout".to_string()), + "statement_timeout" + ); + assert_eq!( + canonicalize_param_name("Statement_Timeout".to_string()), + "statement_timeout" + ); + } +} diff --git a/src/server/server_backend.rs b/src/server/server_backend.rs index 344bbd0aa..6749e73f2 100644 --- a/src/server/server_backend.rs +++ b/src/server/server_backend.rs @@ -173,15 +173,6 @@ pub struct Server { operator_managed_startup_keys: Arc>, } -/// Shared empty key set so pools that don't use `startup_parameters` -/// hand every backend the same `Arc` instead of allocating a -/// new empty `HashSet` per spawn. -fn empty_operator_keys() -> Arc> { - static EMPTY: once_cell::sync::Lazy>> = - once_cell::sync::Lazy::new(|| Arc::new(HashSet::new())); - EMPTY.clone() -} - impl std::fmt::Display for Server { fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result { write!( @@ -781,6 +772,7 @@ impl Server { application_name: String, session_mode: bool, startup_parameters: &std::collections::BTreeMap, + operator_managed_startup_keys: Arc>, ) -> Result { let config = get_config(); @@ -1058,36 +1050,6 @@ impl Server { phase_started.elapsed().as_secs_f64(), ); - // The empty case is a shared `Arc` static, so - // pools that don't use the feature pay zero allocation - // per spawn. The non-empty case still allocates once per - // spawn — the caller in pool/server_pool.rs knows the - // map shape but does not currently pass an already-Arc'd - // HashSet through, and lifting the construction up there - // is a larger refactor than this commit warrants. - let operator_managed_startup_keys: Arc> = if startup_parameters - .is_empty() - { - empty_operator_keys() - } else { - // Canonicalize every operator key the same way - // ServerParameters::set_param does on the - // sync_parameters path. Without this an operator - // value configured as `timezone` would not match - // a client-startup value reported as `TimeZone` - // in compare_params(), letting the client - // override the operator default. See codex - // MED #7 (fresh review). - Arc::new( - startup_parameters - .keys() - .map(|k| { - crate::server::parameters::canonicalize_param_name(k.clone()) - }) - .collect(), - ) - }; - let server = Server { address: address.to_owned(), stream: BufStream::new(stream), diff --git a/src/web/metrics/mod.rs b/src/web/metrics/mod.rs index e6e6ce894..7f32816a5 100644 --- a/src/web/metrics/mod.rs +++ b/src/web/metrics/mod.rs @@ -479,7 +479,7 @@ pub(crate) static LISTENER_REJECTIONS_TOTAL: Lazy = Lazy::new(|| /// back to scanning the M-field for any sent key wrapped in double /// quotes (PG keeps the quote markers across all `lc_messages` locales). /// If both heuristics fail — typically because the PG error is unrelated -/// to the sent map at all — the counter does NOT increment. Operator +/// to the sent map at all — the counter does NOT increment. An operator /// reading a non-zero rate can be confident the issue is on a key they /// configured; the per-line warn log carries the parameter name and /// username for triage. @@ -513,16 +513,14 @@ pub(crate) static LISTENER_REJECTIONS_TOTAL: Lazy = Lazy::new(|| /// `MAX_STARTUP_PACKET_LENGTH` (10 000 bytes). Same drop-all /// behaviour. /// * `auth_query_oversize` — the auth_query `startup_parameters` -/// text column for some username exceeded the operator budget at -/// parse time, so the per-user overlay is ignored. +/// value for some username exceeded the operator budget at parse +/// time, so the per-user overlay is ignored. /// * `auth_query_overlay_oversize` — the merged baseline+overlay was -/// over budget, but the baseline alone fits. Keeps general/pool -/// defaults (statement_timeout, lock_timeout, ...) for that user -/// instead of stripping the whole configured parameter set. +/// over budget, but the baseline alone fits. Backend startup is +/// rejected with `SQLSTATE 53400` rather than sending a partial map. /// * `auth_query_bad_type` — the auth_query `startup_parameters` -/// column has a non-text type (likely `json`/`jsonb`); pg_doorman -/// reads it as text, so the row's overlay is dropped. Cast to -/// `::text` in the auth_query SELECT to fix. +/// column has an unsupported PostgreSQL type. Native `text`, `json`, +/// and `jsonb` are accepted. /// * `auth_query_invalid_json` — the column value is not valid JSON. /// * `auth_query_invalid_shape` — the column parses but the /// top-level value is not a JSON object. @@ -536,8 +534,8 @@ pub(crate) static LISTENER_REJECTIONS_TOTAL: Lazy = Lazy::new(|| /// was dropped. Incremented once per such row. /// /// All cases also emit a `warn!` log line for human triage; the -/// counter exists so dashboards and alerts can spot the silent drop -/// without log scraping. +/// counter lets dashboards and alerts catch these drops without log +/// scraping. pub(crate) static STARTUP_PARAMETERS_DROPPED_TOTAL: Lazy = Lazy::new(|| { let counter = IntCounterVec::new( Opts::new( diff --git a/src/web/routes/collect/config.rs b/src/web/routes/collect/config.rs index b4fcc64ed..5e8cb9966 100644 --- a/src/web/routes/collect/config.rs +++ b/src/web/routes/collect/config.rs @@ -31,18 +31,30 @@ fn is_startup_parameter_key(key: &str) -> bool { key.contains(".startup_parameters.") || key.starts_with("startup_parameters.") } -/// Bind-address fields require a restart; everything else takes effect on -/// the next backend or `RELOAD`. Listed precisely so the UI can render -/// the right "restart_required" pill instead of marking everything -/// reloadable. -const IMMUTABLES: &[&str] = &["host", "port", "connect_timeout"]; +/// Listener bind fields and runtime-construction fields require a +/// restart; most other fields take effect on `RELOAD` or on the next +/// backend. The list uses full flattened keys because `/api/config` +/// emits paths such as `general.host` and `web.host`. +/// `worker_threads`, `unix_socket_dir`, and `backlog` shape the Tokio +/// runtime and listener socket at process start; SIGHUP cannot rebuild +/// those. +const IMMUTABLES: &[&str] = &[ + "general.host", + "general.port", + "general.worker_threads", + "general.unix_socket_dir", + "general.backlog", + "web.host", + "web.port", +]; -/// Flatten a serde JSON value into dotted keys → string values. Operators -/// have asked for a coverage-complete `/api/config` so they can verify -/// TLS / auth_query / pool sizing / prepared cache / web settings during -/// an incident — the previous hand-written `From<&Config>` only exposed -/// host/port/connect_timeout/idle_timeout/shutdown_timeout plus pool -/// users/mode, which DBA P3#7 (codex review) flagged as too thin. +fn is_immutable_key(key: &str) -> bool { + IMMUTABLES.contains(&key) +} + +/// Flatten a serde JSON value into dotted keys → string values. This +/// keeps `/api/config` broad enough for incident checks across TLS, +/// auth_query, pool sizing, prepared cache, and web settings. fn flatten_json(prefix: &str, value: &serde_json::Value, out: &mut HashMap) { match value { serde_json::Value::Object(map) => { @@ -97,10 +109,9 @@ pub(crate) fn collect_config(reveal_startup_values: bool) -> ConfigDto { } flat.retain(|k, _| !is_internal_key(k)); - // Diff against `Config::default()` so the UI can show what is at - // its built-in default vs. what an operator changed via the config - // file. The defaults map is computed once per request — cheap, and - // a stable comparison surface for codex DBA P3#7. + // Diff against `Config::default()` so the UI can show what is still + // at the built-in default and what changed in the config file. The + // defaults map is small enough to compute once per request. let mut defaults: HashMap = HashMap::new(); if let Ok(value) = serde_json::to_value(crate::config::Config::default()) { flatten_json("", &value, &mut defaults); @@ -117,11 +128,7 @@ pub(crate) fn collect_config(reveal_startup_values: bool) -> ConfigDto { .get(&key) .map(|d| if mask { "***".to_string() } else { d.clone() }) .unwrap_or_else(|| "-".to_string()); - let changeable = if IMMUTABLES.iter().any(|c| *c == key) { - "no" - } else { - "yes" - }; + let changeable = if is_immutable_key(&key) { "no" } else { "yes" }; let doc = lookup_doc(&key); ConfigEntry { key, @@ -232,6 +239,33 @@ mod tests { assert!(!super::is_secret_key("not_a_secret_check")); } + #[test] + fn is_immutable_key_matches_bind_addresses() { + assert!(super::is_immutable_key("general.host")); + assert!(super::is_immutable_key("general.port")); + assert!(super::is_immutable_key("web.host")); + assert!(super::is_immutable_key("web.port")); + } + + #[test] + fn is_immutable_key_matches_runtime_construction_fields() { + assert!(super::is_immutable_key("general.worker_threads")); + assert!(super::is_immutable_key("general.unix_socket_dir")); + assert!(super::is_immutable_key("general.backlog")); + } + + #[test] + fn is_immutable_key_rejects_reloadable_fields() { + assert!(!super::is_immutable_key("general.idle_timeout")); + assert!(!super::is_immutable_key("general.shutdown_timeout")); + assert!(!super::is_immutable_key("general.connect_timeout")); + assert!(!super::is_immutable_key("pools.app_db.server_host")); + assert!(!super::is_immutable_key("pools.app_db.users.0.username")); + // Bare segment must not match a flattened key. + assert!(!super::is_immutable_key("host")); + assert!(!super::is_immutable_key("port")); + } + /// Coverage check: the previous implementation only exposed /// host/port/connect_timeout/idle_timeout/shutdown_timeout plus /// pool users/mode. The flattened serde view now surfaces every @@ -242,9 +276,8 @@ mod tests { let dto = super::collect_config(true); let keys: std::collections::HashSet<&str> = dto.config.iter().map(|e| e.key.as_str()).collect(); - // Spot-check four orthogonal areas DBA P3#7 called out: - // TLS server-side, prepared cache size, web listener, shutdown - // timeout. + // Spot-check independent areas: core listener, shutdown timeout, + // server-side TLS, and web listener. assert!(keys.contains("general.host"), "{keys:?}"); assert!(keys.contains("general.shutdown_timeout"), "{keys:?}"); assert!(keys.contains("general.server_tls_mode"), "{keys:?}"); diff --git a/src/web/routes/config.rs b/src/web/routes/config.rs index 3b91fd0d3..7f357a313 100644 --- a/src/web/routes/config.rs +++ b/src/web/routes/config.rs @@ -5,11 +5,13 @@ use crate::web::routes::collect::collect_config; use crate::web::server::Response; pub(crate) fn handle_config(role: Role) -> Response { - // Mirror the /api/pools contract: operator-supplied - // startup_parameter values are masked for anonymous readers because - // they may carry tenant identifiers, audit routing tags, or - // accidental secrets. SSO and Admin keep the full view. - let reveal_startup_values = role >= Role::Sso; + // Operator-supplied startup_parameter values can carry tenant + // identifiers, audit routing tags, or accidental secrets. Only Admin + // sees literal values; SSO readers get the same masked view as + // anonymous (key + source + state, value `***`). Because + // `sso_allowed_users = ["*"]` is the default, promoting SSO readers + // here would expose those values to every user accepted by the IdP. + let reveal_startup_values = role >= Role::Admin; Response::ok_json(&collect_config(reveal_startup_values)) } diff --git a/src/web/routes/pools.rs b/src/web/routes/pools.rs index 0c4e70bc7..ce3e7ba0b 100644 --- a/src/web/routes/pools.rs +++ b/src/web/routes/pools.rs @@ -5,11 +5,11 @@ use crate::web::routes::collect::collect_pools; use crate::web::server::Response; pub(crate) fn handle_pools(role: Role) -> Response { - // Anonymous /api/pools must not leak operator-supplied - // startup_parameter values: they can carry tenant identifiers, audit - // tags, or accidental secrets. SSO/Admin callers keep the full view; - // anonymous viewers get parameter+source only. - let reveal_startup_values = role >= Role::Sso; + // Only Admin reads literal startup_parameter values. SSO is treated + // like anonymous here so a broad `sso_allowed_users` list does not + // expose tenant routing tags, audit identifiers, or accidental + // secrets stored in startup_parameters. + let reveal_startup_values = role >= Role::Admin; Response::ok_json(&collect_pools(reveal_startup_values)) } diff --git a/tests/auth_query_jsonb_fixture.sql b/tests/auth_query_jsonb_fixture.sql new file mode 100644 index 000000000..5e70ab3c9 --- /dev/null +++ b/tests/auth_query_jsonb_fixture.sql @@ -0,0 +1,20 @@ +-- Fixture for auth_query rows that return per-user startup_parameters +-- as native jsonb. This covers the path where pg_doorman decodes the +-- column by PostgreSQL type instead of requiring a ::text cast. + +DROP TABLE IF EXISTS auth_users_jsonb; +CREATE TABLE auth_users_jsonb ( + username TEXT NOT NULL, + password TEXT, + startup_parameters JSONB +); + +SET password_encryption = 'md5'; + +DROP USER IF EXISTS sp_jsonb_user; +CREATE USER sp_jsonb_user WITH PASSWORD 'jsonb_pass'; +INSERT INTO auth_users_jsonb + SELECT rolname, rolpassword, '{"plan_cache_mode":"force_custom_plan"}'::jsonb + FROM pg_authid WHERE rolname = 'sp_jsonb_user'; + +GRANT ALL ON DATABASE postgres TO sp_jsonb_user; diff --git a/tests/bdd/features/startup-parameters.feature b/tests/bdd/features/startup-parameters.feature index 15b047c06..6bed5ed1a 100644 --- a/tests/bdd/features/startup-parameters.feature +++ b/tests/bdd/features/startup-parameters.feature @@ -552,3 +552,44 @@ Feature: Per-pool startup_parameters # cascade fits the operator budget for this fixture. And the command output should contain "statement_timeout|10s|general|applied" And the command output should contain "plan_cache_mode|force_custom_plan|pool|applied" + + # A common auth_query schema stores startup_parameters as jsonb. + # pg_doorman must decode that column directly, without forcing the + # operator to add a ::text cast to the SELECT. + Scenario: auth_query startup_parameters from a jsonb column applies without ::text cast + Given PostgreSQL started with pg_hba.conf: + """ + local all all trust + host all postgres 127.0.0.1/32 trust + host all all 127.0.0.1/32 md5 + host all all ::1/128 trust + """ + And fixtures from "tests/auth_query_jsonb_fixture.sql" applied + And pg_doorman started with config: + """ + [general] + host = "127.0.0.1" + port = ${DOORMAN_PORT} + admin_username = "admin" + admin_password = "admin" + pg_hba.content = "host all all 127.0.0.1/32 md5" + + [pools.postgres] + server_host = "127.0.0.1" + server_port = ${PG_PORT} + pool_mode = "transaction" + + [pools.postgres.startup_parameters] + plan_cache_mode = "auto" + + [pools.postgres.auth_query] + query = "SELECT username, password, startup_parameters FROM auth_users_jsonb WHERE username = $1" + user = "postgres" + password = "" + workers = 1 + pool_size = 5 + cache_ttl = "1h" + cache_failure_ttl = "30s" + min_interval = "0s" + """ + Then psql query "SHOW plan_cache_mode" via pg_doorman as user "sp_jsonb_user" to database "postgres" with password "jsonb_pass" returns "force_custom_plan" diff --git a/tests/bdd/features/web-ui.feature b/tests/bdd/features/web-ui.feature index a4cc14a03..c3b6bd6b3 100644 --- a/tests/bdd/features/web-ui.feature +++ b/tests/bdd/features/web-ui.feature @@ -49,8 +49,8 @@ Feature: Web UI listener curl -s -o /dev/null -w "%{http_code} %{content_type}\n" http://127.0.0.1:9127/ """ Then the command should succeed - And output contains "200" - And output contains "text/html" + And the command output should contain "200" + And the command output should contain "text/html" Scenario: Deep link falls back to the SPA shell When I run shell command: @@ -58,8 +58,8 @@ Feature: Web UI listener curl -s -o /dev/null -w "%{http_code} %{content_type}\n" http://127.0.0.1:9127/pools """ Then the command should succeed - And output contains "200" - And output contains "text/html" + And the command output should contain "200" + And the command output should contain "text/html" Scenario: /api/version is anonymous when ui_anonymous is true When I run shell command: @@ -67,7 +67,7 @@ Feature: Web UI listener curl -s -o /dev/null -w "%{http_code}\n" http://127.0.0.1:9127/api/version """ Then the command should succeed - And output contains "200" + And the command output should contain "200" Scenario: /api/logs is admin-only without credentials When I run shell command: @@ -75,7 +75,7 @@ Feature: Web UI listener curl -s -o /dev/null -w "%{http_code}\n" http://127.0.0.1:9127/api/logs """ Then the command should succeed - And output contains "401" + And the command output should contain "401" Scenario: JSON callers receive 401 without WWW-Authenticate When I run shell command: @@ -83,7 +83,7 @@ Feature: Web UI listener curl -s -i -H "Accept: application/json" http://127.0.0.1:9127/api/logs | grep -ci 'WWW-Authenticate' || echo 0 """ Then the command should succeed - And output contains "0" + And the command output should contain "0" Scenario: Curl-style callers do receive WWW-Authenticate on 401 When I run shell command: @@ -91,7 +91,7 @@ Feature: Web UI listener curl -s -i http://127.0.0.1:9127/api/logs | grep -ci 'WWW-Authenticate' """ Then the command should succeed - And output contains "1" + And the command output should contain "1" Scenario: /api/logs accepts admin basic auth When I run shell command: @@ -99,7 +99,7 @@ Feature: Web UI listener curl -s --user 'admin:webui_bdd' -o /dev/null -w "%{http_code}\n" http://127.0.0.1:9127/api/logs """ Then the command should succeed - And output contains "200" + And the command output should contain "200" Scenario: /metrics still serves Prometheus when the UI is on When I run shell command: @@ -107,7 +107,7 @@ Feature: Web UI listener curl -s http://127.0.0.1:9127/metrics | head -1 """ Then the command should succeed - And output contains "# HELP" + And the command output should contain "# HELP" Scenario: /api/auth/config is anonymous and reports SSO disabled When I run shell command: @@ -115,8 +115,8 @@ Feature: Web UI listener curl -s http://127.0.0.1:9127/api/auth/config """ Then the command should succeed - And output contains "\"sso_enabled\":false" - And output contains "\"current_user\":null" + And the command output should contain "\"sso_enabled\":false" + And the command output should contain "\"current_user\":null" Scenario: /api/auth/config carries admin identity for Basic auth When I run shell command: @@ -124,8 +124,8 @@ Feature: Web UI listener curl -s --user 'admin:webui_bdd' http://127.0.0.1:9127/api/auth/config """ Then the command should succeed - And output contains "\"role\":\"admin\"" - And output contains "\"source\":\"basic\"" + And the command output should contain "\"role\":\"admin\"" + And the command output should contain "\"source\":\"basic\"" Scenario: /api/admin/reload requires Admin and rejects anonymous with 401 When I run shell command: @@ -133,7 +133,7 @@ Feature: Web UI listener curl -s -X POST -o /dev/null -w "%{http_code}\n" http://127.0.0.1:9127/api/admin/reload """ Then the command should succeed - And output contains "401" + And the command output should contain "401" Scenario: Bearer with malformed token is Rejected (401) when SSO is off When I run shell command: @@ -141,7 +141,7 @@ Feature: Web UI listener curl -s -H "Authorization: Bearer not.a.real.jwt" -o /dev/null -w "%{http_code}\n" http://127.0.0.1:9127/api/logs """ Then the command should succeed - And output contains "401" + And the command output should contain "401" Scenario: Access log line surfaces in /api/logs after a probe When I run shell command: @@ -151,3 +151,21 @@ Feature: Web UI listener curl -s --user 'admin:webui_bdd' http://127.0.0.1:9127/api/logs | grep -c 'pg_doorman::web::access' || echo 0 """ Then the command should succeed + + # Bind-address fields require a restart; most other fields take effect + # on RELOAD or on the next backend. /api/config must mark the flattened + # bind-address keys as immutable so the SPA shows restart_required for + # the right rows. + Scenario: /api/config marks bind-address fields as restart-required + When I run shell command: + """ + curl -s --user 'admin:webui_bdd' http://127.0.0.1:9127/api/config | \ + grep -oE '"key":"(general|web)\.(host|port)"[^}]*"changeable":"[^"]*"' | sort -u + """ + Then the command should succeed + # Each bind-address key must appear with changeable=no. + And the command output should contain "\"key\":\"general.host\"" + And the command output should contain "\"key\":\"general.port\"" + And the command output should contain "\"key\":\"web.host\"" + And the command output should contain "\"key\":\"web.port\"" + And the command output should contain "\"changeable\":\"no\""