diff --git a/docs/declarative-row-security.md b/docs/declarative-row-security.md index c75fbe8..280d6cf 100644 --- a/docs/declarative-row-security.md +++ b/docs/declarative-row-security.md @@ -279,8 +279,10 @@ this table-scoped work. 4. **Prove application behavior.** The local Supabase harness now checks real PostgREST requests from two authenticated users and anonymous callers, allowed and denied writes, changed visibility, and rollback after a cancelled apply. - Hosted connection and privilege validation remains a follow-up on a disposable - project; local results do not establish hosted support. + The [hosted suite](../integration/supabase/hosted/README.md) also exercised + preview, apply, convergence, tenant access, and atomic rollback as the owning + `postgres` role on a disposable project. Hosted non-owner privilege validation + remains a follow-up; owner-role results do not establish that boundary. The [inspection tests](../pkg/schemadiff/row_security_integration_test.go), [round-trip tests](../pkg/schemadiff/row_security_roundtrip_integration_test.go), and diff --git a/docs/supabase.md b/docs/supabase.md index 5090c9d..455ea40 100644 --- a/docs/supabase.md +++ b/docs/supabase.md @@ -4,9 +4,12 @@ Your Supabase app runs on PostgreSQL. pg-sprite helps you change its tables as you build: add a column for a new feature, or add an index as your queries grow. Local tests cover column additions and concurrent index builds alongside -Supabase's access policies, Data API, and Realtime subscriptions. Hosted projects -are the next validation step. Changes that need a replacement table are refused -today because copy-and-swap is not implemented yet. +Supabase's access policies, Data API, and Realtime subscriptions. An opt-in [hosted suite](../integration/supabase/hosted/README.md) also exercises +real Auth users and service endpoints. Its Realtime schema change cases initialize +one shared fixture and verify baseline delivery before applying changes, then +check that events and tenant isolation survive on the same connections. Changes +that need a replacement table are refused today because copy-and-swap is not +implemented yet. Start with [your first change](#make-your-first-change). The [test results](#what-works-today) and [roadmap](#where-we-go-next) show how far the current coverage goes. @@ -32,7 +35,7 @@ and the [Data API](https://supabase.com/docs/guides/api) for more detail. This walkthrough adds a nullable `title` column to an existing `public.documents` table. Substitute your own app table and column. Start on a development project; -the compatibility results below come from local Supabase services. +the capability matrix below describes the pinned local Supabase services. ### Install pg-sprite @@ -57,7 +60,7 @@ See [engine-role.md](engine-role.md) for the operation-specific grants. Supavisor offers two pooling modes: - **Session mode** keeps the same PostgreSQL connection for the client's session. - The tested session endpoint works with pg-sprite and can help on IPv4-only networks + The locally tested session endpoint works with pg-sprite and can help on IPv4-only networks - **Transaction mode** can assign a different PostgreSQL connection after each transaction. Do not use it for pg-sprite: execution limits need a stable session @@ -74,9 +77,9 @@ Replace the example value with your connection string. This sets an environment variable without opening a connection. Keep credentials out of source control; your secret manager can also set this variable for you or your agent. -For hosted connections, use `sslmode=verify-full` in the URL to verify the server's -certificate and hostname. If you need to supply a CA certificate separately, -save the certificate for your project and point pg-sprite at it: +For certificate and hostname verification, download the CA from **Database +Settings → SSL Configuration** and follow [Supabase's verification instructions](https://supabase.com/docs/guides/platform/ssl-enforcement#a-note-about-postgres-ssl-modes). +Use `sslmode=verify-full` with `sslrootcert` in the URL, or point pg-sprite at the CA: ```sh export PGSPRITE_CA_CERT='/absolute/path/to/project-ca.crt' @@ -84,8 +87,10 @@ export PGSPRITE_CA_CERT='/absolute/path/to/project-ca.crt' That variable sets the certificate file used by the commands below. A certificate error should be fixed by checking the hostname and trusted certificate, rather -than disabling verification. Hosted certificate handling remains a validation -milestone for this guide. +than disabling verification. Hosted checks passed with the default URI, +`sslmode=require`, and `verify-full` with the downloaded CA. The default connection +used TLS, which alone does not establish server identity verification. An unrelated +CA and a mismatched expected hostname were rejected in the [hosted TLS tests](../integration/supabase/hosted/tls_test.go). ### Preview the change @@ -189,7 +194,9 @@ The result lists the statements committed in one transaction. Applying it again reports that row security already matches. A mixed column/index and RLS edit refuses without committing either part. See [output and failure handling](declarative-row-security.md#apply-the-declaration). This CLI route is covered by PostgreSQL integration tests; local Supabase API -coverage uses the same executor. Hosted validation remains a separate step. +coverage uses the same executor. The [hosted suite](../integration/supabase/hosted/README.md) +also exercised preview, apply, convergence, and atomic rollback as the owning +`postgres` role. Hosted non-owner privilege validation remains a follow-up. ## What works today @@ -257,26 +264,31 @@ outcome; they do not assume every unsuccessful change rolls back completely. - Realtime coverage is limited to INSERT/UPDATE subscriptions during the tested native changes. Deletes, reconnect recovery, column removal, and table replacement need separate validation -- PostgREST cache refresh was exercised with the image's schema-change event triggers; +- PostgREST cache refresh was exercised with the image's schema change event triggers; a deployment without those triggers needs its own reload workflow -Hosted role configuration and TLS remain unverified. The Auth -service runs its schema initialization; JWTs are signed by the test fixture, -so this does not test signup or login. These local results are not an -unrestricted Supabase support claim. +Hosted checks have passed as the owning `postgres` role on PostgreSQL 17.6, +including CA-based TLS verification, the CLI schema lifecycle, real Auth defaults +and foreign keys, and atomic RLS rollback. Hosted poolers remain untested and +Realtime checks cover continuity after fixture initialization; see the +[hosted suite](../integration/supabase/hosted/README.md) for cases and limits. +In the local suite, Auth runs its schema initialization and the fixture signs JWTs, +so those cases do not test login. Hosted cases create +confirmed test users through the Auth admin API and sign in with their passwords; +they do not test public signup, email delivery, or OAuth. Neither suite establishes +unrestricted Supabase support. ## Where we go next -The next milestones build on the local tests. Each needs repeatable evidence +The next milestones build on the local and scoped hosted tests. Each needs repeatable evidence before we expand the support claim: -1. **Validate hosted projects.** Run the same checks on a disposable Supabase - project, including certificate verification, network access, and hosted roles -2. **Validate declarative schema workflows.** Export an existing Supabase schema, - edit the desired SQL files, preview the diff, and apply supported changes. - Verify that the live schema matches the files and a second diff is empty, - while access policies and Realtime subscriptions still work. Make clear which - objects the files describe and which remain managed separately +1. **Finish hosted validation.** Core direct-connection and TLS cases have passed; + validate the hosted pooler endpoints +2. **Publish a reproducible declarative workflow.** The hosted CLI cases cover + create, export, edit, apply, RLS, and convergence. Complete a fresh-project + walkthrough from the published guide, with clear boundaries for managed objects + and Realtime behavior 3. **Cover more app workflows.** Exercise deletes and reconnects in Realtime, Realtime payloads after column renames and removals, and real signup/login flows. Make the limits of desired schema files and access-policy handling diff --git a/integration/supabase/README.md b/integration/supabase/README.md index 34d0dcf..a497a5d 100644 --- a/integration/supabase/README.md +++ b/integration/supabase/README.md @@ -100,3 +100,9 @@ Handle unrelated upgrades in a separate follow-up. When making an upgrade: This check is triggered by agent work, not a scheduled update bot. Fresh CI runners pull the pinned images on each run, which also exposes unavailable pins. + +## Hosted projects + +The separate [hosted suite](hosted/README.md) uses real Auth users and managed +endpoints with an explicit opt-in. Do not redirect this local fixture at a hosted +project: its fixed names, JWT signer, and Docker lifecycle are local-only. diff --git a/integration/supabase/hosted/README.md b/integration/supabase/hosted/README.md new file mode 100644 index 0000000..6a23907 --- /dev/null +++ b/integration/supabase/hosted/README.md @@ -0,0 +1,149 @@ +# Hosted Supabase checks + +Use this suite on a **disposable hosted project**. It creates real Auth users, +application tables, policies, and publication entries, then removes its fixtures. +It leaves project settings alone. No Docker services or locally signed JWTs are used. + +This is an opt-in complement to the [local suite](../README.md), not a replacement +for its larger DDL matrix or a required hosted CI job. Realtime checks establish +working delivery before applying schema changes, then verify continuity on the +same table and connections. Setup failures fail the test; writes are never retried. + +## Run it + +Create a disposable Supabase project and copy its **Direct connection** details +from **Connect**. Use a private PostgreSQL password file or your secret manager; +never check credentials into the repository. Keep the creation defaults if you +want to test the normal new-project experience. + +In **Settings → API Keys**, copy the publishable and secret keys into a local JSON +file with mode `0600`. The secret key is used only to create and delete test users; +application requests use the publishable key and real user access tokens. + +```json +{ + "publishable_key": "YOUR_PUBLISHABLE_KEY", + "secret_key": "YOUR_SECRET_KEY" +} +``` + +Set the paths and project details, then run the suite from the repository root: + +```sh +export SUPABASE_HOSTED_TEST=1 +export SUPABASE_HOSTED_URL='https://YOUR_PROJECT_REF.supabase.co' +export SUPABASE_HOSTED_CREDENTIALS='/absolute/path/to/credentials.json' +export PGSPRITE_URL='postgresql://postgres@db.YOUR_PROJECT_REF.supabase.co:5432/postgres' +export PGPASSFILE='/absolute/path/to/pgpass' +go build -o ./bin/pg-sprite ./cmd/pg-sprite +export SUPABASE_HOSTED_BIN="$PWD/bin/pg-sprite" +go test -count=1 -timeout=10m -v ./integration/supabase/hosted +``` + +A successful case prints `--- PASS: TestHostedDeclarativeRLS`. With the opt-in +unset, cases print a skip reason and make no database or API requests. The baseline +requires a direct hostname matching the API project. Do not run multiple copies +against the same project: publication changes affect its shared Realtime service. + +To test the session pooler, set `SUPABASE_HOSTED_SESSION_URL` to its URL from +**Connect** and add a password-file entry for that exact endpoint. Its host, +port, and username differ from the direct endpoint. A missing URL produces an +explicit skip. + +Transaction-pooler refusal stays in the [controlled local suite](../../../pkg/dbconn/supabase_integration_test.go). +An idle hosted transaction pooler can reuse one backend and pass the session +probe, so a hosted assertion cannot reliably prove refusal. A passing probe +does not make transaction pooling supported; use direct or session connections. + +To run the CLI lifecycle independently of Realtime and pooler tests: + +```sh +go test -count=1 -timeout=5m -v ./integration/supabase/hosted -run '^TestHostedCLI' +``` + +These cases do not change Realtime publication membership. Passing them establishes +only the capabilities below; it does not clear failures in the separate Realtime +cases. For example, a successful column case prints +`--- PASS: TestHostedCLIAddColumn`. + +## TLS checks without API keys + +TLS tests use only `PGSPRITE_URL`, `PGPASSFILE`, and the opt-in. Supply the direct +URL with `sslmode=verify-full` and `sslrootcert` pointing at the CA downloaded from +**Database Settings → SSL Configuration**. Then run: + +```sh +go test -count=1 -timeout=2m -v ./integration/supabase/hosted -run '^TestHostedTLS' +``` + +Expected cases are `TestHostedTLSDefault`, `TestHostedTLSRequire`, +`TestHostedTLSVerifyFull`, `TestHostedTLSRejectUntrustedCA`, and +`TestHostedTLSRejectWrongHostname`. Successful connections assert encryption using +`pg_stat_ssl`; negative cases require the specific certificate-trust or hostname +error, not just any failed connection. Verification cases skip explicitly when +`sslrootcert` is absent. The hostname negative case changes the expected TLS name +while still dialing the real database; it uses the underlying pgx driver to inject +that mismatch. These tests make no schema changes and do not test server-side +rejection of plaintext connections. + +## What the cases establish + +| Case | Passing result | +| --- | --- | +| `TestHostedTableCleanup` | Cleanup handles absent tables and tables committed before a later error | +| `TestHostedCLIAuthDefault` | A desired `auth.uid()` default uses the real HTTP caller; spoofing another owner is denied | +| `TestHostedCLIAuthForeignKey` | A validated reference to real Auth users rejects orphan writes | +| `TestHostedCLIUniqueIndexFailure` | Duplicate data causes a typed failure and a reported invalid index; rows and access survive | +| `TestHostedRLSCancellationPreservesAccess` | Cancellation after a live policy drop rolls back the transaction; read and write isolation survive | +| `TestHostedCLICreateExportAndRLS` | Desired table creation, export/diff convergence, explicit RLS preview/apply, real tenant API isolation, and an unchanged second apply | +| `TestHostedCLIAddColumn` | A defaulted column is applied without losing existing data or RLS access controls | +| `TestHostedCLIAddIndex` | The new index is valid; table identity, rows, policies, and grants remain unchanged | +| `TestHostedCLIFailedNotNull` | Failed validation preserves NULL data and tenant access, leaves the column nullable, and reports the retained unvalidated helper constraint | +| `TestHostedCLIRefuseDropColumn` | Destructive preview and apply are refused without dropping a column | +| `TestHostedCLIRefuseCopySwap` | The unavailable copy-and-swap plan is refused before its safe prefix can run | +| `TestHostedCLIRLSLockFailure` | A lock-budget failure leaves existing policies and tenant access intact | +| `TestHostedRealtimeContinuity/initialize` | An active reader and a delivered baseline for both tenants precede all pg-sprite DDL | +| `TestHostedRealtimeContinuity/no_DDL_control` | Both tenants receive updates before schema changes | +| `TestHostedRealtimeContinuity/refuse_volatile_UUID` | Imperative and desired UUID-default changes are refused without applying their safe prefix; data and events survive | +| `TestHostedRealtimeContinuity/refuse_volatile_timestamp` | The same refusal checks hold for `clock_timestamp()` | +| `TestHostedRealtimeContinuity/add_column` | A nullable column preserves events and API isolation on the same table and sockets | +| `TestHostedRealtimeContinuity/add_index` | A concurrent index preserves events, API isolation, and table identity | +| `TestHostedDeclarativeRLS` | Policy changes alter real users' HTTP access; removing the last policy denies reads while retaining the data | +| `TestHostedSessionPooler` | A column change succeeds through the supplied session endpoint | + +A passing refusal case means **safe rejection**, not support for executing that +DDL. None of these tests establishes support for every Supabase feature or plan. + +Each test uses random fixture names and fails on a collision instead of deleting +an existing table. Cleanup has its own bounded context and deletes only objects +created by that case. If the process is forcibly terminated, inspect the test's +`pgsprite_hosted_…` tables and `pgsprite-…@example.com` users before removing any +leftovers. Do not reset `public` or delete unrelated Auth users. + +## Realtime fixture lifecycle + +`TestHostedRealtimeContinuity` keeps one table, publication membership, two users, +and two sockets for all cases. Setup waits for an active wal2json reader, verifies +the authenticated subscriptions, then sends one baseline update and requires +delivery to both users. Initialization failures fail the test before DDL runs. + +The publication-startup budget is 75 seconds; each event must arrive within 30 +seconds. Protocol heartbeats keep sockets alive during setup. There are no fixed +readiness sleeps, retried writes, or reconnects. This setup predicate is specific +to the fixture, not a public Supabase readiness API. The suite tests whether +pg-sprite preserves established delivery, not Supabase's cold-start guarantees. + +To validate a fixture fix, export the credentials and `SUPABASE_HOSTED_TEST=1` +as shown above, then use the repository's fail-fast, race-enabled procedure: + +```sh +: "${SUPABASE_HOSTED_TEST:?Export SUPABASE_HOSTED_TEST=1 first}" +test "$SUPABASE_HOSTED_TEST" = 1 || exit 1 +scripts/test-flaky.sh TestHostedRealtimeContinuity 10 ./integration/supabase/hosted +``` + +The script stops at the first failure. A successful validation prints +`PASSED all 10 iterations`; it never retries a failed run to obtain a pass. + +A success line is evidence only when the hosted case ran. Without the opt-in, +Go skips it and the flake script can still report success. diff --git a/integration/supabase/hosted/auth_test.go b/integration/supabase/hosted/auth_test.go new file mode 100644 index 0000000..f859a0b --- /dev/null +++ b/integration/supabase/hosted/auth_test.go @@ -0,0 +1,79 @@ +package hosted_test + +import ( + "fmt" + "net/http" + "testing" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// auth.uid() must resolve the real HTTP caller, not the administrator applying DDL. +func TestHostedCLIAuthDefault(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + path := desiredFile(t, fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL DEFAULT auth.uid(), + body text NOT NULL +);`, pgx.Identifier{f.name}.Sanitize())) + f.cli(t, 0, "migrate", "--desired", path, "--json") + f.waitForAPI(t) + status, _, err := f.request(t.Context(), http.MethodPost, "/rest/v1/"+f.name, f.keys.Publishable, f.users[0].Token, map[string]any{"id": 3, "body": "default owner"}) + require.NoError(t, err) + require.Equal(t, http.StatusCreated, status) + var owner string + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT owner_id::text FROM "+f.table+" WHERE id=3").Scan(&owner)) + assert.Equal(t, f.users[0].ID, owner) + status, _, err = f.request(t.Context(), http.MethodPost, "/rest/v1/"+f.name, f.keys.Publishable, f.users[0].Token, map[string]any{"id": 4, "owner_id": f.users[1].ID, "body": "spoofed owner"}) + require.NoError(t, err) + assert.Equal(t, http.StatusForbidden, status) + f.assertRows(t, f.users[0].Token, 1, 3) + f.assertRows(t, f.users[1].Token, 2) + f.assertRows(t, "") + f.assertExportConverges(t) +} + +// Reference users created through Auth; never insert directly into managed auth tables. +func TestHostedCLIAuthForeignKey(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + before := f.snapshot(t) + f.cli(t, 0, "migrate", "--alter", "ALTER TABLE "+f.table+" ADD CONSTRAINT owner_fk FOREIGN KEY (owner_id) REFERENCES auth.users (id)", "--json") + var valid bool + var target string + require.NoError(t, f.pool.QueryRow(t.Context(), `SELECT convalidated, confrelid::regclass::text FROM pg_constraint WHERE conrelid=$1::regclass AND conname='owner_fk'`, f.table).Scan(&valid, &target)) + assert.True(t, valid) + assert.Equal(t, "auth.users", target) + assert.Equal(t, before, f.snapshot(t)) + var orphan string + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT gen_random_uuid()::text").Scan(&orphan)) + var exists bool + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT EXISTS(SELECT 1 FROM auth.users WHERE id=$1)", orphan).Scan(&exists)) + require.False(t, exists) + _, err := f.pool.Exec(t.Context(), "INSERT INTO "+f.table+" VALUES (3,$1,'orphan')", orphan) + var pgErr *pgconn.PgError + require.ErrorAs(t, err, &pgErr) + assert.Equal(t, "23503", pgErr.Code) + f.assertAccess(t) +} + +// A failed concurrent unique build reports its invalid index instead of hiding it. +func TestHostedCLIUniqueIndexFailure(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + f.exec(t, "UPDATE "+f.table+" SET body='duplicate'") + before := f.snapshot(t) + index := pgx.Identifier{f.name + "_unique"}.Sanitize() + result := f.cli(t, 1, "migrate", "--alter", "CREATE UNIQUE INDEX "+index+" ON "+f.table+" (body)", "--json") + assert.Equal(t, "failed", result["outcome"]) + assert.Equal(t, "invalid-index-own-leftover", result["code"]) + var valid bool + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT indisvalid FROM pg_index WHERE indexrelid=$1::regclass", "public."+index).Scan(&valid)) + assert.False(t, valid) + assert.Equal(t, before, f.snapshot(t)) + f.assertAccess(t) +} diff --git a/integration/supabase/hosted/cleanup_test.go b/integration/supabase/hosted/cleanup_test.go new file mode 100644 index 0000000..5f1cbff --- /dev/null +++ b/integration/supabase/hosted/cleanup_test.go @@ -0,0 +1,48 @@ +package hosted_test + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Register before the first write: the CLI can commit DDL and then fail to +// return valid output. Never take cleanup ownership of an existing relation. +func (f *fixture) prepareTableCleanup(t *testing.T) { + t.Helper() + ctx, cancel := context.WithTimeout(t.Context(), 15*time.Second) + defer cancel() + var exists bool + require.NoError(t, f.pool.QueryRow(ctx, "SELECT to_regclass($1) IS NOT NULL", f.table).Scan(&exists)) + require.False(t, exists, "fixture relation must not already exist") + t.Cleanup(func() { + ctx, cancel := context.WithTimeout(context.WithoutCancel(t.Context()), 15*time.Second) + defer cancel() + _, err := f.pool.Exec(ctx, "DROP TABLE IF EXISTS "+f.table) + assert.NoError(t, err) + }) +} + +// Cleanup covers a command that never created its table and one that committed +// the table before a later error. The same hosted opt-in and project checks apply. +func TestHostedTableCleanup(t *testing.T) { + f := newFixture(t) + ctx, cancel := context.WithTimeout(t.Context(), 15*time.Second) + defer cancel() + t.Run("no committed table", func(t *testing.T) { + f.prepareTableCleanup(t) + }) + t.Run("table committed before failure", func(t *testing.T) { + f.prepareTableCleanup(t) + _, err := f.pool.Exec(ctx, "CREATE TABLE "+f.table+" (id integer PRIMARY KEY)") + require.NoError(t, err) + _, err = f.pool.Exec(ctx, "SELECT 1 / 0") + require.Error(t, err) + }) + var exists bool + require.NoError(t, f.pool.QueryRow(ctx, "SELECT to_regclass($1) IS NOT NULL", f.table).Scan(&exists)) + assert.False(t, exists, "cleanup removes the committed table despite the later failure") +} diff --git a/integration/supabase/hosted/ddl_test.go b/integration/supabase/hosted/ddl_test.go new file mode 100644 index 0000000..5463150 --- /dev/null +++ b/integration/supabase/hosted/ddl_test.go @@ -0,0 +1,83 @@ +package hosted_test + +import ( + "fmt" + "testing" + + "github.com/block/pg-sprite/pkg/diffplan" + "github.com/block/pg-sprite/pkg/migrate" + "github.com/block/pg-sprite/pkg/router" + "github.com/block/pg-sprite/pkg/statement" + "github.com/block/pg-sprite/pkg/verdict" + "github.com/jackc/pgx/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func (f *fixture) apply(t *testing.T, sql string) verdict.Verdict { + t.Helper() + st, err := statement.ParseOne(sql) + require.NoError(t, err) + result, err := migrate.Run(t.Context(), f.pool, st, migrate.DefaultOptions()) + require.NoError(t, err) + require.Equal(t, verdict.OutcomeExecuted, result.Outcome) + return result +} + +// Volatile defaults require copy-and-swap; refusal must preserve the live table. +func realtimeRefuseUUID(t *testing.T, f *fixture, first, second *stream) { + f.refuse(t, first, second, "ALTER TABLE "+f.table+" ADD COLUMN token uuid DEFAULT gen_random_uuid()", fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL, + safe_prefix text, + token uuid DEFAULT gen_random_uuid() + )`, pgx.Identifier{f.name}.Sanitize())) +} +func realtimeRefuseTimestamp(t *testing.T, f *fixture, first, second *stream) { + f.refuse(t, first, second, "ALTER TABLE "+f.table+" ADD COLUMN created_at timestamptz DEFAULT clock_timestamp()", fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL, + safe_prefix text, + created_at timestamptz DEFAULT clock_timestamp() + )`, pgx.Identifier{f.name}.Sanitize())) +} + +func (f *fixture) snapshot(t *testing.T) string { + t.Helper() + var result string + require.NoError(t, f.pool.QueryRow(t.Context(), `SELECT jsonb_build_object( + 'columns',(SELECT jsonb_agg(to_jsonb(a) ORDER BY attnum) FROM pg_attribute a WHERE attrelid=$1::regclass AND attnum>0 AND NOT attisdropped), + 'policies',(SELECT jsonb_agg(to_jsonb(p) ORDER BY polname) FROM pg_policy p WHERE polrelid=$1::regclass), + 'table',(SELECT jsonb_build_array(oid,relfilenode,relrowsecurity,relacl) FROM pg_class WHERE oid=$1::regclass), + 'rows',(SELECT jsonb_agg(to_jsonb(r) ORDER BY id) FROM `+f.table+` r) + )::text`, f.table).Scan(&result)) + return result +} +func (f *fixture) refuse(t *testing.T, first, second *stream, sql, desired string) { + t.Helper() + before := f.snapshot(t) + st, err := statement.ParseOne(sql) + require.NoError(t, err) + result, err := migrate.Run(t.Context(), f.pool, st, migrate.DefaultOptions()) + require.NoError(t, err) + assert.Equal(t, verdict.OutcomeRefused, result.Outcome) + assert.Equal(t, verdict.ReasonBackendUnavailable, result.Reason) + assert.Empty(t, result.ExecutedSQL) + ds, err := statement.ParseDesired(desired) + require.NoError(t, err) + plan, err := diffplan.Plan(t.Context(), f.pool, diffplan.Request{Schema: "public", Desired: ds}) + require.NoError(t, err) + require.Equal(t, router.DispositionUnavailable, plan.Disposition) + report, err := migrate.RunDesired(t.Context(), f.pool, migrate.DesiredRequest{Schema: "public", Desired: ds}, migrate.DefaultOptions()) + require.NoError(t, err) + assert.Equal(t, verdict.OutcomeRefused, report.Outcome) + assert.Equal(t, verdict.ReasonBackendUnavailable, report.Reason) + assert.Empty(t, report.Verdicts) + assert.Equal(t, before, f.snapshot(t)) + f.assertAccess(t) + f.exec(t, "UPDATE "+f.table+" SET body=$1", "after "+t.Name()) + first.exactRow(t, "UPDATE", 1) + second.exactRow(t, "UPDATE", 2) +} diff --git a/integration/supabase/hosted/fixture_test.go b/integration/supabase/hosted/fixture_test.go new file mode 100644 index 0000000..1ac5a2b --- /dev/null +++ b/integration/supabase/hosted/fixture_test.go @@ -0,0 +1,216 @@ +package hosted_test + +import ( + "bytes" + "context" + "crypto/rand" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "os" + "strings" + "testing" + "time" + + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type credentials struct { + Publishable string `json:"publishable_key"` + Secret string `json:"secret_key"` +} +type user struct { + ID string `json:"id"` + Token string `json:"access_token"` +} +type fixture struct { + pool *pgxpool.Pool + base string + keys credentials + users []user + client *http.Client + name string + table string +} + +func newFixture(t *testing.T) *fixture { + t.Helper() + if os.Getenv("SUPABASE_HOSTED_TEST") != "1" { + t.Skip("set SUPABASE_HOSTED_TEST=1 for a disposable hosted project") + } + base := os.Getenv("SUPABASE_HOSTED_URL") + parsed, err := url.Parse(base) + require.NoError(t, err) + require.Equal(t, "https", parsed.Scheme) + require.Nil(t, parsed.User) + require.Empty(t, parsed.RawQuery) + require.Empty(t, parsed.Fragment) + require.Empty(t, parsed.Path) + ref := strings.TrimSuffix(parsed.Host, ".supabase.co") + require.NotEmpty(t, ref) + require.NotContains(t, ref, ".") + require.Equal(t, ref+".supabase.co", parsed.Host) + dsn := os.Getenv("PGSPRITE_URL") + cfg, err := pgx.ParseConfig(dsn) + require.NoError(t, err) + require.Equal(t, "db."+ref+".supabase.co", cfg.Host, "use this project's direct database URL for the baseline") + path := os.Getenv("SUPABASE_HOSTED_CREDENTIALS") + info, err := os.Stat(path) + require.NoError(t, err) + require.Zero(t, info.Mode().Perm()&0077, "credentials file must be private") + data, err := os.ReadFile(path) + require.NoError(t, err) + var keys credentials + require.NoError(t, json.Unmarshal(data, &keys)) + require.NotEmpty(t, keys.Publishable) + require.NotEmpty(t, keys.Secret) + pool, err := dbconn.NewPool(t.Context(), dbconn.Config{URL: dsn}) + require.NoError(t, err) + t.Cleanup(pool.Close) + f := &fixture{pool: pool, base: base, keys: keys, client: &http.Client{Timeout: 15 * time.Second}} + for range 2 { + f.users = append(f.users, f.createUser(t)) + } + f.name = "pgsprite_hosted_" + randomID(t) + f.table = pgx.Identifier{"public", f.name}.Sanitize() + return f +} + +func randomID(t *testing.T) string { + t.Helper() + b := make([]byte, 8) + _, err := rand.Read(b) + require.NoError(t, err) + return hex.EncodeToString(b) +} + +// Administrative requests never print their response bodies: those can contain tokens. +func (f *fixture) request(ctx context.Context, method, path, key, token string, payload any) (int, []byte, error) { + body, err := json.Marshal(payload) + if err != nil { + return 0, nil, err + } + req, err := http.NewRequestWithContext(ctx, method, f.base+path, bytes.NewReader(body)) + if err != nil { + return 0, nil, err + } + req.Header.Set("apikey", key) + req.Header.Set("Content-Type", "application/json") + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + req.Header.Set("Prefer", "return=representation") + response, err := f.client.Do(req) + if err != nil { + return 0, nil, err + } + data, readErr := io.ReadAll(io.LimitReader(response.Body, 1<<20)) + closeErr := response.Body.Close() + if readErr != nil { + return 0, nil, readErr + } + if closeErr != nil { + return 0, nil, closeErr + } + return response.StatusCode, data, nil +} + +func (f *fixture) createUser(t *testing.T) user { + t.Helper() + email := "pgsprite-" + randomID(t) + "@example.com" + password := randomID(t) + randomID(t) + status, data, err := f.request(t.Context(), http.MethodPost, "/auth/v1/admin/users", f.keys.Secret, "", map[string]any{"email": email, "password": password, "email_confirm": true}) + require.NoError(t, err) + require.Equal(t, http.StatusOK, status) + var u user + require.NoError(t, json.Unmarshal(data, &u)) + require.NotEmpty(t, u.ID) + t.Cleanup(func() { + ctx, cancel := context.WithTimeout(context.WithoutCancel(t.Context()), 15*time.Second) + defer cancel() + status, _, err := f.request(ctx, http.MethodDelete, "/auth/v1/admin/users/"+url.PathEscape(u.ID), f.keys.Secret, "", nil) + assert.NoError(t, err) + assert.Equal(t, http.StatusOK, status) + }) + status, data, err = f.request(t.Context(), http.MethodPost, "/auth/v1/token?grant_type=password", f.keys.Publishable, "", map[string]string{"email": email, "password": password}) + require.NoError(t, err) + require.Equal(t, http.StatusOK, status) + var session user + require.NoError(t, json.Unmarshal(data, &session)) + require.NotEmpty(t, session.Token) + u.Token = session.Token + return u +} + +func (f *fixture) exec(t *testing.T, sql string, args ...any) { + t.Helper() + _, err := f.pool.Exec(t.Context(), sql, args...) + require.NoError(t, err) +} +func (f *fixture) seed(t *testing.T) { + t.Helper() + f.seedTable(t) + f.exec(t, "ALTER PUBLICATION supabase_realtime ADD TABLE "+f.table) + f.waitForAPI(t) + f.assertAccess(t) +} + +func (f *fixture) seedUnpublished(t *testing.T) { + t.Helper() + f.seedTable(t) + f.waitForAPI(t) + f.assertAccess(t) +} + +func (f *fixture) seedTable(t *testing.T) { + t.Helper() + f.prepareTableCleanup(t) + f.exec(t, fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL + )`, f.table)) + f.exec(t, "ALTER TABLE "+f.table+" ENABLE ROW LEVEL SECURITY") + f.exec(t, "CREATE POLICY own_rows ON "+f.table+" TO authenticated USING (owner_id = auth.uid()) WITH CHECK (owner_id = auth.uid())") + f.exec(t, "GRANT SELECT, INSERT, UPDATE, DELETE ON "+f.table+" TO authenticated, anon") + f.exec(t, "INSERT INTO "+f.table+" VALUES (1,$1,'first'),(2,$2,'second')", f.users[0].ID, f.users[1].ID) +} + +func (f *fixture) waitForAPI(t *testing.T) { + t.Helper() + f.exec(t, "NOTIFY pgrst, 'reload schema'") + const cacheDeadline = 30 * time.Second + require.Eventually(t, func() bool { + status, _, err := f.request(t.Context(), http.MethodGet, "/rest/v1/"+f.name+"?select=id", f.keys.Publishable, f.users[0].Token, nil) + return err == nil && status == http.StatusOK + }, cacheDeadline, 200*time.Millisecond, "Data API must discover the fixture") +} + +func (f *fixture) assertRows(t *testing.T, token string, want ...int) { + t.Helper() + status, data, err := f.request(t.Context(), http.MethodGet, "/rest/v1/"+f.name+"?select=id&order=id", f.keys.Publishable, token, nil) + require.NoError(t, err) + require.Equal(t, http.StatusOK, status) + var rows []struct { + ID int `json:"id"` + } + require.NoError(t, json.Unmarshal(data, &rows)) + ids := make([]int, 0, len(rows)) + for _, r := range rows { + ids = append(ids, r.ID) + } + assert.Equal(t, append([]int{}, want...), ids) +} +func (f *fixture) assertAccess(t *testing.T) { + t.Helper() + f.assertRows(t, f.users[0].Token, 1) + f.assertRows(t, f.users[1].Token, 2) + f.assertRows(t, "") +} diff --git a/integration/supabase/hosted/lifecycle_failure_test.go b/integration/supabase/hosted/lifecycle_failure_test.go new file mode 100644 index 0000000..512d4ef --- /dev/null +++ b/integration/supabase/hosted/lifecycle_failure_test.go @@ -0,0 +1,96 @@ +package hosted_test + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Failed NOT NULL validation retains NULL rows and reports its committed scaffold. +func TestHostedCLIFailedNotNull(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + f.exec(t, "ALTER TABLE "+f.table+" ADD COLUMN label text") + before := f.snapshot(t) + result := f.cli(t, 1, "migrate", "--alter", "ALTER TABLE "+f.table+" ALTER COLUMN label SET NOT NULL", "--json") + assert.Equal(t, "failed", result["outcome"]) + assert.Equal(t, "execution-failed", result["code"]) + require.Len(t, result["executed_sql"], 1, "the NOT VALID scaffold commits before validation fails") + var nulls, scaffolds int + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT count(*) FROM "+f.table+" WHERE label IS NULL").Scan(&nulls)) + require.NoError(t, f.pool.QueryRow(t.Context(), `SELECT count(*) FROM pg_constraint + WHERE conrelid=$1::regclass AND contype='c' AND NOT convalidated`, f.table).Scan(&scaffolds)) + assert.Equal(t, 2, nulls) + assert.Equal(t, 1, scaffolds) + assert.Equal(t, before, f.snapshot(t)) + f.assertAccess(t) +} + +// A destructive desired schema is rejected before dropping data or policies. +func TestHostedCLIRefuseDropColumn(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + before := f.snapshot(t) + path := desiredFile(t, fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL +);`, pgx.Identifier{f.name}.Sanitize())) + preview := f.cli(t, 2, "migrate", "--desired", path, "--dry-run", "--json") + assert.Equal(t, "destructive-change", preview["reason"]) + result := f.cli(t, 2, "migrate", "--desired", path, "--json") + assert.Equal(t, "destructive-change", result["reason"]) + assert.Equal(t, before, f.snapshot(t)) + f.assertAccess(t) +} + +// A refused copy-and-swap plan cannot partially apply its safe column addition. +func TestHostedCLIRefuseCopySwap(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + before := f.snapshot(t) + path := desiredFile(t, fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL, + safe_prefix text, + token uuid DEFAULT gen_random_uuid() +);`, pgx.Identifier{f.name}.Sanitize())) + result := f.cli(t, 2, "migrate", "--desired", path, "--json") + assert.Equal(t, "backend-unavailable", result["reason"]) + assert.Empty(t, result["verdicts"]) + assert.Equal(t, before, f.snapshot(t)) + f.assertAccess(t) +} + +// A held lock prevents removing the last policy; access must remain unchanged. +func TestHostedCLIRLSLockFailure(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + before := f.snapshot(t) + name := pgx.Identifier{f.name}.Sanitize() + path := desiredFile(t, fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL +); +ALTER TABLE %s ENABLE ROW LEVEL SECURITY;`, name, name)) + holder, err := f.pool.Begin(t.Context()) + require.NoError(t, err) + t.Cleanup(func() { + ctx, cancel := context.WithTimeout(context.WithoutCancel(t.Context()), 5*time.Second) + defer cancel() + _ = holder.Rollback(ctx) + }) + _, err = holder.Exec(t.Context(), "LOCK TABLE "+f.table+" IN ACCESS SHARE MODE") + require.NoError(t, err) + result := f.cli(t, 1, "migrate", "--desired", path, "--lock-timeout", "200ms", "--json") + assert.Equal(t, "budget-lock-exceeded", result["code"]) + require.NoError(t, holder.Rollback(t.Context())) + assert.Equal(t, before, f.snapshot(t)) + f.assertAccess(t) +} diff --git a/integration/supabase/hosted/lifecycle_test.go b/integration/supabase/hosted/lifecycle_test.go new file mode 100644 index 0000000..a3f630b --- /dev/null +++ b/integration/supabase/hosted/lifecycle_test.go @@ -0,0 +1,135 @@ +package hosted_test + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "testing" + "time" + + "github.com/block/pg-sprite/pkg/verdict" + "github.com/jackc/pgx/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Run the shipped CLI, not a second implementation of desired-state execution. +func (f *fixture) cli(t *testing.T, wantExit int, args ...string) map[string]any { + t.Helper() + binary := os.Getenv("SUPABASE_HOSTED_BIN") + require.NotEmpty(t, binary, "set SUPABASE_HOSTED_BIN to a freshly built pg-sprite binary") + ctx, cancel := context.WithTimeout(t.Context(), time.Minute) + defer cancel() + cmd := exec.CommandContext(ctx, binary, args...) + var diagnostic bytes.Buffer + cmd.Stderr = &diagnostic + output, err := cmd.Output() + code := 0 + if err != nil { + var exit *exec.ExitError + require.True(t, errors.As(err, &exit), "CLI could not execute") + code = exit.ExitCode() + } + require.Equal(t, wantExit, code, "CLI stdout: %s; stderr: %s", output, diagnostic.String()) + var result map[string]any + if len(args) > 0 && args[len(args)-1] == "--json" { + require.NoError(t, json.Unmarshal(output, &result)) + } + return result +} + +func desiredFile(t *testing.T, sql string) string { + t.Helper() + path := filepath.Join(t.TempDir(), "desired.sql") + require.NoError(t, os.WriteFile(path, []byte(sql), 0600)) + return path +} + +func (f *fixture) assertExportConverges(t *testing.T) { + t.Helper() + dir := filepath.Join(t.TempDir(), "export") + f.cli(t, 0, "pull", "--schema", "public", "--out", dir) + result := f.cli(t, 0, "diff", "--desired", filepath.Join(dir, f.name+".sql"), "--json") + require.Contains(t, result, "statements") + assert.Empty(t, result["statements"]) +} + +// This complete CLI lifecycle uses project defaults, then manages RLS explicitly. +// It does not subscribe to Realtime or change the publication. +func TestHostedCLICreateExportAndRLS(t *testing.T) { + f := newFixture(t) + name := pgx.Identifier{f.name}.Sanitize() + base := fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL +);`, name) + path := desiredFile(t, base) + f.cli(t, 0, "migrate", "--desired", path, "--dry-run", "--json") + f.prepareTableCleanup(t) + f.cli(t, 0, "migrate", "--desired", path, "--json") + f.exec(t, "INSERT INTO "+f.table+" VALUES (1,$1,'first'),(2,$2,'second')", f.users[0].ID, f.users[1].ID) + f.assertExportConverges(t) + + rls := base + fmt.Sprintf(` +ALTER TABLE %s ENABLE ROW LEVEL SECURITY; +CREATE POLICY readers ON %s FOR SELECT TO authenticated + USING (owner_id = auth.uid());`, name, name) + path = desiredFile(t, rls) + preview := f.cli(t, 2, "migrate", "--desired", path, "--dry-run", "--json") + assert.Equal(t, "unsupported-statement", preview["reason"]) + assert.NotEmpty(t, preview["row_security_review"], "preview must describe the RLS changes for review") + f.cli(t, 0, "migrate", "--desired", path, "--json") + f.waitForAPI(t) + f.assertAccess(t) + before := f.snapshot(t) + result := f.cli(t, 0, "migrate", "--desired", path, "--json") + assert.Equal(t, string(verdict.OutcomeExecuted), result["outcome"]) + assert.Empty(t, result["executed_sql"]) + assert.Equal(t, before, f.snapshot(t)) + f.assertExportConverges(t) + f.assertAccess(t) +} + +func TestHostedCLIAddColumn(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + path := desiredFile(t, fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL, + archived boolean NOT NULL DEFAULT false +);`, pgx.Identifier{f.name}.Sanitize())) + f.cli(t, 0, "migrate", "--desired", path, "--json") + var rows int + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT count(*) FROM "+f.table+" WHERE NOT archived AND ((id=1 AND body='first') OR (id=2 AND body='second'))").Scan(&rows)) + assert.Equal(t, 2, rows) + f.assertAccess(t) + f.assertExportConverges(t) +} + +func TestHostedCLIAddIndex(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + before := f.snapshot(t) + name := pgx.Identifier{f.name}.Sanitize() + index := pgx.Identifier{f.name + "_body"}.Sanitize() + path := desiredFile(t, fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL +); +CREATE INDEX %s ON %s (body);`, name, index, name)) + f.cli(t, 0, "migrate", "--desired", path, "--json") + var valid bool + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT indisvalid FROM pg_index WHERE indexrelid=$1::regclass", "public."+index).Scan(&valid)) + assert.True(t, valid) + assert.Equal(t, before, f.snapshot(t), "index addition preserves rows, columns, table identity, policies, and grants") + f.assertAccess(t) + f.assertExportConverges(t) +} diff --git a/integration/supabase/hosted/pooler_test.go b/integration/supabase/hosted/pooler_test.go new file mode 100644 index 0000000..0fff015 --- /dev/null +++ b/integration/supabase/hosted/pooler_test.go @@ -0,0 +1,40 @@ +package hosted_test + +import ( + "os" + "strings" + "testing" + + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/jackc/pgx/v5" + "github.com/stretchr/testify/require" +) + +// Pooler URLs come from Connect, not an inferred region or an assumed port. +func TestHostedSessionPooler(t *testing.T) { + dsn := os.Getenv("SUPABASE_HOSTED_SESSION_URL") + if dsn == "" { + t.Skip("set SUPABASE_HOSTED_SESSION_URL to test the session pooler") + } + if os.Getenv("SUPABASE_HOSTED_TEST") != "1" { + t.Skip("set SUPABASE_HOSTED_TEST=1 for a disposable hosted project") + } + checkPoolerProject(t, os.Getenv("SUPABASE_HOSTED_URL"), dsn) + f := newFixture(t) + f.seedUnpublished(t) + pool, err := dbconn.NewPool(t.Context(), dbconn.Config{URL: dsn}) + require.NoError(t, err) + t.Cleanup(pool.Close) + pooled := *f + pooled.pool = pool + pooled.apply(t, "ALTER TABLE "+f.table+" ADD COLUMN session_note text") + f.assertAccess(t) +} +func checkPoolerProject(t *testing.T, base, dsn string) { + t.Helper() + cfg, err := pgx.ParseConfig(dsn) + require.NoError(t, err) + ref := strings.TrimSuffix(strings.TrimPrefix(base, "https://"), ".supabase.co") + require.Equal(t, "postgres."+ref, cfg.User, "pooler must address the same disposable project") + require.True(t, strings.HasSuffix(cfg.Host, ".pooler.supabase.com"), "use the dashboard's Supavisor endpoint") +} diff --git a/integration/supabase/hosted/realtime_lifecycle_test.go b/integration/supabase/hosted/realtime_lifecycle_test.go new file mode 100644 index 0000000..56c0b37 --- /dev/null +++ b/integration/supabase/hosted/realtime_lifecycle_test.go @@ -0,0 +1,91 @@ +package hosted_test + +import ( + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// One fixture and two sockets cover the entire steady-state schema change lifecycle. +// Initialization is observable and has a separate deadline from event delivery. +func TestHostedRealtimeContinuity(t *testing.T) { + f := newFixture(t) + f.seed(t) + first, second := f.subscribe(t, 0), f.subscribe(t, 1) + if !t.Run("initialize", func(t *testing.T) { + started := time.Now() + // Wait for an active publication reader before sending the baseline write. + const publicationStartupDeadline = 75 * time.Second + require.EventuallyWithT(t, func(c *assert.CollectT) { + var ready bool + err := f.pool.QueryRow(t.Context(), `SELECT EXISTS ( + SELECT 1 FROM pg_replication_slots WHERE plugin='wal2json' AND active + )`).Scan(&ready) + require.NoError(c, err) + assert.True(c, ready, "waiting for hosted publication reader") + }, publicationStartupDeadline, 200*time.Millisecond) + for _, u := range f.users { + var count int + require.NoError(t, f.pool.QueryRow(t.Context(), `SELECT count(*) FROM realtime.subscription + WHERE entity=$1::regclass AND claims->>'sub'=$2 AND claims->>'role'='authenticated'`, f.table, u.ID).Scan(&count)) + require.Equal(t, 1, count) + } + f.exec(t, "UPDATE "+f.table+" SET body='baseline ready'") + first.exactRow(t, "UPDATE", 1) + second.exactRow(t, "UPDATE", 2) + t.Logf("reader and baseline delivery ready after %s", time.Since(started)) + }) { + return + } + update := func(t *testing.T) { + t.Helper() + f.exec(t, "UPDATE "+f.table+" SET body=$1", t.Name()) + first.exactRow(t, "UPDATE", 1) + second.exactRow(t, "UPDATE", 2) + f.assertAccess(t) + } + if !t.Run("no_DDL_control", update) { + return + } + if !t.Run("refuse_volatile_UUID", func(t *testing.T) { realtimeRefuseUUID(t, f, first, second) }) { + return + } + if !t.Run("refuse_volatile_timestamp", func(t *testing.T) { realtimeRefuseTimestamp(t, f, first, second) }) { + return + } + var oid uint32 + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT $1::regclass::oid", f.table).Scan(&oid)) + if !t.Run("add_column", func(t *testing.T) { + f.apply(t, "ALTER TABLE "+f.table+" ADD COLUMN title text") + update(t) + }) { + return + } + if !t.Run("add_index", func(t *testing.T) { + result := f.apply(t, "CREATE INDEX "+pgx.Identifier{f.name + "_body"}.Sanitize()+" ON "+f.table+" (body)") + require.Len(t, result.ExecutedSQL, 1) + assert.Regexp(t, `^CREATE INDEX CONCURRENTLY `, result.ExecutedSQL[0]) + update(t) + }) { + return + } + var current uint32 + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT $1::regclass::oid", f.table).Scan(¤t)) + assert.Equal(t, oid, current) +} + +// Unlike a search for one matching row, reject an unexpected row or operation. +func (s *stream) exactRow(t *testing.T, operation string, id int) { + t.Helper() + s.await(t, func(e event) bool { + if e.Event != "postgres_changes" { + return false + } + require.Equal(t, operation, e.Payload.Data.Type) + require.Equal(t, id, e.Payload.Data.Record.ID) + return true + }) +} diff --git a/integration/supabase/hosted/realtime_test.go b/integration/supabase/hosted/realtime_test.go new file mode 100644 index 0000000..3db0645 --- /dev/null +++ b/integration/supabase/hosted/realtime_test.go @@ -0,0 +1,113 @@ +package hosted_test + +import ( + "context" + "fmt" + "net/url" + "testing" + "time" + + "github.com/gorilla/websocket" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type event struct { + Event string `json:"event"` + Payload struct { + Status string `json:"status"` + Extension string `json:"extension"` + Message string `json:"message"` + Data struct { + Type string `json:"type"` + Record struct { + ID int `json:"id"` + Owner string `json:"owner_id"` + } `json:"record"` + } `json:"data"` + } `json:"payload"` +} +type stream struct { + conn *websocket.Conn + owner string +} + +func (f *fixture) subscribe(t *testing.T, tenant int) *stream { + t.Helper() + endpoint, err := url.Parse(f.base) + require.NoError(t, err) + endpoint.Scheme = "wss" + endpoint.Path = "/realtime/v1/websocket" + endpoint.RawQuery = url.Values{"apikey": {f.keys.Publishable}, "vsn": {"1.0.0"}}.Encode() + conn, response, err := websocket.DefaultDialer.DialContext(t.Context(), endpoint.String(), nil) + if response != nil && response.Body != nil { + require.NoError(t, response.Body.Close()) + } + // Do not include the URL in failure diagnostics; it contains the API key. + require.True(t, err == nil, "Realtime WebSocket connection failed") + t.Cleanup(func() { assert.NoError(t, conn.Close()) }) + conn.SetReadLimit(1 << 20) + require.NoError(t, conn.SetWriteDeadline(time.Now().Add(5*time.Second))) + require.NoError(t, conn.WriteJSON(map[string]any{ + "topic": "realtime:" + f.name, "event": "phx_join", "ref": "1", "join_ref": "1", + "payload": map[string]any{"access_token": f.users[tenant].Token, "config": map[string]any{"postgres_changes": []map[string]any{{"event": "*", "schema": "public", "table": f.name}}}}, + })) + s := &stream{conn: conn, owner: f.users[tenant].ID} + s.startHeartbeat(t) + s.await(t, func(e event) bool { + return e.Event == "system" && e.Payload.Extension == "postgres_changes" && e.Payload.Status == "ok" + }) + return s +} + +func (s *stream) await(t *testing.T, match func(event) bool) { + t.Helper() + const eventDeadline = 30 * time.Second + require.NoError(t, s.conn.SetReadDeadline(time.Now().Add(eventDeadline))) + for { + var e event + err := s.conn.ReadJSON(&e) + require.NoError(t, err, "waiting for hosted Realtime event") + if e.Event == "system" { + t.Logf("Realtime system status=%s message=%s", e.Payload.Status, e.Payload.Message) + } + require.NotEqual(t, "phx_error", e.Event) + require.NotEqual(t, "phx_close", e.Event) + require.NotEqual(t, "error", e.Payload.Status, "Realtime system error: %s", e.Payload.Message) + if e.Event == "postgres_changes" { + require.Equal(t, s.owner, e.Payload.Data.Record.Owner, "cross-tenant event") + } + if match(e) { + return + } + } +} + +// Keep long fixture initialization alive using the protocol heartbeat, not reconnects. +// The writer is stopped and joined before the earlier socket-close cleanup runs. +func (s *stream) startHeartbeat(t *testing.T) { + t.Helper() + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan error, 1) + go func() { + ticker := time.NewTicker(15 * time.Second) + defer ticker.Stop() + for ref := 2; ; ref++ { + select { + case <-ctx.Done(): + done <- nil + return + case <-ticker.C: + if err := s.conn.SetWriteDeadline(time.Now().Add(5 * time.Second)); err != nil { + done <- err + return + } + if err := s.conn.WriteJSON(map[string]any{"topic": "phoenix", "event": "heartbeat", "payload": map[string]any{}, "ref": fmt.Sprint(ref)}); err != nil { + done <- err + return + } + } + } + }() + t.Cleanup(func() { cancel(); assert.NoError(t, <-done, "Realtime heartbeat failed") }) +} diff --git a/integration/supabase/hosted/rls_rollback_test.go b/integration/supabase/hosted/rls_rollback_test.go new file mode 100644 index 0000000..81cea9d --- /dev/null +++ b/integration/supabase/hosted/rls_rollback_test.go @@ -0,0 +1,88 @@ +package hosted_test + +import ( + "context" + "fmt" + "net/http" + "sync/atomic" + "testing" + "time" + + "github.com/block/pg-sprite/pkg/executor" + "github.com/block/pg-sprite/pkg/schemadiff" + "github.com/block/pg-sprite/pkg/statement" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Cancel only after the server has successfully dropped a live policy. This +// exercises partial work inside a real transaction, not a pre-execution refusal. +type cancelAfterRLSDrop struct { + cancel context.CancelFunc + dropSQL string + reached atomic.Bool +} + +type rlsDropQueryKey struct{} + +func (c *cancelAfterRLSDrop) TraceQueryStart(ctx context.Context, _ *pgx.Conn, data pgx.TraceQueryStartData) context.Context { + return context.WithValue(ctx, rlsDropQueryKey{}, data.SQL == c.dropSQL) +} +func (c *cancelAfterRLSDrop) TraceQueryEnd(ctx context.Context, _ *pgx.Conn, data pgx.TraceQueryEndData) { + // Match the fully qualified live statement, not a DROP in a scratch schema. + liveDrop, _ := ctx.Value(rlsDropQueryKey{}).(bool) + if liveDrop && data.Err == nil && data.CommandTag.String() == "DROP POLICY" { + c.reached.Store(true) + c.cancel() + } +} + +func TestHostedRLSCancellationPreservesAccess(t *testing.T) { + f := newFixture(t) + f.seedUnpublished(t) + pool, name := f.pool, f.name + before, err := schemadiff.Introspect(t.Context(), pool, "public", name) + require.NoError(t, err) + f.assertAccess(t) + desired, err := statement.ParseDesiredWithRowSecurity(fmt.Sprintf(`CREATE TABLE %s ( + id int PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL + ); + ALTER TABLE %s ENABLE ROW LEVEL SECURITY; + CREATE POLICY nobody ON %s + FOR SELECT TO authenticated USING (false);`, pgx.Identifier{name}.Sanitize(), pgx.Identifier{name}.Sanitize(), pgx.Identifier{name}.Sanitize())) + require.NoError(t, err) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + trace := &cancelAfterRLSDrop{ + cancel: cancel, + dropSQL: `DROP POLICY "own_rows" ON ` + f.table, + } + cfg := pool.Config() + cfg.ConnConfig.Tracer = trace + changing, err := pgxpool.NewWithConfig(t.Context(), cfg) + require.NoError(t, err) + defer changing.Close() + _, err = executor.ExecuteRowSecurity(ctx, changing, "public", desired, executor.Budget{LockTimeout: 100 * time.Millisecond, StatementTimeout: 5 * time.Second}) + require.True(t, trace.reached.Load(), "cancellation must happen after live DROP POLICY succeeds") + require.ErrorIs(t, err, context.Canceled) + after, err := schemadiff.Introspect(t.Context(), pool, "public", name) + require.NoError(t, err) + assert.Equal(t, before, after) + f.assertAccess(t) + status, _, err := f.request(t.Context(), http.MethodPost, "/rest/v1/"+name, f.keys.Publishable, f.users[0].Token, map[string]any{"id": 90, "owner_id": f.users[1].ID, "body": "cross tenant"}) + require.NoError(t, err) + assert.Equal(t, http.StatusForbidden, status) + status, _, err = f.request(t.Context(), http.MethodPost, "/rest/v1/"+name, f.keys.Publishable, "", map[string]any{"id": 91, "owner_id": f.users[0].ID, "body": "anonymous"}) + require.NoError(t, err) + assert.Equal(t, http.StatusUnauthorized, status) + status, _, err = f.request(t.Context(), http.MethodPost, "/rest/v1/"+name, f.keys.Publishable, f.users[0].Token, map[string]any{"id": 3, "owner_id": f.users[0].ID, "body": "after rollback"}) + require.NoError(t, err) + require.Equal(t, http.StatusCreated, status) + f.assertRows(t, f.users[0].Token, 1, 3) + f.assertRows(t, f.users[1].Token, 2) + f.assertRows(t, "") +} diff --git a/integration/supabase/hosted/rls_test.go b/integration/supabase/hosted/rls_test.go new file mode 100644 index 0000000..263282b --- /dev/null +++ b/integration/supabase/hosted/rls_test.go @@ -0,0 +1,52 @@ +package hosted_test + +import ( + "fmt" + "net/http" + "testing" + "time" + + "github.com/block/pg-sprite/pkg/executor" + "github.com/block/pg-sprite/pkg/statement" + "github.com/jackc/pgx/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Policies change real signed-in users' access; removing the last policy denies all. +func TestHostedDeclarativeRLS(t *testing.T) { + f := newFixture(t) + name := pgx.Identifier{f.name}.Sanitize() + f.seedUnpublished(t) + desired := fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL + ); + ALTER TABLE %s ENABLE ROW LEVEL SECURITY; + CREATE POLICY readers ON %s FOR SELECT TO authenticated + USING (owner_id = auth.uid());`, name, name, name) + definitions, err := statement.ParseDesiredWithRowSecurity(desired) + require.NoError(t, err) + _, err = executor.ExecuteRowSecurity(t.Context(), f.pool, "public", definitions, executor.Budget{LockTimeout: time.Second, StatementTimeout: 5 * time.Second}) + require.NoError(t, err) + f.assertAccess(t) + status, _, err := f.request(t.Context(), http.MethodPost, "/rest/v1/"+f.name, f.keys.Publishable, f.users[0].Token, map[string]any{"id": 3, "owner_id": f.users[0].ID, "body": "write denied"}) + require.NoError(t, err) + assert.Equal(t, http.StatusForbidden, status) + deny, err := statement.ParseDesiredWithRowSecurity(fmt.Sprintf(`CREATE TABLE %s ( + id integer PRIMARY KEY, + owner_id uuid NOT NULL, + body text NOT NULL + ); + ALTER TABLE %s ENABLE ROW LEVEL SECURITY;`, name, name)) + require.NoError(t, err) + _, err = executor.ExecuteRowSecurity(t.Context(), f.pool, "public", deny, executor.Budget{LockTimeout: time.Second, StatementTimeout: 5 * time.Second}) + require.NoError(t, err) + f.assertRows(t, f.users[0].Token) + f.assertRows(t, f.users[1].Token) + f.assertRows(t, "") + var count int + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT count(*) FROM "+f.table).Scan(&count)) + assert.Equal(t, 2, count) +} diff --git a/integration/supabase/hosted/tls_test.go b/integration/supabase/hosted/tls_test.go new file mode 100644 index 0000000..2a24e7b --- /dev/null +++ b/integration/supabase/hosted/tls_test.go @@ -0,0 +1,114 @@ +package hosted_test + +import ( + "context" + "crypto/rand" + "crypto/rsa" + "crypto/x509" + "crypto/x509/pkix" + "encoding/pem" + "math/big" + "net/url" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// TLS checks need only a database credential; they create no tables or Auth users. +func hostedTLSURL(t *testing.T, mode string, keepCA bool) string { + t.Helper() + if os.Getenv("SUPABASE_HOSTED_TEST") != "1" { + t.Skip("set SUPABASE_HOSTED_TEST=1 for hosted checks") + } + u, err := url.Parse(os.Getenv("PGSPRITE_URL")) + require.NoError(t, err) + require.Contains(t, []string{"postgres", "postgresql"}, u.Scheme) + require.True(t, strings.HasPrefix(u.Hostname(), "db.") && strings.HasSuffix(u.Hostname(), ".supabase.co"), "use a hosted direct endpoint") + q := u.Query() + if mode == "" { + q.Del("sslmode") + } else { + q.Set("sslmode", mode) + } + if !keepCA { + q.Del("sslrootcert") + } else if q.Get("sslrootcert") == "" { + t.Skip("set sslrootcert to the downloaded Supabase CA for verification tests") + } + u.RawQuery = q.Encode() + return u.String() +} +func assertHostedTLS(t *testing.T, dsn string) { + t.Helper() + ctx, cancel := context.WithTimeout(t.Context(), 20*time.Second) + defer cancel() + pool, err := dbconn.NewPool(ctx, dbconn.Config{URL: dsn}) + require.NoError(t, err) + defer pool.Close() + var ssl bool + var version string + require.NoError(t, pool.QueryRow(ctx, "SELECT ssl,version FROM pg_stat_ssl WHERE pid=pg_backend_pid()").Scan(&ssl, &version)) + assert.True(t, ssl) + t.Logf("encrypted session: %s", version) +} +func TestHostedTLSDefault(t *testing.T) { + assertHostedTLS(t, hostedTLSURL(t, "", false)) +} +func TestHostedTLSRequire(t *testing.T) { + assertHostedTLS(t, hostedTLSURL(t, "require", false)) +} +func TestHostedTLSVerifyFull(t *testing.T) { + assertHostedTLS(t, hostedTLSURL(t, "verify-full", true)) +} +func TestHostedTLSRejectUntrustedCA(t *testing.T) { + dsn := hostedTLSURL(t, "verify-full", true) + u, err := url.Parse(dsn) + require.NoError(t, err) + q := u.Query() + q.Set("sslrootcert", untrustedCA(t)) + u.RawQuery = q.Encode() + ctx, cancel := context.WithTimeout(t.Context(), 20*time.Second) + defer cancel() + pool, err := dbconn.NewPool(ctx, dbconn.Config{URL: u.String()}) + if pool != nil { + defer pool.Close() + } + var untrusted x509.UnknownAuthorityError + require.ErrorAs(t, err, &untrusted, "network/auth errors are not proof of certificate rejection") +} +func TestHostedTLSRejectWrongHostname(t *testing.T) { + dsn := hostedTLSURL(t, "verify-full", true) + cfg, err := pgxpool.ParseConfig(dsn) + require.NoError(t, err) + require.NotNil(t, cfg.ConnConfig.TLSConfig) + require.False(t, cfg.ConnConfig.TLSConfig.InsecureSkipVerify) + require.Empty(t, cfg.ConnConfig.Fallbacks) + // Keep the real dial address and CA, changing only the expected identity. + cfg.ConnConfig.TLSConfig.ServerName = "wrong-host.invalid" + ctx, cancel := context.WithTimeout(t.Context(), 20*time.Second) + defer cancel() + pool, err := pgxpool.NewWithConfig(ctx, cfg) + require.NoError(t, err) + defer pool.Close() + err = pool.Ping(ctx) + var mismatch x509.HostnameError + require.ErrorAs(t, err, &mismatch) +} +func untrustedCA(t *testing.T) string { + t.Helper() + key, err := rsa.GenerateKey(rand.Reader, 2048) + require.NoError(t, err) + cert := &x509.Certificate{SerialNumber: big.NewInt(1), Subject: pkix.Name{CommonName: "Untrusted test CA"}, NotBefore: time.Now().Add(-time.Hour), NotAfter: time.Now().Add(time.Hour), IsCA: true, BasicConstraintsValid: true, KeyUsage: x509.KeyUsageCertSign} + der, err := x509.CreateCertificate(rand.Reader, cert, cert, &key.PublicKey, key) + require.NoError(t, err) + path := filepath.Join(t.TempDir(), "untrusted.crt") + require.NoError(t, os.WriteFile(path, pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der}), 0600)) + return path +}