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 1c35db0..ab01189 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, 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 new file mode 100644 index 0000000..3901f2e --- /dev/null +++ b/docs/atomic-row-security.md @@ -0,0 +1,86 @@ +# 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. 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. +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. 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. + +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 +scratch inspection. Caller cancellation is kept separate from budget exhaustion. + +The caller needs table-owner privileges and permission to create the temporary +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. + +## Invariants and tests + +- **RS-1:** Read the live baseline only after taking the target lock; never accept a + 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 + 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..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 @@ -176,7 +179,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 +272,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, 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/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..7981932 100644 --- a/docs/execution-model.md +++ b/docs/execution-model.md @@ -289,6 +289,8 @@ 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 | | `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 625bd38..733f59b 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) @@ -396,6 +397,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 | `RenderRowSecurity` from the scratch catalog and executor readback; qualified-helper, quoted-name, 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..29ce674 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_safer_sequence" + 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..6d515e1 100644 --- a/pkg/executor/code.go +++ b/pkg/executor/code.go @@ -23,6 +23,10 @@ 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. CodeBlockingOutcomeUnknown Code = "blocking-outcome-unknown" // CodeInvalidBlockingBudget identifies an unrepresentable or disabled bound. @@ -148,6 +152,8 @@ func Codes() []Code { CodeBudgetLockExceeded, CodeBudgetStatementExceeded, CodeBlockingOutcomeUnknown, + CodeRowSecurityRefused, + CodeRowSecurityOutcomeUnknown, CodeInvalidBlockingBudget, CodeUnsupportedAcceptedBlocking, CodeCancelledByCaller, @@ -195,7 +201,7 @@ func Codes() []Code { // permanent. func (c Code) Permanent() bool { switch c { - case CodeInvalidIndexOtherTable, + case CodeRowSecurityRefused, CodeInvalidIndexOtherTable, CodeInvalidIndexNotDroppable, CodeInvalidBlockingBudget, CodeUnsupportedAcceptedBlocking, @@ -242,6 +248,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 @@ -302,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 5c2e4ec..cdf99f3 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,8 @@ 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, executor.CodeCancelledByCaller: false, diff --git a/pkg/executor/row_security.go b/pkg/executor/row_security.go new file mode 100644 index 0000000..51f1950 --- /dev/null +++ b/pkg/executor/row_security.go @@ -0,0 +1,182 @@ +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" + "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 + } + if desired.Table() == "" || schema == "" { + 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) + 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 + } + 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 +} + +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) + } + 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() + // 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 + } + // 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 + 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) + } + // 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()) + if err != nil { + return RowSecurityReport{}, fmt.Errorf("inspect row security target %s: %w", target, err) + } + if err := admitRowSecurityTable(schema, live, wanted); err != nil { + 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{}, 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{}, fmt.Errorf("render row security for %s: %w: %w", target, ErrRowSecurityRefused, err) + } + 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) + } + } + } + // 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{}, 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) + } + 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_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_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 new file mode 100644 index 0000000..909be2d --- /dev/null +++ b/pkg/executor/row_security_dependencies_integration_test.go @@ -0,0 +1,64 @@ +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" + + "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) + 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 + ); + 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_integration_test.go b/pkg/executor/row_security_integration_test.go new file mode 100644 index 0000000..7b295c4 --- /dev/null +++ b/pkg/executor/row_security_integration_test.go @@ -0,0 +1,381 @@ +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.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 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) + 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) +} + +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.ErrorIs(t, err, executor.ErrTableNotFound) + assert.Equal(t, executor.CodeTableNotFound, executor.OutcomeCode(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.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)) + 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) +} + +// 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..433ee8a --- /dev/null +++ b/pkg/executor/row_security_internal_test.go @@ -0,0 +1,54 @@ +package executor + +import ( + "context" + "testing" + "time" + + "github.com/block/pg-sprite/pkg/statement" + "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)) + }) +} + +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..7204e34 --- /dev/null +++ b/pkg/executor/row_security_owner.go @@ -0,0 +1,35 @@ +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 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, &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) + } + 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) + } + 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) +} diff --git a/pkg/executor/row_security_test.go b/pkg/executor/row_security_test.go new file mode 100644 index 0000000..ab1da61 --- /dev/null +++ b/pkg/executor/row_security_test.go @@ -0,0 +1,29 @@ +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) + 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()) +} 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..1132b15 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. @@ -35,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 { @@ -107,16 +149,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 {