From a3fbe23b4085f57e5e348a83dcdec8434d986630 Mon Sep 17 00:00:00 2001 From: Armand Parajon Date: Mon, 21 Sep 2026 20:00:59 -0400 Subject: [PATCH 1/5] executor: apply row security changes atomically --- .agents/checks/review.md | 5 +- README.md | 5 +- SAFETY.md | 10 +- docs/atomic-row-security.md | 72 ++++ docs/capabilities.md | 5 +- docs/declarative-row-security.md | 27 +- docs/execution-model.md | 1 + docs/invariants.md | 12 + docs/limitations.md | 2 +- docs/tcb-model.md | 2 +- pkg/capabilities/capabilities.yaml | 16 +- pkg/capabilities/capabilities_test.go | 4 +- pkg/executor/code.go | 7 + pkg/executor/code_test.go | 2 + pkg/executor/row_security.go | 138 ++++++++ pkg/executor/row_security_integration_test.go | 330 ++++++++++++++++++ pkg/executor/row_security_test.go | 25 ++ pkg/schemadiff/desired.go | 20 +- pkg/schemadiff/row_security_roundtrip.go | 23 +- pkg/statement/desired_rls.go | 5 +- pkg/statement/row_security_execution.go | 58 +++ pkg/statement/row_security_execution_test.go | 27 ++ 22 files changed, 757 insertions(+), 39 deletions(-) create mode 100644 docs/atomic-row-security.md create mode 100644 pkg/executor/row_security.go create mode 100644 pkg/executor/row_security_integration_test.go create mode 100644 pkg/executor/row_security_test.go create mode 100644 pkg/statement/row_security_execution.go create mode 100644 pkg/statement/row_security_execution_test.go diff --git a/.agents/checks/review.md b/.agents/checks/review.md index f101ffa..9c53432 100644 --- a/.agents/checks/review.md +++ b/.agents/checks/review.md @@ -22,8 +22,9 @@ the reviewer's distillation. package-private constructors — never a raw string or bool that a caller could fabricate. Core code re-verifies its own preconditions; it never trusts that the planner or CLI checked. -- `statement.DesiredWithRowSecurity` is inspection-only: keep it separate from - `DesiredSchema` and refuse it at every live execution boundary. +- `statement.DesiredWithRowSecurity` stays separate from `DesiredSchema`. Its only + live consumer is `executor.ExecuteRowSecurity`, which must enforce RS-1..RS-4; + generic native and create executors must still refuse policy SQL. - Invariant enforcement points carry a `// INV: ` comment matching [docs/invariants.md](../../docs/invariants.md); violations use `ErrInvariantViolation` naming the ID and abort fail-closed — never a warning, never retried. diff --git a/README.md b/README.md index 0b5bdd8..4853ae2 100644 --- a/README.md +++ b/README.md @@ -66,8 +66,9 @@ refusal — never a silently wrong or incomplete result: - **Copy-and-swap** (genuine table rewrites) is not yet available — those changes refuse rather than fall through to a blocking rewrite. -- **Row security** is included in exports when present, with reviewable before/after differences. Applying - policy changes is not supported yet. See the workflow and roadmap in +- **Row security** is included in exports when present, with reviewable before/after differences. The + [atomic Go executor](docs/atomic-row-security.md) can apply RLS-only changes to + existing tables; CLI execution and mixed table/policy changes remain unsupported. See the workflow and roadmap in [declarative-row-security.md](docs/declarative-row-security.md). - **Foreign keys** are out of the declarative model in either direction: desired files cannot declare them, and export refuses both a table that diff --git a/SAFETY.md b/SAFETY.md index 07ddfd2..155bcee 100644 --- a/SAFETY.md +++ b/SAFETY.md @@ -73,9 +73,9 @@ The short version — the full rules live in [docs/tcb-model.md](docs/tcb-model. `preflight.PrivilegedRole`, `preflight.CopySwapTarget`, `dbconn.TableLock`, `checksum.VerifiedShadow`, and `checksum.CleanWatermark`); dangerous APIs accept only proof types — e.g. the planned cutover swap will accept only a `VerifiedShadow`. -- `statement.DesiredWithRowSecurity` proves admission for rolled-back scratch - inspection only. It is distinct from `DesiredSchema` and must never be accepted - by a live executor. +- `statement.DesiredWithRowSecurity` proves declaration syntax, not execution safety. It stays distinct from + `DesiredSchema`. Only `executor.ExecuteRowSecurity` may consume it for live RLS: + that executor locks, checks table equality, and verifies convergence in one transaction. - **Put a limit on everything.** Every loop bounded, every queue bounded, every retry counted, every wait deadlined. An unbounded anything in a core package is a review-blocking defect. - **Assert the positive and the negative space; pair assertions across boundaries.** Invariant @@ -127,3 +127,7 @@ The short version — the full rules live in [docs/tcb-model.md](docs/tcb-model. test-first with the invariant's named test obligation, small diffs, careful review. - **Outside the core: more AI, less steering.** Iterate at inference speed; the boundary means a bug in the periphery cannot corrupt data. + +The atomic RLS executor also admits `pkg/schemadiff` scratch introspection, table +comparison, and render admission into the core. Those calls refuse mixed or +unsupported table shapes; final catalog comparison gates commit (RS-1..RS-4). diff --git a/docs/atomic-row-security.md b/docs/atomic-row-security.md new file mode 100644 index 0000000..0af3255 --- /dev/null +++ b/docs/atomic-row-security.md @@ -0,0 +1,72 @@ +# Atomic row security changes + +The dedicated Go executor converges the complete RLS definition of one existing +ordinary table. It does not change columns, indexes, or constraints and does not +create missing tables. `diff` remains a preview; no new CLI flags are required. + +## Call the Go API + +Use `statement.ParseDesiredWithRowSecurity` on the same SQL you preview with +`diff`, then call the dedicated executor with an existing `dbconn` pool: + +```go +desired, err := statement.ParseDesiredWithRowSecurity(` + CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (owner_id = 7); +`) +if err != nil { + return err +} +report, err := executor.ExecuteRowSecurity(ctx, pool, "public", desired, executor.Budget{ + LockTimeout: 100 * time.Millisecond, + StatementTimeout: 5 * time.Second, +}) +if err != nil { + return err +} +// report.Statements contains the committed statements; an empty slice means no change. +``` + +`StatementTimeout` also caps the entire attempt. There is no approval token or +saved fingerprint. Orchestrators retain their own replan and consent rules. +Changing an RLS definition can widen access even when no data is deleted. + +## Execution contract + +1. Validate the parsed declaration and nonzero budgets. +2. Begin one transaction with a deadline for the entire attempt and transaction-local + lock and statement timeouts. Acquire `ACCESS EXCLUSIVE` on the target. +3. Read the live definition and materialize the desired SQL in a rolled-back + savepoint on the same connection. Refuse unsupported table shapes or any table delta. +4. If RLS already matches, finish without policy DDL. Otherwise replace the complete + policy set (including unchanged policies), apply ENABLE/DISABLE and FORCE/NO FORCE, and preserve policy comments. +5. Read back the complete definition. Commit only if it matches the desired model. + +The lock blocks reads and writes briefly; this is bounded metadata DDL, not an +online copy. Other sessions never see the intermediate policy set. A failure before +commit rolls back every change. A lost commit response is an unknown outcome: +inspect the database before retrying. There are no automatic retries. + +The caller needs table-owner privileges and permission to create the temporary +scratch schema. Roles and qualified helpers must already exist. Grants, role +membership, helper bodies, authentication, and Supabase-managed schemas are outside +this operation. Application authorization tests are still needed. Concurrent +administration of those dependencies is not serialized by the table lock. + +## Invariants and tests + +- **RS-1:** Read the live baseline only after taking the target lock; never accept a + caller-supplied diff as execution authority. Refuse mixed changes before live DDL. +- **RS-2:** All policy/settings changes and the final catalog comparison share one + transaction. Fault injection after live DDL must prove the original state survives. +- **RS-3:** Bound lock waits, individual statements, and the entire attempt. Cancellation + or lock exhaustion must leave the original policies intact. +- **RS-4:** Execute only admitted RLS statements against the qualified target, with a + pg_catalog-only search path. Scratch objects never persist. Verify convergence before commit. + +These checks belong to the executor, regardless of whether a CLI or orchestrator +calls it. SchemaBot can retain its existing replan and consent workflow. diff --git a/docs/capabilities.md b/docs/capabilities.md index 5da4f5f..2eda030 100644 --- a/docs/capabilities.md +++ b/docs/capabilities.md @@ -176,7 +176,7 @@ the canonical example. > `make gen-capabilities`; do not edit the generated regions below by hand. -**53 operations: 18 supported today, 19 planned behind a typed refusal, 14 out of scope +**54 operations: 19 supported today, 19 planned behind a typed refusal, 14 out of scope by design, and 2 with no online mechanism in PostgreSQL to build on.** @@ -269,7 +269,8 @@ review the object warrants) · | PL/pgSQL function bodies (`CREATE OR REPLACE FUNCTION`) | ⚪ | — | No — owner tooling | Transactional catalog work that takes no lock on any relation; nothing for an online engine to add. No peer online executor owns it either | | Triggers (`CREATE TRIGGER`) | ⚪ | — | No — owner tooling | Catalog work — no scan, no rewrite — but it takes a brief `SHARE ROW EXCLUSIVE` on the table, queues behind long-running queries, and blocks writers while it waits — run it under a `lock_timeout` | | Extensions (`CREATE EXTENSION`) | ⚪ | — | No — owner tooling | Same: catalog bootstrap, owner tooling | -| Grants, roles, row-level-security policies | 🔵 | — | No — provisioning / IaC | Access control changes remain with provisioning. Export includes RLS when present; explicit RLS declarations support comparison and review-only deltas with advisory access warnings. Plain table files leave access control separately managed. RLS execution remains unsupported. See the [workflow and roadmap](declarative-row-security.md). See [engine-role.md](engine-role.md) for the engine's own role | +| Grants, roles, and CLI policy DDL | 🔵 | — | No — provisioning / IaC | Access control changes remain with provisioning. Export includes RLS when present; explicit RLS declarations support comparison and review-only deltas with advisory access warnings. Plain table files leave access control separately managed. The separate [atomic RLS Go executor](atomic-row-security.md) supports RLS-only changes on existing supported tables; it does not route through these CLI front doors. See the [workflow and roadmap](declarative-row-security.md). See [engine-role.md](engine-role.md) for the engine's own role | +| Complete table-local RLS definition (Go API only) | ✅ | native, as-is | Yes | The [atomic RLS executor](atomic-row-security.md) locks an existing supported table, refuses structural changes, replaces policies/settings, and verifies convergence before committing. Lock waits and the whole transaction are bounded. No CLI execution or new diff flags | | Standalone sequences | ⚪ | — | No — owner tooling | Transactional catalog work on an object with no readers-and-writers problem | | Publications, subscriptions | 🔵 | — | No — replication provisioning / IaC | Replication provisioning, not table shape (`ALTER PUBLICATION ... ADD TABLE` also takes `SHARE UPDATE EXCLUSIVE` on the table) | diff --git a/docs/declarative-row-security.md b/docs/declarative-row-security.md index e39a009..0770546 100644 --- a/docs/declarative-row-security.md +++ b/docs/declarative-row-security.md @@ -3,8 +3,9 @@ You can export a table's RLS settings and policies, keep them alongside its SQL, and verify that the live definition still matches. `pull` includes RLS when the live table has settings or policies; ordinary tables get no extra SQL. -**Applying changes to RLS is not supported yet.** A difference produces a review -of the captured definitions and a refusal, not SQL to execute. +`diff` shows the captured definitions and refuses to emit execution SQL for RLS +changes. The dedicated [atomic Go executor](atomic-row-security.md) can apply an +RLS-only declaration to an existing supported table. No new CLI flags are needed. ## Export and compare @@ -106,8 +107,8 @@ Equal definitions do not prove equal access: grants, role membership, helper function bodies, and authentication configuration are outside this comparison. Library callers use `RenderWithRowSecurity`, `ParseDesiredWithRowSecurity`, and -`diffplan.PlanWithRowSecurity`. The parser returns a separate inspection-only type -that the existing live executors cannot accept. +`diffplan.PlanWithRowSecurity`. The parser returns a separate declaration type. Only the dedicated atomic RLS +executor accepts it for live execution; generic native and create executors do not. ## Review a difference @@ -183,7 +184,7 @@ complete support matrix. pg-sprite already compares a live table with desired SQL materialized inside a rolled-back scratch transaction. This work extends that model instead of importing another schema engine. PostgreSQL should resolve SQL and supply its catalog representation. -Before enabling policy execution, settle these boundaries: +The execution contract preserves these boundaries: - **Explicit ownership.** Existing table-only files keep access control separately managed. A caller must opt into managing a table's complete RLS definition. @@ -202,7 +203,7 @@ Before enabling policy execution, settle these boundaries: removing a restrictive one, disabling RLS, or changing a role can widen access without deleting data. A data-destruction label is not a complete authorization contract. Do not claim to prove arbitrary predicates equivalent. -- **Atomic transitions.** Recheck the reviewed state under the appropriate lock +- **Atomic transitions.** Derive the change from live state under the appropriate lock and apply a table's policy transition in one bounded transaction. Replacement must not leave a committed intermediate access rule. Refuse mixed table/policy plans until their execution strategy preserves that guarantee. @@ -217,11 +218,11 @@ this table-scoped work. incomplete export. Implemented here; table-only diff behavior stays intact. 2. **Round-trip the declaration.** Admit and export SQL under explicit RLS scope; materialize it in scratch and prove the unchanged definition produces an empty - diff. Implemented through automatic export and SQL declarations; execution remains refused, - including greenfield creation. -3. **Plan and execute transitions.** Add typed security changes, exact-state - revalidation, lock budgets, atomic application, dependency handling, and reports - suitable for users and orchestrators. + diff. Implemented through automatic export and SQL declarations; greenfield + creation remains refused. +3. **Execute transitions atomically.** The Go executor now locks, derives, applies, + and verifies one table's RLS state in one bounded transaction. Mixed changes and + policy relation dependencies remain unsupported. CLI integration is a follow-up. 4. **Prove application behavior.** Extend the local Supabase harness with real authenticated and anonymous requests, two users, allowed and denied writes, and interrupted transitions. Then validate hosted connection and privilege @@ -231,7 +232,9 @@ The [inspection tests](../pkg/schemadiff/row_security_integration_test.go), [round-trip tests](../pkg/schemadiff/row_security_roundtrip_integration_test.go), and [Supabase auth test](../integration/supabase/row_security_test.go) use real databases and readable DDL. They prove catalog fidelity, round trips, and refusal boundaries, -**not support for applying policies**. The existing PostgreSQL CI matrix and +**not support for applying policies**. The separate +[executor tests](../pkg/executor/row_security_integration_test.go) cover atomic +application, rollback, lock waits, and non-owner access on PostgreSQL. The existing PostgreSQL CI matrix and Supabase compatibility job both run `pkg/schemadiff`; no separate runner is needed. Hosted validation is not a prerequisite for the local steps, nor replaced by them. diff --git a/docs/execution-model.md b/docs/execution-model.md index fae5169..5d2469d 100644 --- a/docs/execution-model.md +++ b/docs/execution-model.md @@ -289,6 +289,7 @@ concurrently is. | --- | --- | --- | | `budget-lock-exceeded` | no | The lock was not granted within `lock_timeout`; nothing executed | | `budget-statement-exceeded` | no | The statement ran past `statement_timeout` and was cancelled | +| `row-security-outcome-unknown` | no | The atomic RLS commit response is uncertain; inspect the catalog before retrying | | `blocking-outcome-unknown` | no | The accepted blocking transaction reached an ambiguous client boundary; inspect the catalog before retrying | | `invalid-blocking-budget` | yes | An accepted blocking bound is disabled or cannot be represented by PostgreSQL | | `unsupported-accepted-blocking` | yes | The statement is outside the accepted blocking executor's narrow index-maintenance set, or the server will not run it inside the engine-owned transaction (`REINDEX` on a partitioned relation, SQLSTATE `25001`) | diff --git a/docs/invariants.md b/docs/invariants.md index 9c7b086..49f6810 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -22,6 +22,7 @@ several of these unrepresentable, and the in-TCB engineering rules live in - [Correctness (CO)](#correctness-co) - [Locking and concurrency (LK)](#locking-and-concurrency-lk) - [Accepted blocking execution (AB)](#accepted-blocking-execution-ab) +- [Atomic row security (RS)](#atomic-row-security-rs) - [State, checkpoint, and resume (ST)](#state-checkpoint-and-resume-st) - [Refusals and preflight (RF)](#refusals-and-preflight-rf) - [Orchestration / control-plane (OC)](#orchestration--control-plane-oc) @@ -391,6 +392,17 @@ Every acceptance is audited at warn level before execution, regardless of `--deb `TestMigrateAcceptBlockingRunsDropIndex`, `demo/tour.sh` (`execute_accepted`). *Source:* [lock-budgeted passthrough](lock-budgeted-passthrough.md#exit-codes). +## Atomic row security (RS) + +The [atomic RLS contract](atomic-row-security.md) defines these executor obligations: + +| ID | Must hold | Enforcement and test obligation | +| --- | --- | --- | +| RS-1 | Lock the live target before deriving the change; refuse mixed table changes | `ExecuteRowSecurity`; contention and mixed-change tests | +| RS-2 | Policy changes and convergence verification commit together or roll back | `ExecuteRowSecurity`; real DDL fault injection restores original policies | +| RS-3 | Bound lock waits, statements, and the whole attempt | `ExecuteRowSecurity`; lock contention and deadline cancellation | +| RS-4 | Run only admitted, qualified RLS DDL; keep scratch disposable | `SecurityStatements` and executor readback; quoted names and scratch-cleanup tests | + ## State, checkpoint, and resume (ST) ### ST-1 — The checkpoint is one row per target, written atomically diff --git a/docs/limitations.md b/docs/limitations.md index 11e0105..ec8e4a7 100644 --- a/docs/limitations.md +++ b/docs/limitations.md @@ -25,7 +25,7 @@ with a typed refusal — never a silently wrong or incomplete result: | Not modeled | Current behavior | | --- | --- | -| Row-level security (RLS) settings and policies | `pull` exports settings, policies, roles, and policy comments; `diff` verifies unchanged definitions and shows review-only RLS deltas before refusing changes. Policy relation subqueries, any changed definition, and greenfield creation in this mode are refused. Tables without RLS get no extra SQL; table-only desired files leave access control separately managed. See the [RLS workflow and roadmap](declarative-row-security.md). | +| Row-level security (RLS) settings and policies | `pull` exports settings, policies, roles, and policy comments; `diff` verifies unchanged definitions and shows review-only RLS deltas before refusing changes. The [atomic Go executor](atomic-row-security.md) supports RLS-only changes to existing supported tables. Policy relation subqueries, mixed table changes, greenfield creation, and CLI policy execution remain refused. Tables without RLS get no extra SQL; table-only desired files leave access control separately managed. See the [RLS workflow and roadmap](declarative-row-security.md). | | Foreign keys (either direction) | Unsupported in the declarative model, in both directions. A desired file cannot declare a `REFERENCES` clause (refused at parse), and export refuses both sides of a foreign-key relationship: a table whose definition carries foreign-key constraints surfaces the parse gate's typed error, and a table that other tables reference refuses with its own typed error — a single-table baseline cannot carry incoming foreign-key topology, so rendering one would silently drop the relationship. Foreign-key **DDL is still supported through the statement front door** (`ADD FOREIGN KEY` routes to the online `NOT VALID` + `VALIDATE` sequence). Because a desired file can never declare a foreign key, `diff` on a live table that carries one plans a **destructive** `DROP CONSTRAINT` for it — gated like every destructive change, never auto-executed — and incoming foreign keys are invisible to a single-table diff entirely, so tables participating in foreign-key relationships in either direction should not be managed declaratively yet. | | Partitioned tables | Partitioned parents and their partitions are introspectable, and the statement front door supports in-place changes on them (see the partitioned-parent rows above), but they cannot be expressed in or exported to a desired file: the model captures the partition key and attachment only to refuse — it does not carry partition bounds or the parent/partition topology. A partitioning mismatch between live and desired is a typed `diff` refusal, never a zero diff. | | Classic table inheritance (`INHERITS`) | Both parents and children are refused during export. The model records classic inheritance edges in both directions only to fail closed: rendering a child would flatten inherited columns and rendering a parent would omit its children. Declarative partitions share `pg_inherits` but remain classified by `relispartition` and use the partition refusal above. | diff --git a/docs/tcb-model.md b/docs/tcb-model.md index 12acb94..49ac8e3 100644 --- a/docs/tcb-model.md +++ b/docs/tcb-model.md @@ -91,7 +91,7 @@ to obtain the type is through the function that validates it. | --- | --- | --- | --- | | `string` (one SQL statement) | `statement.ParseOne` | `statement.Statement` (carries the text, kind, and target relation of the one statement the real grammar admitted; the executor accepts nothing else) | ST-7 | | `string` (desired-state schema file) | `statement.ParseDesired` | `statement.DesiredSchema` (carries the single CREATE TABLE and its indexes in execution order; every consumer replays that order) | ST-8 | -| `string` (complete table-local RLS declaration) | `statement.ParseDesiredWithRowSecurity` | `statement.DesiredWithRowSecurity` (scratch inspection only; cannot enter the existing live executors) | Explicit RLS scope, including an empty policy set | +| `string` (complete table-local RLS declaration) | `statement.ParseDesiredWithRowSecurity` | `statement.DesiredWithRowSecurity` (declaration syntax; only the dedicated atomic RLS executor may apply it) | Explicit RLS scope, including an empty policy set | | `string` (user SQL) | `statement.ParseOne` / `statement.ParseOps`, then `planner.Classify` | `planner.Plan` / `planner.Decision` | CO-7 — classification consumes parsed operation descriptors | | table name | preflight | `PreflightedTable` (carries the proven facts: PK, no FKs/views, replica identity, headroom) | ST-6, RF-* | | table name (create target) | `preflight.CheckTableAbsent` | `AbsentTarget` (carries the resolved creation schema and the verified-free name; time-of-check — minted inside the apply session, never carried across a plan boundary, and re-verified at use the way ST-7 re-verifies `PreflightedTable`) | ST-6 for the create path | diff --git a/pkg/capabilities/capabilities.yaml b/pkg/capabilities/capabilities.yaml index 1209b85..89679df 100644 --- a/pkg/capabilities/capabilities.yaml +++ b/pkg/capabilities/capabilities.yaml @@ -509,7 +509,7 @@ rows: reason_notes: "Same: catalog bootstrap, owner tooling" - id: grants-roles-row-level-security-policies area: "types_and_non_table_objects" - operation: "Grants, roles, row-level-security policies" + operation: "Grants, roles, and CLI policy DDL" tier: "t3" status_mark: "🔵" engine_path: "none" @@ -519,7 +519,19 @@ rows: migrate: refused diff: refused refusal_reason: unsupported-statement - reason_notes: "Access control changes remain with provisioning. Export includes RLS when present; explicit RLS declarations support comparison and review-only deltas with advisory access warnings. Plain table files leave access control separately managed. RLS execution remains unsupported. See the [workflow and roadmap](declarative-row-security.md). See [engine-role.md](engine-role.md) for the engine's own role" + reason_notes: "Access control changes remain with provisioning. Export includes RLS when present; explicit RLS declarations support comparison and review-only deltas with advisory access warnings. Plain table files leave access control separately managed. The separate [atomic RLS Go executor](atomic-row-security.md) supports RLS-only changes on existing supported tables; it does not route through these CLI front doors. See the [workflow and roadmap](declarative-row-security.md). See [engine-role.md](engine-role.md) for the engine's own role" + - id: declarative-row-security-library + area: "types_and_non_table_objects" + operation: "Complete table-local RLS definition (Go API only)" + tier: "t1" + status_mark: "✅" + engine_path: "native_as_is" + online_safety_problem: true + front_doors: + migrate: refused + diff: refused + refusal_reason: unsupported-statement + reason_notes: "The [atomic RLS executor](atomic-row-security.md) locks an existing supported table, refuses structural changes, replaces policies/settings, and verifies convergence before committing. Lock waits and the whole transaction are bounded. No CLI execution or new diff flags" - id: standalone-sequences area: "types_and_non_table_objects" operation: "Standalone sequences" diff --git a/pkg/capabilities/capabilities_test.go b/pkg/capabilities/capabilities_test.go index 29b1f07..9502e86 100644 --- a/pkg/capabilities/capabilities_test.go +++ b/pkg/capabilities/capabilities_test.go @@ -21,7 +21,7 @@ func validRow() Row { func TestEmbeddedCapabilitiesValidate(t *testing.T) { rows, err := Rows() require.NoError(t, err) - assert.Len(t, rows, 53) + assert.Len(t, rows, 54) } // The engine refuses a partitioned parent with its own target-fact reason @@ -93,7 +93,7 @@ func TestValidateRules(t *testing.T) { } // An enum error names the row, the field, and the offending value, so a -// typo in a 53-row file is found without diffing the vocabulary by hand. +// typo in the capability file is found without diffing the vocabulary by hand. func TestValidateNamesTheFieldAndValue(t *testing.T) { tests := map[string]struct { mutate func(*Row) diff --git a/pkg/executor/code.go b/pkg/executor/code.go index cb84478..898c900 100644 --- a/pkg/executor/code.go +++ b/pkg/executor/code.go @@ -23,6 +23,8 @@ const ( // CodeBudgetStatementExceeded: the statement ran past // statement_timeout and was cancelled; the change does real work. CodeBudgetStatementExceeded Code = "budget-statement-exceeded" + // CodeRowSecurityOutcomeUnknown requires catalog inspection before retry. + CodeRowSecurityOutcomeUnknown Code = "row-security-outcome-unknown" // CodeBlockingOutcomeUnknown requires catalog inspection before retry. CodeBlockingOutcomeUnknown Code = "blocking-outcome-unknown" // CodeInvalidBlockingBudget identifies an unrepresentable or disabled bound. @@ -148,6 +150,7 @@ func Codes() []Code { CodeBudgetLockExceeded, CodeBudgetStatementExceeded, CodeBlockingOutcomeUnknown, + CodeRowSecurityOutcomeUnknown, CodeInvalidBlockingBudget, CodeUnsupportedAcceptedBlocking, CodeCancelledByCaller, @@ -242,6 +245,10 @@ func OutcomeCode(err error) Code { if errors.As(err, &budgetErr) { return budgetErr.Code() } + var rlsUnknown *RowSecurityOutcomeUnknownError + if errors.As(err, &rlsUnknown) { + return CodeRowSecurityOutcomeUnknown + } var unknownErr *BlockingOutcomeUnknownError if errors.As(err, &unknownErr) { return CodeBlockingOutcomeUnknown diff --git a/pkg/executor/code_test.go b/pkg/executor/code_test.go index 5c2e4ec..47c7d36 100644 --- a/pkg/executor/code_test.go +++ b/pkg/executor/code_test.go @@ -24,6 +24,7 @@ func TestOutcomeCodeMapsTypedOutcomes(t *testing.T) { err error want executor.Code }{ + {name: "row security unknown commit", err: &executor.RowSecurityOutcomeUnknownError{Err: errors.New("connection lost")}, want: executor.CodeRowSecurityOutcomeUnknown}, {name: "nil error has no code", err: nil, want: executor.Code("")}, { name: "lock budget", @@ -162,6 +163,7 @@ func TestCodePermanentClassifiesEveryCode(t *testing.T) { executor.CodeBudgetLockExceeded: false, executor.CodeBudgetStatementExceeded: false, executor.CodeBlockingOutcomeUnknown: false, + executor.CodeRowSecurityOutcomeUnknown: false, executor.CodeInvalidBlockingBudget: true, executor.CodeUnsupportedAcceptedBlocking: true, executor.CodeCancelledByCaller: false, diff --git a/pkg/executor/row_security.go b/pkg/executor/row_security.go new file mode 100644 index 0000000..ac7d636 --- /dev/null +++ b/pkg/executor/row_security.go @@ -0,0 +1,138 @@ +package executor + +import ( + "context" + "fmt" + "strconv" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/schemadiff" + "github.com/block/pg-sprite/pkg/statement" +) + +// RowSecurityReport contains only committed RLS statements, in execution order. +// A no-op has an empty Statements slice. Errors return a zero report. +type RowSecurityReport struct { + Schema string `json:"schema"` + Table string `json:"table"` + Statements []string `json:"statements"` +} + +// RowSecurityOutcomeUnknownError means the commit response was lost or failed. +// Inspect the catalog before retrying; an error does not prove rollback here. +type RowSecurityOutcomeUnknownError struct{ Err error } + +// Error describes the uncertain commit boundary. +func (e *RowSecurityOutcomeUnknownError) Error() string { + return fmt.Sprintf("row security commit outcome is unknown; inspect the catalog: %v", e.Err) +} + +// Unwrap retains the underlying commit error. +func (e *RowSecurityOutcomeUnknownError) Unwrap() error { return e.Err } + +// ExecuteRowSecurity converges only the complete RLS definition of an existing +// ordinary table. It locks and inspects the target itself and refuses table +// deltas. Budget.StatementTimeout also bounds the entire attempt, including +// connection acquisition, scratch inspection, and lock waits. No retries occur. +// The caller must have table-owner and scratch-schema creation privileges. +func ExecuteRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, desired statement.DesiredWithRowSecurity, b Budget) (RowSecurityReport, error) { + if err := b.validate(); err != nil { + return RowSecurityReport{}, err + } + security, err := desired.SecurityStatements(schema) + if err != nil { + return RowSecurityReport{}, err + } + // INV: RS-3 — the whole attempt has one deadline, not a fresh budget per policy. + ctx, cancel := context.WithTimeout(ctx, b.StatementTimeout) + defer cancel() + tx, err := pool.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.ReadCommitted}) + if err != nil { + return RowSecurityReport{}, fmt.Errorf("begin row security change: %w", err) + } + defer func() { + cleanup, stop := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) + defer stop() + _ = tx.Rollback(cleanup) + }() + settings := "SET LOCAL lock_timeout = " + strconv.FormatInt(b.LockTimeout.Milliseconds(), 10) + "; SET LOCAL statement_timeout = " + strconv.FormatInt(b.StatementTimeout.Milliseconds(), 10) + if _, err := tx.Exec(ctx, settings); err != nil { + return RowSecurityReport{}, fmt.Errorf("set row security budgets: %w", err) + } + target := pgx.Identifier{schema, desired.Table()}.Sanitize() + // INV: RS-1 — all live comparison and DDL occur after this exclusive lock. + if _, err := tx.Exec(ctx, "LOCK TABLE ONLY "+target+" IN ACCESS EXCLUSIVE MODE"); err != nil { + return RowSecurityReport{}, fmt.Errorf("lock row security target %s: %w", target, err) + } + live, err := schemadiff.IntrospectTx(ctx, tx, schema, desired.Table()) + if err != nil { + return RowSecurityReport{}, err + } + wanted, err := schemadiff.IntrospectDesiredWithRowSecurityTx(ctx, tx, desired) + if err != nil { + return RowSecurityReport{}, err + } + if err := admitRowSecurityTable(schema, live, wanted); err != nil { + return RowSecurityReport{}, err + } + report := RowSecurityReport{Schema: schema, Table: desired.Table(), Statements: []string{}} + if _, err := schemadiff.DiffWithRowSecurity(schema, live, wanted); err != nil { + // INV: RS-4 — helpers are explicitly qualified; no target-schema function + // may shadow a built-in while replaying policy expressions. + if _, err := tx.Exec(ctx, dbconn.LocalSearchPath("pg_catalog")); err != nil { + return RowSecurityReport{}, err + } + for _, policy := range live.RowSecurity.Policies { + report.Statements = append(report.Statements, "DROP POLICY "+pgx.Identifier{policy.Name}.Sanitize()+" ON "+target) + } + report.Statements = append(report.Statements, security...) + // An omitted FORCE declaration means NO FORCE, not "leave it alone". + force := "NO FORCE" + if wanted.RowSecurity.Forced { + force = "FORCE" + } + report.Statements = append(report.Statements, "ALTER TABLE "+target+" "+force+" ROW LEVEL SECURITY") + for _, sql := range report.Statements { + if _, err := tx.Exec(ctx, sql); err != nil { + return RowSecurityReport{}, fmt.Errorf("apply row security on %s: %w", target, err) + } + } + } + // INV: RS-2 — verify convergence before committing, while retaining the lock. + actual, err := schemadiff.IntrospectTx(ctx, tx, schema, desired.Table()) + if err != nil { + return RowSecurityReport{}, err + } + if err := admitRowSecurityTable(schema, actual, wanted); err != nil { + return RowSecurityReport{}, fmt.Errorf("%w: RS-2: target shape changed: %w", ErrInvariantViolation, err) + } + if _, err := schemadiff.DiffWithRowSecurity(schema, actual, wanted); err != nil { + return RowSecurityReport{}, fmt.Errorf("%w: RS-2: row security did not converge: %w", ErrInvariantViolation, err) + } + if err := tx.Commit(ctx); err != nil { + return RowSecurityReport{}, &RowSecurityOutcomeUnknownError{Err: err} + } + return report, nil +} + +func admitRowSecurityTable(schema string, live, wanted schemadiff.Model) error { + // Render refuses unsupported shapes even when two models happen to match. + table := live + table.RowSecurity = schemadiff.RowSecurity{} + if _, err := schemadiff.Render(table); err != nil { + return fmt.Errorf("row security target: %w", err) + } + wanted.RowSecurity = schemadiff.RowSecurity{} + changes, err := schemadiff.Diff(schema, table, wanted) + if err != nil { + return err + } + if len(changes) != 0 { + return fmt.Errorf("mixed table and row security changes: %w", schemadiff.ErrUnsupportedChange) + } + return nil +} diff --git a/pkg/executor/row_security_integration_test.go b/pkg/executor/row_security_integration_test.go new file mode 100644 index 0000000..4bd94d5 --- /dev/null +++ b/pkg/executor/row_security_integration_test.go @@ -0,0 +1,330 @@ +package executor_test + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/internal/testutil" + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/executor" + "github.com/block/pg-sprite/pkg/schemadiff" + "github.com/block/pg-sprite/pkg/statement" +) + +func rlsFixture(t *testing.T) (*pgxpool.Pool, string) { + t.Helper() + pool, err := dbconn.NewPool(t.Context(), dbconn.Config{URL: testutil.StartPostgres(t)}) + require.NoError(t, err) + t.Cleanup(pool.Close) + schema := testutil.NewSchema(t, pool) + _, err = pool.Exec(t.Context(), fmt.Sprintf(`CREATE TABLE %s.documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE %s.documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON %s.documents FOR SELECT USING (owner_id = 7);`, schema, schema, schema)) + require.NoError(t, err) + return pool, schema +} + +func applyRLS(t *testing.T, pool *pgxpool.Pool, schema, sql string) (executor.RowSecurityReport, error) { + t.Helper() + desired, err := statement.ParseDesiredWithRowSecurity(sql) + require.NoError(t, err) + return executor.ExecuteRowSecurity(t.Context(), pool, schema, desired, executor.Budget{LockTimeout: 100 * time.Millisecond, StatementTimeout: 5 * time.Second}) +} + +func TestExecuteRowSecurityReplacesPoliciesAtomically(t *testing.T) { + pool, schema := rlsFixture(t) + sql := `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + ALTER TABLE documents FORCE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (owner_id = 9); + CREATE POLICY writers ON documents FOR INSERT WITH CHECK (owner_id = 9); + COMMENT ON POLICY readers ON documents IS 'Only your documents';` + report, err := applyRLS(t, pool, schema, sql) + require.NoError(t, err) + assert.NotEmpty(t, report.Statements) + actual, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.True(t, actual.RowSecurity.Enabled) + assert.True(t, actual.RowSecurity.Forced) + require.Len(t, actual.RowSecurity.Policies, 2) + assert.Equal(t, "(owner_id = 9)", *actual.RowSecurity.Policies[0].Using) + assert.Equal(t, "Only your documents", *actual.RowSecurity.Policies[0].Comment) + assert.Equal(t, "(owner_id = 9)", *actual.RowSecurity.Policies[1].WithCheck) + again, err := applyRLS(t, pool, schema, sql) + require.NoError(t, err) + assert.Empty(t, again.Statements, "a converged retry must do no live DDL") +} + +func TestExecuteRowSecurityRemovesLastPolicyAndForce(t *testing.T) { + pool, schema := rlsFixture(t) + _, err := pool.Exec(t.Context(), fmt.Sprintf(`ALTER TABLE %s.documents FORCE ROW LEVEL SECURITY;`, schema)) + require.NoError(t, err) + _, err = applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents DISABLE ROW LEVEL SECURITY;`) + require.NoError(t, err) + actual, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.False(t, actual.RowSecurity.Enabled) + assert.False(t, actual.RowSecurity.Forced) + assert.Empty(t, actual.RowSecurity.Policies) +} + +func TestExecuteRowSecurityRefusesMixedChanges(t *testing.T) { + pool, schema := rlsFixture(t) + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + report, err := applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL, + title text + ); + ALTER TABLE documents DISABLE ROW LEVEL SECURITY;`) + require.ErrorIs(t, err, schemadiff.ErrUnsupportedChange) + assert.Empty(t, report.Statements) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} + +func TestExecuteRowSecurityLockContentionLeavesPoliciesIntact(t *testing.T) { + pool, schema := rlsFixture(t) + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + tx, err := pool.Begin(t.Context()) + require.NoError(t, err) + defer func() { _ = tx.Rollback(context.WithoutCancel(t.Context())) }() + _, err = tx.Exec(t.Context(), "LOCK TABLE "+pgx.Identifier{schema, "documents"}.Sanitize()+" IN ACCESS SHARE MODE") + require.NoError(t, err) + _, err = applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents DISABLE ROW LEVEL SECURITY;`) + var pgErr *pgconn.PgError + require.ErrorAs(t, err, &pgErr) + assert.Equal(t, "55P03", pgErr.Code) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} + +func TestExecuteRowSecurityRollsBackAfterLiveDDLFailure(t *testing.T) { + pool, schema := rlsFixture(t) + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + // A real server-side failure after DROP POLICY and CREATE POLICY proves that + // the transaction restores the original policy, rather than just refusing early. + trigger := pgx.Identifier{schema + "_fail_policy"}.Sanitize() + _, err = pool.Exec(t.Context(), fmt.Sprintf(`CREATE FUNCTION %s.fail_policy() RETURNS event_trigger + LANGUAGE plpgsql AS $$ + BEGIN + IF EXISTS ( + SELECT 1 FROM pg_event_trigger_ddl_commands() d + JOIN pg_policy p ON d.classid = 'pg_policy'::regclass AND d.objid = p.oid + JOIN pg_class c ON c.oid = p.polrelid + JOIN pg_namespace n ON n.oid = c.relnamespace + WHERE n.nspname = '%s' AND d.command_tag = 'CREATE POLICY' + ) THEN + RAISE EXCEPTION 'injected policy failure'; + END IF; + END; + $$; + CREATE EVENT TRIGGER %s ON ddl_command_end EXECUTE FUNCTION %s.fail_policy();`, schema, schema, trigger, schema)) + require.NoError(t, err) + t.Cleanup(func() { + _, err := pool.Exec(context.WithoutCancel(t.Context()), "DROP EVENT TRIGGER "+trigger) + require.NoError(t, err) + }) + report, err := applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents DISABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (owner_id = 9);`) + var pgErr *pgconn.PgError + require.ErrorAs(t, err, &pgErr) + assert.Equal(t, "P0001", pgErr.Code) + assert.Empty(t, report.Statements) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} + +func TestExecuteRowSecurityEnabledWithoutPoliciesDeniesRows(t *testing.T) { + pool, schema := rlsFixture(t) + _, err := pool.Exec(t.Context(), fmt.Sprintf(`INSERT INTO %s.documents VALUES (1, 7), (2, 9);`, schema)) + require.NoError(t, err) + _, err = applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY;`) + require.NoError(t, err) + // Real non-owner role: catalog convergence alone cannot prove default deny. + role := pgx.Identifier{schema + "_reader"}.Sanitize() + _, err = pool.Exec(t.Context(), "CREATE ROLE "+role) + require.NoError(t, err) + t.Cleanup(func() { + ctx := context.WithoutCancel(t.Context()) + _, err := pool.Exec(ctx, "DROP OWNED BY "+role+"; DROP ROLE "+role) + require.NoError(t, err) + }) + _, err = pool.Exec(t.Context(), fmt.Sprintf("GRANT USAGE ON SCHEMA %s TO %s; GRANT SELECT ON %s.documents TO %s", schema, role, schema, role)) + require.NoError(t, err) + tx, err := pool.Begin(t.Context()) + require.NoError(t, err) + defer func() { _ = tx.Rollback(context.WithoutCancel(t.Context())) }() + _, err = tx.Exec(t.Context(), "SET LOCAL ROLE "+role) + require.NoError(t, err) + var count int + require.NoError(t, tx.QueryRow(t.Context(), "SELECT count(*) FROM "+pgx.Identifier{schema, "documents"}.Sanitize()).Scan(&count)) + assert.Zero(t, count) +} + +func TestExecuteRowSecurityRefusesMissingTable(t *testing.T) { + pool, schema := rlsFixture(t) + _, err := applyRLS(t, pool, schema, `CREATE TABLE missing ( + id bigint PRIMARY KEY + ); + ALTER TABLE missing ENABLE ROW LEVEL SECURITY;`) + require.Error(t, err) + assert.Empty(t, relationKind(t, pool, schema, "missing")) +} + +func TestExecuteRowSecurityReadsStateAfterWaitingForLock(t *testing.T) { + pool, schema := rlsFixture(t) + desired, err := statement.ParseDesiredWithRowSecurity(`CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (owner_id = 9);`) + require.NoError(t, err) + holder, err := pool.Begin(t.Context()) + require.NoError(t, err) + defer func() { _ = holder.Rollback(context.WithoutCancel(t.Context())) }() + target := pgx.Identifier{schema, "documents"}.Sanitize() + _, err = holder.Exec(t.Context(), "LOCK TABLE "+target+" IN ACCESS EXCLUSIVE MODE") + require.NoError(t, err) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + done := make(chan error, 1) + go func() { + _, err := executor.ExecuteRowSecurity(ctx, pool, schema, desired, executor.Budget{LockTimeout: 3 * time.Second, StatementTimeout: 5 * time.Second}) + done <- err + }() + // Observe the actual lock queue, not a timing guess. The earlier policy read + // would miss this new policy and fail convergence instead of applying cleanly. + const queueDeadline = 2 * time.Second + const queuePoll = 10 * time.Millisecond + require.Eventually(t, func() bool { + var waiting bool + err := pool.QueryRow(t.Context(), `SELECT EXISTS ( + SELECT 1 FROM pg_locks WHERE relation = $1::regclass + AND mode = 'AccessExclusiveLock' AND NOT granted + )`, target).Scan(&waiting) + return err == nil && waiting + }, queueDeadline, queuePoll) + _, err = holder.Exec(t.Context(), "CREATE POLICY concurrent_reader ON "+target+" FOR SELECT USING (owner_id = 11)") + require.NoError(t, err) + require.NoError(t, holder.Commit(t.Context())) + require.NoError(t, <-done) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + require.Len(t, after.RowSecurity.Policies, 1) + assert.Equal(t, "readers", after.RowSecurity.Policies[0].Name) + assert.Equal(t, "(owner_id = 9)", *after.RowSecurity.Policies[0].Using) +} + +func TestExecuteRowSecurityUsesOneConnectionAndQuotedNames(t *testing.T) { + pool, err := dbconn.NewPool(t.Context(), dbconn.Config{URL: testutil.StartPostgres(t), MaxConns: 1}) + require.NoError(t, err) + t.Cleanup(pool.Close) + schema := testutil.NewSchema(t, pool) + _, err = pool.Exec(t.Context(), fmt.Sprintf(`CREATE TABLE %s."odd""table" ( + id bigint PRIMARY KEY + );`, schema)) + require.NoError(t, err) + desired, err := statement.ParseDesiredWithRowSecurity(`CREATE TABLE "odd""table" ( + id bigint PRIMARY KEY + ); + ALTER TABLE "odd""table" ENABLE ROW LEVEL SECURITY; + CREATE POLICY "read""only" ON "odd""table" FOR SELECT USING (id = 7); + COMMENT ON POLICY "read""only" ON "odd""table" IS 'a quote: '' and semicolon;';`) + require.NoError(t, err) + _, err = executor.ExecuteRowSecurity(t.Context(), pool, schema, desired, executor.Budget{LockTimeout: time.Second, StatementTimeout: 5 * time.Second}) + require.NoError(t, err) + actual, err := schemadiff.Introspect(t.Context(), pool, schema, `odd"table`) + require.NoError(t, err) + require.Len(t, actual.RowSecurity.Policies, 1) + assert.Equal(t, `read"only`, actual.RowSecurity.Policies[0].Name) + assert.Equal(t, "a quote: ' and semicolon;", *actual.RowSecurity.Policies[0].Comment) + var scratch int + require.NoError(t, pool.QueryRow(t.Context(), `SELECT count(*) FROM pg_namespace WHERE nspname LIKE 'pgsprite_scratch_%'`).Scan(&scratch)) + assert.Zero(t, scratch, "savepoints must not commit scratch objects with the live change") +} + +func TestExecuteRowSecurityDeadlineRollsBackLiveDDL(t *testing.T) { + pool, schema := rlsFixture(t) + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + // Stall after live policy creation to exhaust the whole attempt budget. + // The policy drops and additions must all roll back on cancellation. + _, err = pool.Exec(t.Context(), "CREATE SEQUENCE "+pgx.Identifier{schema, "fault_reached"}.Sanitize()) + require.NoError(t, err) + trigger := pgx.Identifier{schema + "_fail_policy"}.Sanitize() + _, err = pool.Exec(t.Context(), fmt.Sprintf(`CREATE FUNCTION %s.fail_policy() RETURNS event_trigger + LANGUAGE plpgsql AS $$ + BEGIN + IF EXISTS ( + SELECT 1 FROM pg_event_trigger_ddl_commands() d + JOIN pg_policy p ON d.classid = 'pg_policy'::regclass AND d.objid = p.oid + JOIN pg_class c ON c.oid = p.polrelid + JOIN pg_namespace n ON n.oid = c.relnamespace + WHERE n.nspname = '%s' AND d.command_tag = 'CREATE POLICY' + ) THEN + PERFORM nextval('%s.fault_reached'); + PERFORM pg_sleep(10); + END IF; + END; + $$; + CREATE EVENT TRIGGER %s ON ddl_command_end EXECUTE FUNCTION %s.fail_policy();`, schema, schema, schema, trigger, schema)) + require.NoError(t, err) + t.Cleanup(func() { + _, err := pool.Exec(context.WithoutCancel(t.Context()), "DROP EVENT TRIGGER "+trigger) + require.NoError(t, err) + }) + desired, err := statement.ParseDesiredWithRowSecurity(`CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents DISABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (owner_id = 9);`) + require.NoError(t, err) + report, err := executor.ExecuteRowSecurity(t.Context(), pool, schema, desired, executor.Budget{LockTimeout: 50 * time.Millisecond, StatementTimeout: 2 * time.Second}) + require.Error(t, err) + assert.Empty(t, report.Statements) + var reached bool + require.NoError(t, pool.QueryRow(t.Context(), "SELECT is_called FROM "+pgx.Identifier{schema, "fault_reached"}.Sanitize()).Scan(&reached)) + assert.True(t, reached, "the deadline must fire after live DDL, not during setup") + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} diff --git a/pkg/executor/row_security_test.go b/pkg/executor/row_security_test.go new file mode 100644 index 0000000..bffaf58 --- /dev/null +++ b/pkg/executor/row_security_test.go @@ -0,0 +1,25 @@ +package executor_test + +import ( + "testing" + "time" + + "github.com/block/pg-sprite/pkg/executor" + "github.com/block/pg-sprite/pkg/statement" + "github.com/stretchr/testify/require" +) + +func TestExecuteRowSecurityRejectsInvalidInputsBeforeConnecting(t *testing.T) { + desired, err := statement.ParseDesiredWithRowSecurity(`CREATE TABLE documents ( + id bigint PRIMARY KEY + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY;`) + require.NoError(t, err) + _, err = executor.ExecuteRowSecurity(t.Context(), nil, "public", desired, executor.Budget{}) + require.Error(t, err) + budget := executor.Budget{LockTimeout: time.Millisecond, StatementTimeout: time.Second} + _, err = executor.ExecuteRowSecurity(t.Context(), nil, "public", statement.DesiredWithRowSecurity{}, budget) + require.ErrorIs(t, err, statement.ErrRowSecurityDeclaration) + _, err = executor.ExecuteRowSecurity(t.Context(), nil, "", desired, budget) + require.ErrorIs(t, err, statement.ErrRowSecurityDeclaration) +} diff --git a/pkg/schemadiff/desired.go b/pkg/schemadiff/desired.go index fc58088..ec8302e 100644 --- a/pkg/schemadiff/desired.go +++ b/pkg/schemadiff/desired.go @@ -5,6 +5,7 @@ import ( "crypto/rand" "encoding/hex" "fmt" + "time" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" @@ -25,22 +26,31 @@ func IntrospectDesired(ctx context.Context, db *pgxpool.Pool, desired statement. } func introspectDesiredStatements(ctx context.Context, db *pgxpool.Pool, table string, statements []statement.Statement) (Model, error) { - scratch, err := scratchSchemaName() - if err != nil { - return Model{}, err - } tx, err := db.Begin(ctx) if err != nil { return Model{}, fmt.Errorf("begin scratch transaction: %w", err) } + return introspectDesiredTransaction(ctx, tx, table, statements) +} + +// introspectDesiredTransaction owns tx, which may be a savepoint. Every exit +// rolls it back, leaving its parent's live target and settings untouched. +func introspectDesiredTransaction(ctx context.Context, tx pgx.Tx, table string, statements []statement.Statement) (Model, error) { // The scratch transaction is never committed: rollback is the cleanup // path for success and failure alike, so the redundant-closer exception // does not apply — this rollback is load-bearing and its error is // surfaced on the success path below. defer func() { - _ = tx.Rollback(context.WithoutCancel(ctx)) + cleanup, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) + defer cancel() + _ = tx.Rollback(cleanup) }() + scratch, err := scratchSchemaName() + if err != nil { + return Model{}, err + } + if _, err := tx.Exec(ctx, "CREATE SCHEMA "+pgx.Identifier{scratch}.Sanitize()); err != nil { return Model{}, fmt.Errorf("create scratch schema: %w", err) } diff --git a/pkg/schemadiff/row_security_roundtrip.go b/pkg/schemadiff/row_security_roundtrip.go index ab4f471..fa2c42a 100644 --- a/pkg/schemadiff/row_security_roundtrip.go +++ b/pkg/schemadiff/row_security_roundtrip.go @@ -22,6 +22,19 @@ func IntrospectDesiredWithRowSecurity(ctx context.Context, db *pgxpool.Pool, des return introspectDesiredStatements(ctx, db, desired.Table(), desired.Statements()) } +// IntrospectDesiredWithRowSecurityTx materializes the declaration in a nested +// transaction (savepoint) and always rolls it back. The parent's locks survive. +func IntrospectDesiredWithRowSecurityTx(ctx context.Context, tx pgx.Tx, desired statement.DesiredWithRowSecurity) (Model, error) { + if desired.Table() == "" { + return Model{}, statement.ErrRowSecurityDeclaration + } + scratch, err := tx.Begin(ctx) + if err != nil { + return Model{}, fmt.Errorf("begin desired row security savepoint: %w", err) + } + return introspectDesiredTransaction(ctx, scratch, desired.Table(), desired.Statements()) +} + // RenderWithRowSecurity exports the complete table-local RLS definition, // including explicit DISABLE when RLS is off. The result is admitted only by // ParseDesiredWithRowSecurity and cannot be passed to the live create executor. @@ -107,16 +120,16 @@ func renderPolicy(b *strings.Builder, target string, p Policy) error { return nil } -// DiffWithRowSecurity checks a complete RLS-scoped declaration. Only unchanged -// state is admitted in this phase: security deltas and mixed table/security -// changes are refused as ErrUnsupportedChange. No executable policy SQL is -// produced. Table-only callers retain the independent Diff contract. +// DiffWithRowSecurity checks equality of a complete RLS-scoped declaration. +// Security deltas and mixed table/security changes return ErrUnsupportedChange; +// this comparison never emits executable policy SQL. The dedicated atomic +// executor uses it to verify convergence. Table-only callers retain Diff. func DiffWithRowSecurity(schema string, live, desired Model) ([]Change, error) { if live.Table != desired.Table { return nil, ErrDifferentTables } if !rowSecurityEqual(live.RowSecurity, desired.RowSecurity) { - return nil, fmt.Errorf("row security differs; policy execution is not supported: %w", ErrUnsupportedChange) + return nil, fmt.Errorf("row security differs; use the dedicated atomic RLS executor: %w", ErrUnsupportedChange) } live.RowSecurity = RowSecurity{} desired.RowSecurity = RowSecurity{} diff --git a/pkg/statement/desired_rls.go b/pkg/statement/desired_rls.go index 15058fd..e1a3518 100644 --- a/pkg/statement/desired_rls.go +++ b/pkg/statement/desired_rls.go @@ -18,8 +18,9 @@ var ErrRowSecurityDeclaration = errors.New("desired row security requires one ex // desired-state materialization can preserve their dependency identities. var ErrPolicyRelationDependency = errors.New("policy relation dependencies are not supported in desired row security") -// DesiredWithRowSecurity proves admission for scratch inspection only. It is -// deliberately distinct from DesiredSchema: live executors cannot accept it. +// DesiredWithRowSecurity proves declaration syntax for scratch inspection and +// the dedicated atomic RLS executor. It is distinct from DesiredSchema; generic +// native and create executors cannot accept it. Live safety needs executor checks. // The explicit parse entry point owns the complete table-local RLS definition, // including an empty policy set. Only ParseDesiredWithRowSecurity constructs it. type DesiredWithRowSecurity struct { diff --git a/pkg/statement/row_security_execution.go b/pkg/statement/row_security_execution.go new file mode 100644 index 0000000..48dcd23 --- /dev/null +++ b/pkg/statement/row_security_execution.go @@ -0,0 +1,58 @@ +package statement + +import ( + "fmt" + + pganalyze "github.com/pganalyze/pg_query_go/v6" + pgquery "github.com/wasilibs/go-pgquery" +) + +// SecurityStatements returns only the declaration's admitted RLS statements, +// qualified for the requested schema. Table/index definitions are never replayed. +// This does not prove live execution is safe: the dedicated executor must lock, +// compare the table definitions, enforce budgets, and verify the resulting state. +func (d DesiredWithRowSecurity) SecurityStatements(schema string) ([]string, error) { + if d.table == "" || schema == "" { + return nil, ErrRowSecurityDeclaration + } + var result []string + for _, st := range d.statements { + if st.kind != KindProvisioning { + continue + } + tree, err := pgquery.Parse(st.sql) + if err != nil { + return nil, fmt.Errorf("parse row security statement: %w", err) + } + if len(tree.GetStmts()) != 1 { + return nil, ErrNotOneStatement + } + node := tree.GetStmts()[0].GetStmt() + enabled, forced := 0, 0 + target, err := admitRowSecurity(node, &enabled, &forced) + if err != nil { + return nil, err + } + if target != d.table { + return nil, ErrRowSecurityDeclaration + } + switch { + case node.GetAlterTableStmt() != nil: + node.GetAlterTableStmt().Relation.Schemaname = schema + case node.GetCreatePolicyStmt() != nil: + node.GetCreatePolicyStmt().Table.Schemaname = schema + case node.GetCommentStmt() != nil: + list := node.GetCommentStmt().Object.GetList() + name := &pganalyze.Node{Node: &pganalyze.Node_String_{String_: &pganalyze.String{Sval: schema}}} + list.Items = append([]*pganalyze.Node{name}, list.Items...) + default: + return nil, ErrDisallowedStatement + } + sql, err := deparseOne(node) + if err != nil { + return nil, err + } + result = append(result, sql) + } + return result, nil +} diff --git a/pkg/statement/row_security_execution_test.go b/pkg/statement/row_security_execution_test.go new file mode 100644 index 0000000..56aef60 --- /dev/null +++ b/pkg/statement/row_security_execution_test.go @@ -0,0 +1,27 @@ +package statement_test + +import ( + "testing" + + "github.com/block/pg-sprite/pkg/statement" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestSecurityStatementsQualifiesOnlyRLS(t *testing.T) { + desired, err := statement.ParseDesiredWithRowSecurity(`CREATE TABLE documents ( + id bigint PRIMARY KEY + ); + CREATE INDEX idx_id ON documents (id); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (id = 7); + COMMENT ON POLICY readers ON documents IS 'readers';`) + require.NoError(t, err) + sql, err := desired.SecurityStatements(`odd"schema`) + require.NoError(t, err) + assert.Equal(t, []string{ + `ALTER TABLE "odd""schema".documents ENABLE ROW LEVEL SECURITY`, + `CREATE POLICY readers ON "odd""schema".documents FOR SELECT TO public USING (id = 7)`, + `COMMENT ON POLICY readers ON "odd""schema".documents IS 'readers'`, + }, sql) +} From 92caed444b4b87109ccc297b987f22372de60077 Mon Sep 17 00:00:00 2001 From: Armand Parajon Date: Tue, 22 Sep 2026 16:22:40 -0400 Subject: [PATCH 2/5] executor: tighten atomic RLS safety boundaries Signed-off-by: Armand Parajon --- SAFETY.md | 2 +- docs/atomic-row-security.md | 5 ++ docs/capabilities.md | 2 +- docs/invariants.md | 2 +- pkg/capabilities/capabilities.yaml | 2 +- pkg/executor/row_security.go | 48 +++++++++++---- pkg/executor/row_security_integration_test.go | 51 +++++++++++++++- pkg/executor/row_security_internal_test.go | 44 ++++++++++++++ pkg/schemadiff/row_security_roundtrip.go | 53 +++++++++++++---- pkg/statement/row_security_execution.go | 58 ------------------- pkg/statement/row_security_execution_test.go | 27 --------- 11 files changed, 182 insertions(+), 112 deletions(-) create mode 100644 pkg/executor/row_security_internal_test.go delete mode 100644 pkg/statement/row_security_execution.go delete mode 100644 pkg/statement/row_security_execution_test.go diff --git a/SAFETY.md b/SAFETY.md index 155bcee..ac1d7a7 100644 --- a/SAFETY.md +++ b/SAFETY.md @@ -129,5 +129,5 @@ The short version — the full rules live in [docs/tcb-model.md](docs/tcb-model. a bug in the periphery cannot corrupt data. The atomic RLS executor also admits `pkg/schemadiff` scratch introspection, table -comparison, and render admission into the core. Those calls refuse mixed or +comparison, render admission, and catalog-derived RLS rendering into the core. Those calls refuse mixed or unsupported table shapes; final catalog comparison gates commit (RS-1..RS-4). diff --git a/docs/atomic-row-security.md b/docs/atomic-row-security.md index 0af3255..0d1b319 100644 --- a/docs/atomic-row-security.md +++ b/docs/atomic-row-security.md @@ -44,6 +44,7 @@ Changing an RLS definition can widen access even when no data is deleted. savepoint on the same connection. Refuse unsupported table shapes or any table delta. 4. If RLS already matches, finish without policy DDL. Otherwise replace the complete policy set (including unchanged policies), apply ENABLE/DISABLE and FORCE/NO FORCE, and preserve policy comments. + Live SQL is rendered from the inspected desired catalog, not replayed from the input. 5. Read back the complete definition. Commit only if it matches the desired model. The lock blocks reads and writes briefly; this is bounded metadata DDL, not an @@ -51,6 +52,10 @@ online copy. Other sessions never see the intermediate policy set. A failure bef commit rolls back every change. A lost commit response is an unknown outcome: inspect the database before retrying. There are no automatic retries. +Lock exhaustion reports `budget-lock-exceeded`; the statement or whole-attempt +deadline reports `budget-statement-exceeded`. A missing target reports +`table-not-found`. Caller cancellation is kept separate from budget exhaustion. + The caller needs table-owner privileges and permission to create the temporary scratch schema. Roles and qualified helpers must already exist. Grants, role membership, helper bodies, authentication, and Supabase-managed schemas are outside diff --git a/docs/capabilities.md b/docs/capabilities.md index 2eda030..d158028 100644 --- a/docs/capabilities.md +++ b/docs/capabilities.md @@ -270,7 +270,7 @@ review the object warrants) · | Triggers (`CREATE TRIGGER`) | ⚪ | — | No — owner tooling | Catalog work — no scan, no rewrite — but it takes a brief `SHARE ROW EXCLUSIVE` on the table, queues behind long-running queries, and blocks writers while it waits — run it under a `lock_timeout` | | Extensions (`CREATE EXTENSION`) | ⚪ | — | No — owner tooling | Same: catalog bootstrap, owner tooling | | Grants, roles, and CLI policy DDL | 🔵 | — | No — provisioning / IaC | Access control changes remain with provisioning. Export includes RLS when present; explicit RLS declarations support comparison and review-only deltas with advisory access warnings. Plain table files leave access control separately managed. The separate [atomic RLS Go executor](atomic-row-security.md) supports RLS-only changes on existing supported tables; it does not route through these CLI front doors. See the [workflow and roadmap](declarative-row-security.md). See [engine-role.md](engine-role.md) for the engine's own role | -| Complete table-local RLS definition (Go API only) | ✅ | native, as-is | Yes | The [atomic RLS executor](atomic-row-security.md) locks an existing supported table, refuses structural changes, replaces policies/settings, and verifies convergence before committing. Lock waits and the whole transaction are bounded. No CLI execution or new diff flags | +| Complete table-local RLS definition (Go API only) | ✅ | native, safer sequence | Yes | The [atomic RLS executor](atomic-row-security.md) locks an existing supported table, refuses structural changes, replaces policies/settings, and verifies convergence before committing. Lock waits and the whole transaction are bounded. No CLI execution or new diff flags | | Standalone sequences | ⚪ | — | No — owner tooling | Transactional catalog work on an object with no readers-and-writers problem | | Publications, subscriptions | 🔵 | — | No — replication provisioning / IaC | Replication provisioning, not table shape (`ALTER PUBLICATION ... ADD TABLE` also takes `SHARE UPDATE EXCLUSIVE` on the table) | diff --git a/docs/invariants.md b/docs/invariants.md index 49f6810..10ace7a 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -401,7 +401,7 @@ The [atomic RLS contract](atomic-row-security.md) defines these executor obligat | RS-1 | Lock the live target before deriving the change; refuse mixed table changes | `ExecuteRowSecurity`; contention and mixed-change tests | | RS-2 | Policy changes and convergence verification commit together or roll back | `ExecuteRowSecurity`; real DDL fault injection restores original policies | | RS-3 | Bound lock waits, statements, and the whole attempt | `ExecuteRowSecurity`; lock contention and deadline cancellation | -| RS-4 | Run only admitted, qualified RLS DDL; keep scratch disposable | `SecurityStatements` and executor readback; quoted names and scratch-cleanup tests | +| RS-4 | Run only admitted, qualified RLS DDL; keep scratch disposable | `RenderRowSecurity` from the scratch catalog and executor readback; qualified-helper, quoted-name, and scratch-cleanup tests | ## State, checkpoint, and resume (ST) diff --git a/pkg/capabilities/capabilities.yaml b/pkg/capabilities/capabilities.yaml index 89679df..29ce674 100644 --- a/pkg/capabilities/capabilities.yaml +++ b/pkg/capabilities/capabilities.yaml @@ -525,7 +525,7 @@ rows: operation: "Complete table-local RLS definition (Go API only)" tier: "t1" status_mark: "✅" - engine_path: "native_as_is" + engine_path: "native_safer_sequence" online_safety_problem: true front_doors: migrate: refused diff --git a/pkg/executor/row_security.go b/pkg/executor/row_security.go index ac7d636..a6afaf3 100644 --- a/pkg/executor/row_security.go +++ b/pkg/executor/row_security.go @@ -2,11 +2,13 @@ package executor import ( "context" + "errors" "fmt" "strconv" "time" "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" "github.com/jackc/pgx/v5/pgxpool" "github.com/block/pg-sprite/pkg/dbconn" @@ -43,13 +45,37 @@ func ExecuteRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, if err := b.validate(); err != nil { return RowSecurityReport{}, err } - security, err := desired.SecurityStatements(schema) - if err != nil { - return RowSecurityReport{}, err + if desired.Table() == "" || schema == "" { + return RowSecurityReport{}, statement.ErrRowSecurityDeclaration } // INV: RS-3 — the whole attempt has one deadline, not a fresh budget per policy. - ctx, cancel := context.WithTimeout(ctx, b.StatementTimeout) + attempt, cancel := context.WithTimeout(ctx, b.StatementTimeout) defer cancel() + report, err := executeRowSecurity(attempt, pool, schema, desired, b) + return report, rowSecurityError(ctx, attempt, err, b) +} + +func rowSecurityError(caller, attempt context.Context, err error, b Budget) error { + if err == nil { + return nil + } + var unknown *RowSecurityOutcomeUnknownError + if errors.As(err, &unknown) { + return err + } + if caller.Err() != nil && errors.Is(err, caller.Err()) { + return fmt.Errorf("row security caller stopped: %w", caller.Err()) + } + if errors.Is(err, context.DeadlineExceeded) && errors.Is(attempt.Err(), context.DeadlineExceeded) { + return &BudgetError{Cause: CauseStatement, Budget: b.StatementTimeout, cause: err} + } + if budgetErr := asBudgetError(err, b); budgetErr != nil { + return budgetErr + } + return err +} + +func executeRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, desired statement.DesiredWithRowSecurity, b Budget) (RowSecurityReport, error) { tx, err := pool.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.ReadCommitted}) if err != nil { return RowSecurityReport{}, fmt.Errorf("begin row security change: %w", err) @@ -66,6 +92,10 @@ func ExecuteRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, target := pgx.Identifier{schema, desired.Table()}.Sanitize() // INV: RS-1 — all live comparison and DDL occur after this exclusive lock. if _, err := tx.Exec(ctx, "LOCK TABLE ONLY "+target+" IN ACCESS EXCLUSIVE MODE"); err != nil { + var pgErr *pgconn.PgError + if errors.As(err, &pgErr) && pgErr.Code == sqlstateUndefinedTable { + return RowSecurityReport{}, fmt.Errorf("lock row security target %s: %w: %w", target, ErrTableNotFound, err) + } return RowSecurityReport{}, fmt.Errorf("lock row security target %s: %w", target, err) } live, err := schemadiff.IntrospectTx(ctx, tx, schema, desired.Table()) @@ -89,13 +119,11 @@ func ExecuteRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, for _, policy := range live.RowSecurity.Policies { report.Statements = append(report.Statements, "DROP POLICY "+pgx.Identifier{policy.Name}.Sanitize()+" ON "+target) } - report.Statements = append(report.Statements, security...) - // An omitted FORCE declaration means NO FORCE, not "leave it alone". - force := "NO FORCE" - if wanted.RowSecurity.Forced { - force = "FORCE" + security, err := schemadiff.RenderRowSecurity(schema, wanted) + if err != nil { + return RowSecurityReport{}, err } - report.Statements = append(report.Statements, "ALTER TABLE "+target+" "+force+" ROW LEVEL SECURITY") + report.Statements = append(report.Statements, security...) for _, sql := range report.Statements { if _, err := tx.Exec(ctx, sql); err != nil { return RowSecurityReport{}, fmt.Errorf("apply row security on %s: %w", target, err) diff --git a/pkg/executor/row_security_integration_test.go b/pkg/executor/row_security_integration_test.go index 4bd94d5..4013d95 100644 --- a/pkg/executor/row_security_integration_test.go +++ b/pkg/executor/row_security_integration_test.go @@ -120,6 +120,7 @@ func TestExecuteRowSecurityLockContentionLeavesPoliciesIntact(t *testing.T) { var pgErr *pgconn.PgError require.ErrorAs(t, err, &pgErr) assert.Equal(t, "55P03", pgErr.Code) + assert.Equal(t, executor.CodeBudgetLockExceeded, executor.OutcomeCode(err)) after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") require.NoError(t, err) assert.Equal(t, before, after) @@ -204,7 +205,8 @@ func TestExecuteRowSecurityRefusesMissingTable(t *testing.T) { id bigint PRIMARY KEY ); ALTER TABLE missing ENABLE ROW LEVEL SECURITY;`) - require.Error(t, err) + require.ErrorIs(t, err, executor.ErrTableNotFound) + assert.Equal(t, executor.CodeTableNotFound, executor.OutcomeCode(err)) assert.Empty(t, relationKind(t, pool, schema, "missing")) } @@ -320,6 +322,7 @@ func TestExecuteRowSecurityDeadlineRollsBackLiveDDL(t *testing.T) { require.NoError(t, err) report, err := executor.ExecuteRowSecurity(t.Context(), pool, schema, desired, executor.Budget{LockTimeout: 50 * time.Millisecond, StatementTimeout: 2 * time.Second}) require.Error(t, err) + assert.Equal(t, executor.CodeBudgetStatementExceeded, executor.OutcomeCode(err)) assert.Empty(t, report.Statements) var reached bool require.NoError(t, pool.QueryRow(t.Context(), "SELECT is_called FROM "+pgx.Identifier{schema, "fault_reached"}.Sanitize()).Scan(&reached)) @@ -328,3 +331,49 @@ func TestExecuteRowSecurityDeadlineRollsBackLiveDDL(t *testing.T) { require.NoError(t, err) assert.Equal(t, before, after) } + +// Qualified application helpers keep their binding; a target-schema function +// with a built-in's name must not change the policy during live execution. +func TestExecuteRowSecurityPreservesHelperResolution(t *testing.T) { + pool, schema := rlsFixture(t) + helpers := testutil.NewSchema(t, pool) + _, err := pool.Exec(t.Context(), fmt.Sprintf(` + CREATE FUNCTION %s.allowed_owner() RETURNS bigint + LANGUAGE sql IMMUTABLE AS 'SELECT 9::bigint'; + CREATE FUNCTION %s.abs(bigint) RETURNS bigint + LANGUAGE sql IMMUTABLE AS 'SELECT 999::bigint'; + INSERT INTO %s.documents VALUES (1, -9), (2, -7);`, helpers, schema, schema)) + require.NoError(t, err) + sql := fmt.Sprintf(`CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT + USING (abs(owner_id) = %s.allowed_owner());`, helpers) + report, err := applyRLS(t, pool, schema, sql) + require.NoError(t, err) + actual, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + require.Len(t, actual.RowSecurity.Policies, 1) + assert.Equal(t, "(abs(owner_id) = "+helpers+".allowed_owner())", *actual.RowSecurity.Policies[0].Using) + // SQL in the committed report comes from PostgreSQL's canonical catalog form. + assert.Contains(t, report.Statements[3], "USING ((abs(owner_id) = "+helpers+".allowed_owner()))") + role := pgx.Identifier{schema + "_reader"}.Sanitize() + _, err = pool.Exec(t.Context(), fmt.Sprintf(`CREATE ROLE %s; + GRANT USAGE ON SCHEMA %s, %s TO %s; + GRANT SELECT ON %s.documents TO %s;`, role, schema, helpers, role, schema, role)) + require.NoError(t, err) + t.Cleanup(func() { + _, err := pool.Exec(context.WithoutCancel(t.Context()), "DROP OWNED BY "+role+"; DROP ROLE "+role) + require.NoError(t, err) + }) + tx, err := pool.Begin(t.Context()) + require.NoError(t, err) + defer func() { _ = tx.Rollback(context.WithoutCancel(t.Context())) }() + _, err = tx.Exec(t.Context(), "SET LOCAL ROLE "+role) + require.NoError(t, err) + var ids []int64 + require.NoError(t, tx.QueryRow(t.Context(), "SELECT array_agg(id ORDER BY id) FROM "+pgx.Identifier{schema, "documents"}.Sanitize()).Scan(&ids)) + assert.Equal(t, []int64{1}, ids) +} diff --git a/pkg/executor/row_security_internal_test.go b/pkg/executor/row_security_internal_test.go new file mode 100644 index 0000000..4c013d9 --- /dev/null +++ b/pkg/executor/row_security_internal_test.go @@ -0,0 +1,44 @@ +package executor + +import ( + "context" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgconn" + "github.com/stretchr/testify/assert" +) + +func TestRowSecurityErrorClassification(t *testing.T) { + budget := Budget{LockTimeout: time.Second, StatementTimeout: time.Second} + t.Run("server statement timeout", func(t *testing.T) { + err := rowSecurityError(t.Context(), t.Context(), &pgconn.PgError{Code: "57014"}, budget) + assert.Equal(t, CodeBudgetStatementExceeded, OutcomeCode(err)) + }) + t.Run("attempt deadline", func(t *testing.T) { + attempt, cancel := context.WithDeadline(t.Context(), time.Time{}) + defer cancel() + err := rowSecurityError(t.Context(), attempt, context.DeadlineExceeded, budget) + assert.Equal(t, CodeBudgetStatementExceeded, OutcomeCode(err)) + }) + t.Run("caller cancellation is not budget exhaustion", func(t *testing.T) { + caller, cancel := context.WithCancel(t.Context()) + cancel() + err := rowSecurityError(caller, caller, context.Canceled, budget) + assert.ErrorIs(t, err, context.Canceled) + var budgetErr *BudgetError + assert.NotErrorAs(t, err, &budgetErr) + }) + t.Run("unrelated failure survives expired deadline", func(t *testing.T) { + attempt, cancel := context.WithDeadline(t.Context(), time.Time{}) + defer cancel() + err := rowSecurityError(t.Context(), attempt, ErrTableNotFound, budget) + assert.Equal(t, CodeTableNotFound, OutcomeCode(err)) + }) + t.Run("commit uncertainty survives expired deadline", func(t *testing.T) { + attempt, cancel := context.WithDeadline(t.Context(), time.Time{}) + defer cancel() + err := rowSecurityError(t.Context(), attempt, &RowSecurityOutcomeUnknownError{Err: context.DeadlineExceeded}, budget) + assert.Equal(t, CodeRowSecurityOutcomeUnknown, OutcomeCode(err)) + }) +} diff --git a/pkg/schemadiff/row_security_roundtrip.go b/pkg/schemadiff/row_security_roundtrip.go index fa2c42a..1132b15 100644 --- a/pkg/schemadiff/row_security_roundtrip.go +++ b/pkg/schemadiff/row_security_roundtrip.go @@ -48,26 +48,55 @@ func RenderWithRowSecurity(m Model) (string, error) { var b strings.Builder b.WriteString(base) target := pgx.Identifier{m.Table}.Sanitize() + if err := renderRowSecurity(&b, target, m.RowSecurity); err != nil { + return "", err + } + out := b.String() + if _, err := statement.ParseDesiredWithRowSecurity(out); err != nil { + return "", fmt.Errorf("render row security for %q: %w", m.Table, err) + } + return out, nil +} + +// RenderRowSecurity renders only settings, policies, and policy comments from +// a catalog model. It does not authorize execution; callers must admit the model +// and execute with pg_catalog as the search path, as used during introspection. +func RenderRowSecurity(schema string, m Model) ([]string, error) { + if schema == "" || m.Table == "" { + return nil, ErrUnrenderableRowSecurity + } + var b strings.Builder + if err := renderRowSecurity(&b, pgx.Identifier{schema, m.Table}.Sanitize(), m.RowSecurity); err != nil { + return nil, err + } + statements, err := statement.Split(b.String()) + if err != nil { + return nil, err + } + result := make([]string, len(statements)) + for i, st := range statements { + result[i] = st.SQL + } + return result, nil +} + +func renderRowSecurity(b *strings.Builder, target string, security RowSecurity) error { state := "DISABLE" - if m.RowSecurity.Enabled { + if security.Enabled { state = "ENABLE" } - fmt.Fprintf(&b, "\nALTER TABLE %s %s ROW LEVEL SECURITY;\n", target, state) + fmt.Fprintf(b, "\nALTER TABLE %s %s ROW LEVEL SECURITY;\n", target, state) force := "NO FORCE" - if m.RowSecurity.Forced { + if security.Forced { force = "FORCE" } - fmt.Fprintf(&b, "ALTER TABLE %s %s ROW LEVEL SECURITY;\n", target, force) - for _, policy := range m.RowSecurity.Policies { - if err := renderPolicy(&b, target, policy); err != nil { - return "", err + fmt.Fprintf(b, "ALTER TABLE %s %s ROW LEVEL SECURITY;\n", target, force) + for _, policy := range security.Policies { + if err := renderPolicy(b, target, policy); err != nil { + return err } } - out := b.String() - if _, err := statement.ParseDesiredWithRowSecurity(out); err != nil { - return "", fmt.Errorf("render row security for %q: %w", m.Table, err) - } - return out, nil + return nil } func renderPolicy(b *strings.Builder, target string, p Policy) error { diff --git a/pkg/statement/row_security_execution.go b/pkg/statement/row_security_execution.go deleted file mode 100644 index 48dcd23..0000000 --- a/pkg/statement/row_security_execution.go +++ /dev/null @@ -1,58 +0,0 @@ -package statement - -import ( - "fmt" - - pganalyze "github.com/pganalyze/pg_query_go/v6" - pgquery "github.com/wasilibs/go-pgquery" -) - -// SecurityStatements returns only the declaration's admitted RLS statements, -// qualified for the requested schema. Table/index definitions are never replayed. -// This does not prove live execution is safe: the dedicated executor must lock, -// compare the table definitions, enforce budgets, and verify the resulting state. -func (d DesiredWithRowSecurity) SecurityStatements(schema string) ([]string, error) { - if d.table == "" || schema == "" { - return nil, ErrRowSecurityDeclaration - } - var result []string - for _, st := range d.statements { - if st.kind != KindProvisioning { - continue - } - tree, err := pgquery.Parse(st.sql) - if err != nil { - return nil, fmt.Errorf("parse row security statement: %w", err) - } - if len(tree.GetStmts()) != 1 { - return nil, ErrNotOneStatement - } - node := tree.GetStmts()[0].GetStmt() - enabled, forced := 0, 0 - target, err := admitRowSecurity(node, &enabled, &forced) - if err != nil { - return nil, err - } - if target != d.table { - return nil, ErrRowSecurityDeclaration - } - switch { - case node.GetAlterTableStmt() != nil: - node.GetAlterTableStmt().Relation.Schemaname = schema - case node.GetCreatePolicyStmt() != nil: - node.GetCreatePolicyStmt().Table.Schemaname = schema - case node.GetCommentStmt() != nil: - list := node.GetCommentStmt().Object.GetList() - name := &pganalyze.Node{Node: &pganalyze.Node_String_{String_: &pganalyze.String{Sval: schema}}} - list.Items = append([]*pganalyze.Node{name}, list.Items...) - default: - return nil, ErrDisallowedStatement - } - sql, err := deparseOne(node) - if err != nil { - return nil, err - } - result = append(result, sql) - } - return result, nil -} diff --git a/pkg/statement/row_security_execution_test.go b/pkg/statement/row_security_execution_test.go deleted file mode 100644 index 56aef60..0000000 --- a/pkg/statement/row_security_execution_test.go +++ /dev/null @@ -1,27 +0,0 @@ -package statement_test - -import ( - "testing" - - "github.com/block/pg-sprite/pkg/statement" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" -) - -func TestSecurityStatementsQualifiesOnlyRLS(t *testing.T) { - desired, err := statement.ParseDesiredWithRowSecurity(`CREATE TABLE documents ( - id bigint PRIMARY KEY - ); - CREATE INDEX idx_id ON documents (id); - ALTER TABLE documents ENABLE ROW LEVEL SECURITY; - CREATE POLICY readers ON documents FOR SELECT USING (id = 7); - COMMENT ON POLICY readers ON documents IS 'readers';`) - require.NoError(t, err) - sql, err := desired.SecurityStatements(`odd"schema`) - require.NoError(t, err) - assert.Equal(t, []string{ - `ALTER TABLE "odd""schema".documents ENABLE ROW LEVEL SECURITY`, - `CREATE POLICY readers ON "odd""schema".documents FOR SELECT TO public USING (id = 7)`, - `COMMENT ON POLICY readers ON "odd""schema".documents IS 'readers'`, - }, sql) -} From 1803dba0e1f51f0948baaf69547102dae7302354 Mon Sep 17 00:00:00 2001 From: Armand Parajon Date: Tue, 22 Sep 2026 20:26:40 -0400 Subject: [PATCH 3/5] executor: classify RLS refusals and verify ownership Signed-off-by: Armand Parajon --- docs/atomic-row-security.md | 10 +- docs/execution-model.md | 1 + pkg/executor/code.go | 7 +- pkg/executor/code_test.go | 1 + pkg/executor/row_security.go | 29 +++- ...row_security_admission_integration_test.go | 152 ++++++++++++++++++ pkg/executor/row_security_integration_test.go | 2 + pkg/executor/row_security_internal_test.go | 10 ++ pkg/executor/row_security_owner.go | 31 ++++ pkg/executor/row_security_test.go | 4 + 10 files changed, 236 insertions(+), 11 deletions(-) create mode 100644 pkg/executor/row_security_admission_integration_test.go create mode 100644 pkg/executor/row_security_owner.go diff --git a/docs/atomic-row-security.md b/docs/atomic-row-security.md index 0d1b319..f11911c 100644 --- a/docs/atomic-row-security.md +++ b/docs/atomic-row-security.md @@ -54,10 +54,13 @@ inspect the database before retrying. There are no automatic retries. Lock exhaustion reports `budget-lock-exceeded`; the statement or whole-attempt deadline reports `budget-statement-exceeded`. A missing target reports -`table-not-found`. Caller cancellation is kept separate from budget exhaustion. +`table-not-found`. Invalid declarations, unsupported targets, and insufficient +privileges report permanent `row-security-refused` outcomes, preserving the underlying +cause. Caller cancellation is kept separate from budget exhaustion. The caller needs table-owner privileges and permission to create the temporary -scratch schema. Roles and qualified helpers must already exist. Grants, role +scratch schema. Ownership is checked before locking and checked again under the +lock. Roles and qualified helpers must already exist. Grants, role membership, helper bodies, authentication, and Supabase-managed schemas are outside this operation. Application authorization tests are still needed. Concurrent administration of those dependencies is not serialized by the table lock. @@ -65,7 +68,8 @@ administration of those dependencies is not serialized by the table lock. ## Invariants and tests - **RS-1:** Read the live baseline only after taking the target lock; never accept a - caller-supplied diff as execution authority. Refuse mixed changes before live DDL. + caller-supplied diff as execution authority. Verify ownership before and after locking. + Refuse mixed changes and unsupported table shapes before live DDL. - **RS-2:** All policy/settings changes and the final catalog comparison share one transaction. Fault injection after live DDL must prove the original state survives. - **RS-3:** Bound lock waits, individual statements, and the entire attempt. Cancellation diff --git a/docs/execution-model.md b/docs/execution-model.md index 5d2469d..7981932 100644 --- a/docs/execution-model.md +++ b/docs/execution-model.md @@ -289,6 +289,7 @@ concurrently is. | --- | --- | --- | | `budget-lock-exceeded` | no | The lock was not granted within `lock_timeout`; nothing executed | | `budget-statement-exceeded` | no | The statement ran past `statement_timeout` and was cancelled | +| `row-security-refused` | yes | Change the declaration, unsupported target shape, or privileges before retrying | | `row-security-outcome-unknown` | no | The atomic RLS commit response is uncertain; inspect the catalog before retrying | | `blocking-outcome-unknown` | no | The accepted blocking transaction reached an ambiguous client boundary; inspect the catalog before retrying | | `invalid-blocking-budget` | yes | An accepted blocking bound is disabled or cannot be represented by PostgreSQL | diff --git a/pkg/executor/code.go b/pkg/executor/code.go index 898c900..6d515e1 100644 --- a/pkg/executor/code.go +++ b/pkg/executor/code.go @@ -23,6 +23,8 @@ const ( // CodeBudgetStatementExceeded: the statement ran past // statement_timeout and was cancelled; the change does real work. CodeBudgetStatementExceeded Code = "budget-statement-exceeded" + // CodeRowSecurityRefused requires a declaration, target, or privilege change. + CodeRowSecurityRefused Code = "row-security-refused" // CodeRowSecurityOutcomeUnknown requires catalog inspection before retry. CodeRowSecurityOutcomeUnknown Code = "row-security-outcome-unknown" // CodeBlockingOutcomeUnknown requires catalog inspection before retry. @@ -150,6 +152,7 @@ func Codes() []Code { CodeBudgetLockExceeded, CodeBudgetStatementExceeded, CodeBlockingOutcomeUnknown, + CodeRowSecurityRefused, CodeRowSecurityOutcomeUnknown, CodeInvalidBlockingBudget, CodeUnsupportedAcceptedBlocking, @@ -198,7 +201,7 @@ func Codes() []Code { // permanent. func (c Code) Permanent() bool { switch c { - case CodeInvalidIndexOtherTable, + case CodeRowSecurityRefused, CodeInvalidIndexOtherTable, CodeInvalidIndexNotDroppable, CodeInvalidBlockingBudget, CodeUnsupportedAcceptedBlocking, @@ -309,6 +312,8 @@ func sentinelCode(err error) Code { return CodeUnsupportedCreateStep case errors.Is(err, ErrPoolTooSmall): return CodePoolTooSmall + case errors.Is(err, ErrRowSecurityRefused): + return CodeRowSecurityRefused case errors.Is(err, ErrTableNotFound): return CodeTableNotFound default: diff --git a/pkg/executor/code_test.go b/pkg/executor/code_test.go index 47c7d36..cdf99f3 100644 --- a/pkg/executor/code_test.go +++ b/pkg/executor/code_test.go @@ -163,6 +163,7 @@ func TestCodePermanentClassifiesEveryCode(t *testing.T) { executor.CodeBudgetLockExceeded: false, executor.CodeBudgetStatementExceeded: false, executor.CodeBlockingOutcomeUnknown: false, + executor.CodeRowSecurityRefused: true, executor.CodeRowSecurityOutcomeUnknown: false, executor.CodeInvalidBlockingBudget: true, executor.CodeUnsupportedAcceptedBlocking: true, diff --git a/pkg/executor/row_security.go b/pkg/executor/row_security.go index a6afaf3..77363c7 100644 --- a/pkg/executor/row_security.go +++ b/pkg/executor/row_security.go @@ -46,7 +46,7 @@ func ExecuteRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, return RowSecurityReport{}, err } if desired.Table() == "" || schema == "" { - return RowSecurityReport{}, statement.ErrRowSecurityDeclaration + return RowSecurityReport{}, fmt.Errorf("%w: %w", ErrRowSecurityRefused, statement.ErrRowSecurityDeclaration) } // INV: RS-3 — the whole attempt has one deadline, not a fresh budget per policy. attempt, cancel := context.WithTimeout(ctx, b.StatementTimeout) @@ -72,6 +72,13 @@ func rowSecurityError(caller, attempt context.Context, err error, b Budget) erro if budgetErr := asBudgetError(err, b); budgetErr != nil { return budgetErr } + if errors.Is(err, statement.ErrPolicyRelationDependency) || errors.Is(err, statement.ErrRowSecurityDeclaration) { + return fmt.Errorf("%w: %w", ErrRowSecurityRefused, err) + } + var pgErr *pgconn.PgError + if errors.As(err, &pgErr) && pgErr.Code == "42501" { + return fmt.Errorf("%w: %w", ErrRowSecurityRefused, err) + } return err } @@ -90,6 +97,10 @@ func executeRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, return RowSecurityReport{}, fmt.Errorf("set row security budgets: %w", err) } target := pgx.Identifier{schema, desired.Table()}.Sanitize() + // Reject roles without owner privileges before taking an application-blocking lock. + if err := checkRowSecurityOwner(ctx, tx, schema, desired.Table()); err != nil { + return RowSecurityReport{}, err + } // INV: RS-1 — all live comparison and DDL occur after this exclusive lock. if _, err := tx.Exec(ctx, "LOCK TABLE ONLY "+target+" IN ACCESS EXCLUSIVE MODE"); err != nil { var pgErr *pgconn.PgError @@ -98,30 +109,34 @@ func executeRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, } return RowSecurityReport{}, fmt.Errorf("lock row security target %s: %w", target, err) } + // Ownership may have changed while waiting for the lock. Recheck it under lock. + if err := checkRowSecurityOwner(ctx, tx, schema, desired.Table()); err != nil { + return RowSecurityReport{}, err + } live, err := schemadiff.IntrospectTx(ctx, tx, schema, desired.Table()) if err != nil { - return RowSecurityReport{}, err + return RowSecurityReport{}, fmt.Errorf("inspect row security target %s: %w", target, err) } wanted, err := schemadiff.IntrospectDesiredWithRowSecurityTx(ctx, tx, desired) if err != nil { - return RowSecurityReport{}, err + return RowSecurityReport{}, fmt.Errorf("inspect desired row security for %s: %w", target, err) } if err := admitRowSecurityTable(schema, live, wanted); err != nil { - return RowSecurityReport{}, err + return RowSecurityReport{}, fmt.Errorf("admit row security target %s: %w: %w", target, ErrRowSecurityRefused, err) } report := RowSecurityReport{Schema: schema, Table: desired.Table(), Statements: []string{}} if _, err := schemadiff.DiffWithRowSecurity(schema, live, wanted); err != nil { // INV: RS-4 — helpers are explicitly qualified; no target-schema function // may shadow a built-in while replaying policy expressions. if _, err := tx.Exec(ctx, dbconn.LocalSearchPath("pg_catalog")); err != nil { - return RowSecurityReport{}, err + return RowSecurityReport{}, fmt.Errorf("set row security search path for %s: %w", target, err) } for _, policy := range live.RowSecurity.Policies { report.Statements = append(report.Statements, "DROP POLICY "+pgx.Identifier{policy.Name}.Sanitize()+" ON "+target) } security, err := schemadiff.RenderRowSecurity(schema, wanted) if err != nil { - return RowSecurityReport{}, err + return RowSecurityReport{}, fmt.Errorf("render row security for %s: %w: %w", target, ErrRowSecurityRefused, err) } report.Statements = append(report.Statements, security...) for _, sql := range report.Statements { @@ -133,7 +148,7 @@ func executeRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, // INV: RS-2 — verify convergence before committing, while retaining the lock. actual, err := schemadiff.IntrospectTx(ctx, tx, schema, desired.Table()) if err != nil { - return RowSecurityReport{}, err + return RowSecurityReport{}, fmt.Errorf("verify row security target %s: %w", target, err) } if err := admitRowSecurityTable(schema, actual, wanted); err != nil { return RowSecurityReport{}, fmt.Errorf("%w: RS-2: target shape changed: %w", ErrInvariantViolation, err) diff --git a/pkg/executor/row_security_admission_integration_test.go b/pkg/executor/row_security_admission_integration_test.go new file mode 100644 index 0000000..a7ab90a --- /dev/null +++ b/pkg/executor/row_security_admission_integration_test.go @@ -0,0 +1,152 @@ +package executor_test + +import ( + "context" + "fmt" + "testing" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/executor" + "github.com/block/pg-sprite/pkg/schemadiff" +) + +func requireRLSShapeRefusal(t *testing.T, pool *pgxpool.Pool, schema string, cause error) { + t.Helper() + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + report, err := applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents DISABLE ROW LEVEL SECURITY;`) + require.ErrorIs(t, err, cause) + assert.Equal(t, executor.CodeRowSecurityRefused, executor.OutcomeCode(err)) + assert.True(t, executor.OutcomeCode(err).Permanent()) + assert.Empty(t, report.Statements) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} + +func TestExecuteRowSecurityRefusesReferencedTable(t *testing.T) { + pool, schema := rlsFixture(t) + _, err := pool.Exec(t.Context(), fmt.Sprintf(`CREATE TABLE %s.references_documents ( + document_id bigint REFERENCES %s.documents(id) + );`, schema, schema)) + require.NoError(t, err) + requireRLSShapeRefusal(t, pool, schema, schemadiff.ErrUnrenderableForeignKey) +} + +func TestExecuteRowSecurityRefusesInheritanceParent(t *testing.T) { + pool, schema := rlsFixture(t) + _, err := pool.Exec(t.Context(), fmt.Sprintf(`CREATE TABLE %s.child_documents () INHERITS (%s.documents);`, schema, schema)) + require.NoError(t, err) + requireRLSShapeRefusal(t, pool, schema, schemadiff.ErrUnrenderableInheritance) +} + +func TestExecuteRowSecurityRefusesInheritanceChild(t *testing.T) { + pool, schema := rlsFixture(t) + _, err := pool.Exec(t.Context(), fmt.Sprintf(`CREATE TABLE %s.parent_documents ( + id bigint NOT NULL, + owner_id bigint NOT NULL + ); + ALTER TABLE %s.documents INHERIT %s.parent_documents;`, schema, schema, schema)) + require.NoError(t, err) + requireRLSShapeRefusal(t, pool, schema, schemadiff.ErrUnrenderableInheritance) +} + +func TestExecuteRowSecurityRefusesUnloggedTable(t *testing.T) { + pool, schema := rlsFixture(t) + _, err := pool.Exec(t.Context(), fmt.Sprintf(`ALTER TABLE %s.documents SET UNLOGGED;`, schema)) + require.NoError(t, err) + requireRLSShapeRefusal(t, pool, schema, schemadiff.ErrUnrenderableUnlogged) +} + +func TestExecuteRowSecurityRefusesPartitionedParent(t *testing.T) { + pool, schema := rlsFixture(t) + _, err := pool.Exec(t.Context(), fmt.Sprintf(`DROP TABLE %s.documents; + CREATE TABLE %s.documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ) PARTITION BY RANGE (id); + ALTER TABLE %s.documents ENABLE ROW LEVEL SECURITY;`, schema, schema, schema)) + require.NoError(t, err) + requireRLSShapeRefusal(t, pool, schema, schemadiff.ErrUnrenderablePartition) +} + +func TestExecuteRowSecurityRefusesNonOwnerBeforeLock(t *testing.T) { + pool, schema := rlsFixture(t) + role := pgx.Identifier{schema + "_writer"}.Sanitize() + _, err := pool.Exec(t.Context(), fmt.Sprintf(`CREATE ROLE %s; + GRANT USAGE ON SCHEMA %s TO %s; + GRANT UPDATE ON %s.documents TO %s;`, role, schema, role, schema, role)) + require.NoError(t, err) + t.Cleanup(func() { + _, err := pool.Exec(context.WithoutCancel(t.Context()), "DROP OWNED BY "+role+"; DROP ROLE "+role) + require.NoError(t, err) + }) + cfg := pool.Config() + cfg.AfterConnect = func(ctx context.Context, conn *pgx.Conn) error { + _, err := conn.Exec(ctx, "SET ROLE "+role) + return err + } + restricted, err := pgxpool.NewWithConfig(t.Context(), cfg) + require.NoError(t, err) + defer restricted.Close() + // A conflicting lock makes the assertion distinguish early owner refusal + // from an attempt that first waits for ACCESS EXCLUSIVE and times out. + tx, err := pool.Begin(t.Context()) + require.NoError(t, err) + defer func() { _ = tx.Rollback(context.WithoutCancel(t.Context())) }() + _, err = tx.Exec(t.Context(), "LOCK TABLE "+pgx.Identifier{schema, "documents"}.Sanitize()+" IN ACCESS SHARE MODE") + require.NoError(t, err) + _, err = applyRLS(t, restricted, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents DISABLE ROW LEVEL SECURITY;`) + require.ErrorIs(t, err, executor.ErrRowSecurityRefused) + assert.Equal(t, executor.CodeRowSecurityRefused, executor.OutcomeCode(err)) + assert.True(t, executor.OutcomeCode(err).Permanent()) +} + +func TestExecuteRowSecurityAcceptsInheritedOwnership(t *testing.T) { + pool, schema := rlsFixture(t) + owner := pgx.Identifier{schema + "_owner"}.Sanitize() + role := pgx.Identifier{schema + "_engine"}.Sanitize() + var database string + require.NoError(t, pool.QueryRow(t.Context(), "SELECT current_database()").Scan(&database)) + _, err := pool.Exec(t.Context(), fmt.Sprintf(`CREATE ROLE %s; + CREATE ROLE %s INHERIT; + GRANT %s TO %s; + GRANT USAGE ON SCHEMA %s TO %s; + GRANT CREATE ON DATABASE %s TO %s; + ALTER TABLE %s.documents OWNER TO %s;`, owner, role, owner, role, schema, role, pgx.Identifier{database}.Sanitize(), role, schema, owner)) + require.NoError(t, err) + t.Cleanup(func() { + _, err := pool.Exec(context.WithoutCancel(t.Context()), "DROP OWNED BY "+role+"; DROP OWNED BY "+owner+" CASCADE; DROP ROLE "+role+"; DROP ROLE "+owner) + require.NoError(t, err) + }) + cfg := pool.Config() + cfg.AfterConnect = func(ctx context.Context, conn *pgx.Conn) error { + _, err := conn.Exec(ctx, "SET ROLE "+role) + return err + } + engine, err := pgxpool.NewWithConfig(t.Context(), cfg) + require.NoError(t, err) + defer engine.Close() + _, err = applyRLS(t, engine, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents DISABLE ROW LEVEL SECURITY;`) + require.NoError(t, err) + actual, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.False(t, actual.RowSecurity.Enabled) + assert.Empty(t, actual.RowSecurity.Policies) +} diff --git a/pkg/executor/row_security_integration_test.go b/pkg/executor/row_security_integration_test.go index 4013d95..7b295c4 100644 --- a/pkg/executor/row_security_integration_test.go +++ b/pkg/executor/row_security_integration_test.go @@ -97,6 +97,8 @@ func TestExecuteRowSecurityRefusesMixedChanges(t *testing.T) { ); ALTER TABLE documents DISABLE ROW LEVEL SECURITY;`) require.ErrorIs(t, err, schemadiff.ErrUnsupportedChange) + assert.Equal(t, executor.CodeRowSecurityRefused, executor.OutcomeCode(err)) + assert.True(t, executor.OutcomeCode(err).Permanent()) assert.Empty(t, report.Statements) after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") require.NoError(t, err) diff --git a/pkg/executor/row_security_internal_test.go b/pkg/executor/row_security_internal_test.go index 4c013d9..433ee8a 100644 --- a/pkg/executor/row_security_internal_test.go +++ b/pkg/executor/row_security_internal_test.go @@ -5,6 +5,7 @@ import ( "testing" "time" + "github.com/block/pg-sprite/pkg/statement" "github.com/jackc/pgx/v5/pgconn" "github.com/stretchr/testify/assert" ) @@ -42,3 +43,12 @@ func TestRowSecurityErrorClassification(t *testing.T) { assert.Equal(t, CodeRowSecurityOutcomeUnknown, OutcomeCode(err)) }) } + +func TestRowSecurityAdmissionErrorsArePermanent(t *testing.T) { + for _, cause := range []error{statement.ErrPolicyRelationDependency, statement.ErrRowSecurityDeclaration, &pgconn.PgError{Code: "42501"}} { + err := rowSecurityError(t.Context(), t.Context(), cause, Budget{}) + assert.ErrorIs(t, err, cause) + assert.Equal(t, CodeRowSecurityRefused, OutcomeCode(err)) + assert.True(t, OutcomeCode(err).Permanent()) + } +} diff --git a/pkg/executor/row_security_owner.go b/pkg/executor/row_security_owner.go new file mode 100644 index 0000000..547af9d --- /dev/null +++ b/pkg/executor/row_security_owner.go @@ -0,0 +1,31 @@ +package executor + +import ( + "context" + "errors" + "fmt" + + "github.com/jackc/pgx/v5" +) + +// ErrRowSecurityRefused means the declaration, target shape, or privileges must +// change before retrying. Wrapped causes retain the specific admission failure. +var ErrRowSecurityRefused = errors.New("row security change refused") + +func checkRowSecurityOwner(ctx context.Context, tx pgx.Tx, schema, table string) error { + var owner bool + err := tx.QueryRow(ctx, `SELECT pg_catalog.pg_has_role(current_user, c.relowner, 'USAGE') + FROM pg_catalog.pg_class c JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace + WHERE n.nspname = $1 AND c.relname = $2`, schema, table).Scan(&owner) + target := pgx.Identifier{schema, table}.Sanitize() + if errors.Is(err, pgx.ErrNoRows) { + return fmt.Errorf("check row security owner for %s: %w", target, ErrTableNotFound) + } + if err != nil { + return fmt.Errorf("check row security owner for %s: %w", target, err) + } + if !owner { + return fmt.Errorf("row security on %s requires owner privileges: %w", target, ErrRowSecurityRefused) + } + return nil +} diff --git a/pkg/executor/row_security_test.go b/pkg/executor/row_security_test.go index bffaf58..ab1da61 100644 --- a/pkg/executor/row_security_test.go +++ b/pkg/executor/row_security_test.go @@ -20,6 +20,10 @@ func TestExecuteRowSecurityRejectsInvalidInputsBeforeConnecting(t *testing.T) { budget := executor.Budget{LockTimeout: time.Millisecond, StatementTimeout: time.Second} _, err = executor.ExecuteRowSecurity(t.Context(), nil, "public", statement.DesiredWithRowSecurity{}, budget) require.ErrorIs(t, err, statement.ErrRowSecurityDeclaration) + require.Equal(t, executor.CodeRowSecurityRefused, executor.OutcomeCode(err)) + require.True(t, executor.OutcomeCode(err).Permanent()) _, err = executor.ExecuteRowSecurity(t.Context(), nil, "", desired, budget) require.ErrorIs(t, err, statement.ErrRowSecurityDeclaration) + require.Equal(t, executor.CodeRowSecurityRefused, executor.OutcomeCode(err)) + require.True(t, executor.OutcomeCode(err).Permanent()) } From 64d16a9777257b82e2de7247aeb08fdf948debf6 Mon Sep 17 00:00:00 2001 From: Armand Parajon Date: Tue, 22 Sep 2026 21:23:32 -0400 Subject: [PATCH 4/5] executor: cover RLS privilege races and dependencies Signed-off-by: Armand Parajon --- docs/atomic-row-security.md | 7 +- pkg/executor/row_security.go | 10 +- ..._security_dependencies_integration_test.go | 56 +++++++++ pkg/executor/row_security_desired_error.go | 20 +++ .../row_security_desired_error_test.go | 22 ++++ pkg/executor/row_security_owner.go | 12 +- .../row_security_owner_integration_test.go | 116 ++++++++++++++++++ 7 files changed, 231 insertions(+), 12 deletions(-) create mode 100644 pkg/executor/row_security_dependencies_integration_test.go create mode 100644 pkg/executor/row_security_desired_error.go create mode 100644 pkg/executor/row_security_desired_error_test.go create mode 100644 pkg/executor/row_security_owner_integration_test.go diff --git a/docs/atomic-row-security.md b/docs/atomic-row-security.md index f11911c..ec0e0e7 100644 --- a/docs/atomic-row-security.md +++ b/docs/atomic-row-security.md @@ -56,11 +56,12 @@ Lock exhaustion reports `budget-lock-exceeded`; the statement or whole-attempt deadline reports `budget-statement-exceeded`. A missing target reports `table-not-found`. Invalid declarations, unsupported targets, and insufficient privileges report permanent `row-security-refused` outcomes, preserving the underlying -cause. Caller cancellation is kept separate from budget exhaustion. +cause. This includes unresolved policy roles and qualified helper functions during +scratch inspection. Caller cancellation is kept separate from budget exhaustion. The caller needs table-owner privileges and permission to create the temporary -scratch schema. Ownership is checked before locking and checked again under the -lock. Roles and qualified helpers must already exist. Grants, role +scratch schema. Both privileges are checked before locking and checked again under +the lock. Roles and qualified helpers must already exist. Grants, role membership, helper bodies, authentication, and Supabase-managed schemas are outside this operation. Application authorization tests are still needed. Concurrent administration of those dependencies is not serialized by the table lock. diff --git a/pkg/executor/row_security.go b/pkg/executor/row_security.go index 77363c7..8950b7b 100644 --- a/pkg/executor/row_security.go +++ b/pkg/executor/row_security.go @@ -97,8 +97,8 @@ func executeRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, return RowSecurityReport{}, fmt.Errorf("set row security budgets: %w", err) } target := pgx.Identifier{schema, desired.Table()}.Sanitize() - // Reject roles without owner privileges before taking an application-blocking lock. - if err := checkRowSecurityOwner(ctx, tx, schema, desired.Table()); err != nil { + // Reject missing owner or scratch privileges before taking an application-blocking lock. + if err := checkRowSecurityPrivileges(ctx, tx, schema, desired.Table()); err != nil { return RowSecurityReport{}, err } // INV: RS-1 — all live comparison and DDL occur after this exclusive lock. @@ -109,8 +109,8 @@ func executeRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, } return RowSecurityReport{}, fmt.Errorf("lock row security target %s: %w", target, err) } - // Ownership may have changed while waiting for the lock. Recheck it under lock. - if err := checkRowSecurityOwner(ctx, tx, schema, desired.Table()); err != nil { + // Privileges may have changed while waiting for the lock. Recheck them under lock. + if err := checkRowSecurityPrivileges(ctx, tx, schema, desired.Table()); err != nil { return RowSecurityReport{}, err } live, err := schemadiff.IntrospectTx(ctx, tx, schema, desired.Table()) @@ -119,7 +119,7 @@ func executeRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, } wanted, err := schemadiff.IntrospectDesiredWithRowSecurityTx(ctx, tx, desired) if err != nil { - return RowSecurityReport{}, fmt.Errorf("inspect desired row security for %s: %w", target, err) + return RowSecurityReport{}, fmt.Errorf("inspect desired row security for %s: %w", target, classifyRowSecurityDesiredError(err)) } if err := admitRowSecurityTable(schema, live, wanted); err != nil { return RowSecurityReport{}, fmt.Errorf("admit row security target %s: %w: %w", target, ErrRowSecurityRefused, err) diff --git a/pkg/executor/row_security_dependencies_integration_test.go b/pkg/executor/row_security_dependencies_integration_test.go new file mode 100644 index 0000000..d65f7a6 --- /dev/null +++ b/pkg/executor/row_security_dependencies_integration_test.go @@ -0,0 +1,56 @@ +package executor_test + +import ( + "fmt" + "testing" + + "github.com/jackc/pgx/v5/pgconn" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/executor" + "github.com/block/pg-sprite/pkg/schemadiff" +) + +func TestExecuteRowSecurityRefusesMissingPolicyRole(t *testing.T) { + pool, schema := rlsFixture(t) + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + report, err := applyRLS(t, pool, schema, fmt.Sprintf(`CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT TO %s_missing_role USING (true);`, schema)) + assertPermanentRLSDependency(t, err, "42704") + assert.Empty(t, report.Statements) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} + +func TestExecuteRowSecurityRefusesMissingPolicyHelper(t *testing.T) { + pool, schema := rlsFixture(t) + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + report, err := applyRLS(t, pool, schema, fmt.Sprintf(`CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (%s.missing_helper(owner_id));`, schema)) + assertPermanentRLSDependency(t, err, "42883") + assert.Empty(t, report.Statements) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} + +func assertPermanentRLSDependency(t *testing.T, err error, sqlstate string) { + t.Helper() + var pgErr *pgconn.PgError + require.ErrorAs(t, err, &pgErr) + assert.Equal(t, sqlstate, pgErr.Code) + assert.Equal(t, executor.CodeRowSecurityRefused, executor.OutcomeCode(err)) + assert.True(t, executor.OutcomeCode(err).Permanent()) +} diff --git a/pkg/executor/row_security_desired_error.go b/pkg/executor/row_security_desired_error.go new file mode 100644 index 0000000..3ae1063 --- /dev/null +++ b/pkg/executor/row_security_desired_error.go @@ -0,0 +1,20 @@ +package executor + +import ( + "errors" + "fmt" + "strings" + + "github.com/jackc/pgx/v5/pgconn" +) + +// Scratch SQLSTATE class 42 means an invalid definition or insufficient access, +// including unresolved roles, columns, and helper functions. Retrying the same +// declaration cannot fix it. Operational failures retain their original class. +func classifyRowSecurityDesiredError(err error) error { + var pgErr *pgconn.PgError + if errors.As(err, &pgErr) && strings.HasPrefix(pgErr.Code, "42") { + return fmt.Errorf("%w: %w", ErrRowSecurityRefused, err) + } + return err +} diff --git a/pkg/executor/row_security_desired_error_test.go b/pkg/executor/row_security_desired_error_test.go new file mode 100644 index 0000000..0f049d7 --- /dev/null +++ b/pkg/executor/row_security_desired_error_test.go @@ -0,0 +1,22 @@ +package executor + +import ( + "context" + "testing" + + "github.com/jackc/pgx/v5/pgconn" + "github.com/stretchr/testify/assert" +) + +func TestRowSecurityDesiredOperationalErrorsRemainRetryable(t *testing.T) { + for _, err := range []error{ + context.DeadlineExceeded, + &pgconn.PgError{Code: "57014"}, // query cancelled + &pgconn.PgError{Code: "55P03"}, // lock timeout + &pgconn.PgError{Code: "40001"}, // serialization failure + } { + actual := classifyRowSecurityDesiredError(err) + assert.Equal(t, err, actual) + assert.False(t, OutcomeCode(actual).Permanent()) + } +} diff --git a/pkg/executor/row_security_owner.go b/pkg/executor/row_security_owner.go index 547af9d..7204e34 100644 --- a/pkg/executor/row_security_owner.go +++ b/pkg/executor/row_security_owner.go @@ -12,11 +12,12 @@ import ( // change before retrying. Wrapped causes retain the specific admission failure. var ErrRowSecurityRefused = errors.New("row security change refused") -func checkRowSecurityOwner(ctx context.Context, tx pgx.Tx, schema, table string) error { - var owner bool - err := tx.QueryRow(ctx, `SELECT pg_catalog.pg_has_role(current_user, c.relowner, 'USAGE') +func checkRowSecurityPrivileges(ctx context.Context, tx pgx.Tx, schema, table string) error { + var owner, createSchema bool + err := tx.QueryRow(ctx, `SELECT pg_catalog.pg_has_role(current_user, c.relowner, 'USAGE'), + pg_catalog.has_database_privilege(current_user, pg_catalog.current_database(), 'CREATE') FROM pg_catalog.pg_class c JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace - WHERE n.nspname = $1 AND c.relname = $2`, schema, table).Scan(&owner) + WHERE n.nspname = $1 AND c.relname = $2`, schema, table).Scan(&owner, &createSchema) target := pgx.Identifier{schema, table}.Sanitize() if errors.Is(err, pgx.ErrNoRows) { return fmt.Errorf("check row security owner for %s: %w", target, ErrTableNotFound) @@ -27,5 +28,8 @@ func checkRowSecurityOwner(ctx context.Context, tx pgx.Tx, schema, table string) if !owner { return fmt.Errorf("row security on %s requires owner privileges: %w", target, ErrRowSecurityRefused) } + if !createSchema { + return fmt.Errorf("row security on %s requires CREATE on the database for scratch inspection: %w", target, ErrRowSecurityRefused) + } return nil } diff --git a/pkg/executor/row_security_owner_integration_test.go b/pkg/executor/row_security_owner_integration_test.go new file mode 100644 index 0000000..76f063e --- /dev/null +++ b/pkg/executor/row_security_owner_integration_test.go @@ -0,0 +1,116 @@ +package executor_test + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/executor" + "github.com/block/pg-sprite/pkg/schemadiff" + "github.com/block/pg-sprite/pkg/statement" +) + +func rlsOwnerPool(t *testing.T, admin *pgxpool.Pool, schema string) (*pgxpool.Pool, string) { + t.Helper() + role := pgx.Identifier{schema + "_owner"}.Sanitize() + _, err := admin.Exec(t.Context(), fmt.Sprintf(`CREATE ROLE %s; + GRANT USAGE ON SCHEMA %s TO %s; + ALTER TABLE %s.documents OWNER TO %s;`, role, schema, role, schema, role)) + require.NoError(t, err) + t.Cleanup(func() { + _, err := admin.Exec(context.WithoutCancel(t.Context()), "DROP OWNED BY "+role+" CASCADE; DROP ROLE "+role) + require.NoError(t, err) + }) + cfg := admin.Config() + cfg.AfterConnect = func(ctx context.Context, conn *pgx.Conn) error { + _, err := conn.Exec(ctx, "SET ROLE "+role) + return err + } + pool, err := pgxpool.NewWithConfig(t.Context(), cfg) + require.NoError(t, err) + t.Cleanup(pool.Close) + return pool, role +} + +func TestExecuteRowSecurityRefusesMissingScratchPrivilegeBeforeLock(t *testing.T) { + admin, schema := rlsFixture(t) + pool, _ := rlsOwnerPool(t, admin, schema) + var canCreate bool + require.NoError(t, pool.QueryRow(t.Context(), "SELECT has_database_privilege(current_user, current_database(), 'CREATE')").Scan(&canCreate)) + require.False(t, canCreate, "the fixture role must lack database CREATE") + holder, err := admin.Begin(t.Context()) + require.NoError(t, err) + defer func() { _ = holder.Rollback(context.WithoutCancel(t.Context())) }() + _, err = holder.Exec(t.Context(), "LOCK TABLE "+pgx.Identifier{schema, "documents"}.Sanitize()+" IN ACCESS SHARE MODE") + require.NoError(t, err) + // The lock would time out if scratch privileges were checked only after locking. + report, err := applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents DISABLE ROW LEVEL SECURITY;`) + require.ErrorIs(t, err, executor.ErrRowSecurityRefused) + assert.Equal(t, executor.CodeRowSecurityRefused, executor.OutcomeCode(err)) + assert.True(t, executor.OutcomeCode(err).Permanent()) + assert.Empty(t, report.Statements) +} + +func TestExecuteRowSecurityRechecksOwnerAfterLockWait(t *testing.T) { + admin, schema := rlsFixture(t) + pool, role := rlsOwnerPool(t, admin, schema) + var database string + require.NoError(t, admin.QueryRow(t.Context(), "SELECT current_database()").Scan(&database)) + _, err := admin.Exec(t.Context(), "GRANT CREATE ON DATABASE "+pgx.Identifier{database}.Sanitize()+" TO "+role) + require.NoError(t, err) + before, err := schemadiff.Introspect(t.Context(), admin, schema, "documents") + require.NoError(t, err) + // This already matches. Without the recheck the no-op would succeed, so the + // test cannot accidentally pass because PostgreSQL rejects a later owner-only DDL. + desired, err := statement.ParseDesiredWithRowSecurity(`CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (owner_id = 7);`) + require.NoError(t, err) + holder, err := admin.Begin(t.Context()) + require.NoError(t, err) + defer func() { _ = holder.Rollback(context.WithoutCancel(t.Context())) }() + target := pgx.Identifier{schema, "documents"}.Sanitize() + _, err = holder.Exec(t.Context(), "LOCK TABLE "+target+" IN ACCESS EXCLUSIVE MODE") + require.NoError(t, err) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + done := make(chan error, 1) + go func() { + _, err := executor.ExecuteRowSecurity(ctx, pool, schema, desired, executor.Budget{LockTimeout: 3 * time.Second, StatementTimeout: 5 * time.Second}) + done <- err + }() + const queueDeadline = 2 * time.Second + const queuePoll = 10 * time.Millisecond + require.Eventually(t, func() bool { + var waiting bool + err := admin.QueryRow(t.Context(), `SELECT EXISTS ( + SELECT 1 FROM pg_locks WHERE relation = $1::regclass + AND mode = 'AccessExclusiveLock' AND NOT granted + )`, target).Scan(&waiting) + return err == nil && waiting + }, queueDeadline, queuePoll) + // Transfer to the administrator but retain enough privilege for the queued + // LOCK to succeed. Only the executor's post-lock owner check must refuse. + _, err = holder.Exec(t.Context(), "ALTER TABLE "+target+" OWNER TO CURRENT_USER; GRANT SELECT, UPDATE ON "+target+" TO "+role) + require.NoError(t, err) + require.NoError(t, holder.Commit(t.Context())) + err = <-done + require.ErrorIs(t, err, executor.ErrRowSecurityRefused) + assert.Equal(t, executor.CodeRowSecurityRefused, executor.OutcomeCode(err)) + after, err := schemadiff.Introspect(t.Context(), admin, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} From 17290fbcf2f7ca1fb886920f1a1373e126b44eb9 Mon Sep 17 00:00:00 2001 From: Armand Parajon Date: Wed, 23 Sep 2026 13:36:32 -0400 Subject: [PATCH 5/5] executor: pin RLS commit guards and shorten locking Signed-off-by: Armand Parajon --- docs/atomic-row-security.md | 16 ++- docs/capabilities.md | 3 + pkg/executor/row_security.go | 9 +- .../row_security_commit_integration_test.go | 114 ++++++++++++++++++ ..._security_dependencies_integration_test.go | 8 ++ 5 files changed, 140 insertions(+), 10 deletions(-) create mode 100644 pkg/executor/row_security_commit_integration_test.go diff --git a/docs/atomic-row-security.md b/docs/atomic-row-security.md index ec0e0e7..3901f2e 100644 --- a/docs/atomic-row-security.md +++ b/docs/atomic-row-security.md @@ -39,9 +39,10 @@ Changing an RLS definition can widen access even when no data is deleted. 1. Validate the parsed declaration and nonzero budgets. 2. Begin one transaction with a deadline for the entire attempt and transaction-local - lock and statement timeouts. Acquire `ACCESS EXCLUSIVE` on the target. -3. Read the live definition and materialize the desired SQL in a rolled-back - savepoint on the same connection. Refuse unsupported table shapes or any table delta. + lock and statement timeouts. Materialize the desired SQL in a rolled-back savepoint + on the same connection before blocking the target. +3. Acquire `ACCESS EXCLUSIVE`, recheck privileges, and read the live definition. + Refuse unsupported table shapes or any table delta. 4. If RLS already matches, finish without policy DDL. Otherwise replace the complete policy set (including unchanged policies), apply ENABLE/DISABLE and FORCE/NO FORCE, and preserve policy comments. Live SQL is rendered from the inspected desired catalog, not replayed from the input. @@ -50,10 +51,13 @@ Changing an RLS definition can widen access even when no data is deleted. The lock blocks reads and writes briefly; this is bounded metadata DDL, not an online copy. Other sessions never see the intermediate policy set. A failure before commit rolls back every change. A lost commit response is an unknown outcome: -inspect the database before retrying. There are no automatic retries. +inspect the database before retrying. This is conservative: any commit error is +reported as unknown, even if cancellation may have prevented COMMIT from being sent. +There are no automatic retries. -Lock exhaustion reports `budget-lock-exceeded`; the statement or whole-attempt -deadline reports `budget-statement-exceeded`. A missing target reports +Expiration of PostgreSQL’s lock timer reports `budget-lock-exceeded`. The statement +or whole-attempt deadline reports `budget-statement-exceeded`, including when the +whole-attempt deadline expires during a lock wait. The first limit reached wins. A missing target reports `table-not-found`. Invalid declarations, unsupported targets, and insufficient privileges report permanent `row-security-refused` outcomes, preserving the underlying cause. This includes unresolved policy roles and qualified helper functions during diff --git a/docs/capabilities.md b/docs/capabilities.md index d158028..2fed70f 100644 --- a/docs/capabilities.md +++ b/docs/capabilities.md @@ -29,6 +29,9 @@ refused form would take, what an operator who accepts a maintenance window can d - [Why typed refusal, not passthrough](#why-typed-refusal-not-passthrough) - [Deliberately operator-owned](#deliberately-operator-owned) +A ✅ marks an implemented capability; check the front-door columns for CLI access. +The atomic RLS executor is currently a Go API only, with no `migrate` or `diff` execution. + ## Query the matrix The marker-delimited regions of this page are generated from diff --git a/pkg/executor/row_security.go b/pkg/executor/row_security.go index 8950b7b..51f1950 100644 --- a/pkg/executor/row_security.go +++ b/pkg/executor/row_security.go @@ -101,6 +101,11 @@ func executeRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, if err := checkRowSecurityPrivileges(ctx, tx, schema, desired.Table()); err != nil { return RowSecurityReport{}, err } + // Materialize the declaration before blocking application reads and writes. + wanted, err := schemadiff.IntrospectDesiredWithRowSecurityTx(ctx, tx, desired) + if err != nil { + return RowSecurityReport{}, fmt.Errorf("inspect desired row security for %s: %w", target, classifyRowSecurityDesiredError(err)) + } // INV: RS-1 — all live comparison and DDL occur after this exclusive lock. if _, err := tx.Exec(ctx, "LOCK TABLE ONLY "+target+" IN ACCESS EXCLUSIVE MODE"); err != nil { var pgErr *pgconn.PgError @@ -117,10 +122,6 @@ func executeRowSecurity(ctx context.Context, pool *pgxpool.Pool, schema string, if err != nil { return RowSecurityReport{}, fmt.Errorf("inspect row security target %s: %w", target, err) } - wanted, err := schemadiff.IntrospectDesiredWithRowSecurityTx(ctx, tx, desired) - if err != nil { - return RowSecurityReport{}, fmt.Errorf("inspect desired row security for %s: %w", target, classifyRowSecurityDesiredError(err)) - } if err := admitRowSecurityTable(schema, live, wanted); err != nil { return RowSecurityReport{}, fmt.Errorf("admit row security target %s: %w: %w", target, ErrRowSecurityRefused, err) } diff --git a/pkg/executor/row_security_commit_integration_test.go b/pkg/executor/row_security_commit_integration_test.go new file mode 100644 index 0000000..87b89f6 --- /dev/null +++ b/pkg/executor/row_security_commit_integration_test.go @@ -0,0 +1,114 @@ +package executor_test + +import ( + "context" + "fmt" + "testing" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/executor" + "github.com/block/pg-sprite/pkg/schemadiff" +) + +// Fire only for live CREATE POLICY, never for the disposable scratch declaration. +func onLiveRLSPolicy(t *testing.T, pool *pgxpool.Pool, schema, action string) { + t.Helper() + trigger := pgx.Identifier{schema + "_policy_fault"}.Sanitize() + _, err := pool.Exec(t.Context(), fmt.Sprintf(`CREATE FUNCTION %s.policy_fault() RETURNS event_trigger + LANGUAGE plpgsql AS $$ + BEGIN + IF EXISTS ( + SELECT 1 FROM pg_event_trigger_ddl_commands() d + JOIN pg_policy p ON d.classid = 'pg_policy'::regclass AND d.objid = p.oid + JOIN pg_class c ON c.oid = p.polrelid + JOIN pg_namespace n ON n.oid = c.relnamespace + WHERE n.nspname = '%s' AND d.command_tag = 'CREATE POLICY' + ) THEN + %s + END IF; + END; + $$; + CREATE EVENT TRIGGER %s ON ddl_command_end EXECUTE FUNCTION %s.policy_fault();`, schema, schema, action, trigger, schema)) + require.NoError(t, err) + t.Cleanup(func() { + _, err := pool.Exec(context.WithoutCancel(t.Context()), "DROP EVENT TRIGGER "+trigger) + require.NoError(t, err) + }) +} + +func TestExecuteRowSecurityRollsBackPolicyDivergence(t *testing.T) { + pool, schema := rlsFixture(t) + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + onLiveRLSPolicy(t, pool, schema, fmt.Sprintf(`ALTER POLICY readers ON %s.documents USING (false);`, schema)) + report, err := applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (owner_id = 9);`) + require.ErrorIs(t, err, executor.ErrInvariantViolation) + assert.Equal(t, executor.CodeInvariantViolation, executor.OutcomeCode(err)) + assert.Empty(t, report.Statements) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} + +func TestExecuteRowSecurityRollsBackTableDivergence(t *testing.T) { + pool, schema := rlsFixture(t) + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + onLiveRLSPolicy(t, pool, schema, fmt.Sprintf(`CREATE TABLE %s.unexpected_reference ( + document_id bigint REFERENCES %s.documents(id) + );`, schema, schema)) + report, err := applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (owner_id = 9);`) + require.ErrorIs(t, err, executor.ErrInvariantViolation) + assert.Empty(t, report.Statements) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} + +func TestExecuteRowSecurityWrapsCommitFailure(t *testing.T) { + pool, schema := rlsFixture(t) + before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + // Defer a real server error until COMMIT, after all convergence checks pass. + _, err = pool.Exec(t.Context(), fmt.Sprintf(`CREATE TABLE %s.commit_probe (id integer); + CREATE FUNCTION %s.fail_commit() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + RAISE EXCEPTION 'injected deferred commit failure'; + END; + $$; + CREATE CONSTRAINT TRIGGER fail_commit AFTER INSERT ON %s.commit_probe + DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION %s.fail_commit();`, schema, schema, schema, schema)) + require.NoError(t, err) + onLiveRLSPolicy(t, pool, schema, fmt.Sprintf(`INSERT INTO %s.commit_probe VALUES (1);`, schema)) + report, err := applyRLS(t, pool, schema, `CREATE TABLE documents ( + id bigint PRIMARY KEY, + owner_id bigint NOT NULL + ); + ALTER TABLE documents ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON documents FOR SELECT USING (owner_id = 9);`) + var unknown *executor.RowSecurityOutcomeUnknownError + require.ErrorAs(t, err, &unknown) + var pgErr *pgconn.PgError + require.ErrorAs(t, err, &pgErr) + assert.Equal(t, "P0001", pgErr.Code) + assert.Equal(t, executor.CodeRowSecurityOutcomeUnknown, executor.OutcomeCode(err)) + assert.Empty(t, report.Statements) + after, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") + require.NoError(t, err) + assert.Equal(t, before, after) +} diff --git a/pkg/executor/row_security_dependencies_integration_test.go b/pkg/executor/row_security_dependencies_integration_test.go index d65f7a6..909be2d 100644 --- a/pkg/executor/row_security_dependencies_integration_test.go +++ b/pkg/executor/row_security_dependencies_integration_test.go @@ -1,9 +1,11 @@ package executor_test import ( + "context" "fmt" "testing" + "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgconn" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -16,6 +18,12 @@ func TestExecuteRowSecurityRefusesMissingPolicyRole(t *testing.T) { pool, schema := rlsFixture(t) before, err := schemadiff.Introspect(t.Context(), pool, schema, "documents") require.NoError(t, err) + holder, err := pool.Begin(t.Context()) + require.NoError(t, err) + defer func() { _ = holder.Rollback(context.WithoutCancel(t.Context())) }() + _, err = holder.Exec(t.Context(), "LOCK TABLE "+pgx.Identifier{schema, "documents"}.Sanitize()+" IN ACCESS SHARE MODE") + require.NoError(t, err) + // Invalid desired SQL must be refused without waiting for the target lock. report, err := applyRLS(t, pool, schema, fmt.Sprintf(`CREATE TABLE documents ( id bigint PRIMARY KEY, owner_id bigint NOT NULL