From 4ccd370a7925912ff848837515a2e4428ec3997b Mon Sep 17 00:00:00 2001 From: Sean Keever <33592180+swkeever@users.noreply.github.com> Date: Fri, 18 Sep 2026 15:54:41 -0400 Subject: [PATCH 1/3] fix(realtime): fetch rows from custom Postgres schemas --- docs/realtime.md | 5 +++-- src/volcano_sdk/realtime.py | 7 +++++-- tests/unit/test_realtime.py | 33 ++++++++++++++++++--------------- 3 files changed, 26 insertions(+), 19 deletions(-) diff --git a/docs/realtime.md b/docs/realtime.md index e1bc3525..754cccb8 100644 --- a/docs/realtime.md +++ b/docs/realtime.md @@ -102,11 +102,12 @@ await client.realtime.remove_channel("public:messages", channel_type="postgres") Insert a row from another client during the listening period. Automatic row lookup requires a primary key named `id`. It reads the current row after the notification; rapid updates may have already changed that row. -Matching insert and update notifications in the `public` schema can fetch full rows using the subscription's user token. +Matching insert and update notifications can fetch full rows using the subscription's user token. Compatible row lookups are batched while publication order is preserved. Defaults are a 20 millisecond window and 50 rows; set `fetch_batch_window_ms` and `fetch_max_batch_size` on the channel to change them. Set `auto_fetch=False` when first creating the channel or `set_database_name(None)` to retain lightweight notifications without row lookups. -Missing rows, failed lookups, and non-public schemas retain the lightweight notification. +For a custom schema, use its `schema:table` channel name and matching schema/table filters; row lookups preserve that schema. +Missing rows and failed lookups retain the lightweight notification. Deletes use `old_record` or the row ID and do not query the database. ## Handle connection changes diff --git a/src/volcano_sdk/realtime.py b/src/volcano_sdk/realtime.py index cee15fb7..fa1a9dd8 100644 --- a/src/volcano_sdk/realtime.py +++ b/src/volcano_sdk/realtime.py @@ -768,7 +768,6 @@ def _postgres_fetch_request( not self._fetch_config.enabled or change.mode != "lightweight" or change.type == "DELETE" - or change.schema != "public" or change.id is None or database_name is None ): @@ -776,7 +775,11 @@ def _postgres_fetch_request( return _PostgresFetchRequest( database_name=database_name, access_token=self._realtime._connection_token(), - table=change.table, + table=( + change.table + if change.schema == "public" + else f"{change.schema}.{change.table}" + ), row_id=change.id, ) diff --git a/tests/unit/test_realtime.py b/tests/unit/test_realtime.py index bb621873..a9fa3f28 100644 --- a/tests/unit/test_realtime.py +++ b/tests/unit/test_realtime.py @@ -1041,17 +1041,19 @@ async def scenario() -> None: table="messages", row_id=42, ) - assert ( - channel._postgres_fetch_request( - realtime_module.PostgresChange( - type="INSERT", - schema="private", - table="messages", - id=42, - mode="lightweight", - ) + assert channel._postgres_fetch_request( + realtime_module.PostgresChange( + type="INSERT", + schema="private", + table="messages", + id=42, + mode="lightweight", ) - is None + ) == realtime_module._PostgresFetchRequest( + database_name="next", + access_token="access-1", + table="private.messages", + row_id=42, ) assert ( channel._postgres_fetch_request( @@ -1083,7 +1085,8 @@ async def scenario() -> None: asyncio.run(scenario()) -def test_realtime_fetches_lightweight_postgres_rows() -> None: +@pytest.mark.parametrize("schema", ["public", "private"]) +def test_realtime_fetches_lightweight_postgres_rows(schema: str) -> None: official = FakeCentrifugeClient() transport = RealtimeDatabaseTransport([{"id": 42, "body": "fetched"}]) client = VolcanoClient( @@ -1098,7 +1101,7 @@ async def scenario() -> None: changes: list[Any] = [] received = asyncio.Event() channel = client.realtime.channel( - "public:messages", + f"{schema}:messages", channel_type="postgres", ) @@ -1108,7 +1111,7 @@ def on_insert(change: Any) -> None: channel.on_postgres_changes( "INSERT", - schema="public", + schema=schema, table="messages", callback=on_insert, ) @@ -1119,7 +1122,7 @@ def on_insert(change: Any) -> None: await subscription.emit( { "type": "INSERT", - "schema": "public", + "schema": schema, "table": "messages", "id": 42, "mode": "lightweight", @@ -1136,7 +1139,7 @@ def on_insert(change: Any) -> None: "authorization": "access-1", "database_name": "app", "body": { - "table": "messages", + "table": "messages" if schema == "public" else f"{schema}.messages", "filters": [{"column": "id", "operator": "in", "value": [42]}], "limit": 1, }, From 9cee2861393cf9b9be3a1d025a05b6a2020ca74c Mon Sep 17 00:00:00 2001 From: Sean Keever <33592180+swkeever@users.noreply.github.com> Date: Fri, 18 Sep 2026 16:02:48 -0400 Subject: [PATCH 2/3] docs(realtime): align custom-schema fetching guidance --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 9d271d28..86f41d33 100644 --- a/README.md +++ b/README.md @@ -768,7 +768,7 @@ stop_changes() ``` Binding a database automatically fetches the matching row for lightweight -`INSERT` and `UPDATE` notifications in the `public` schema. The fetch uses the +`INSERT` and `UPDATE` notifications in public or custom schemas. The fetch uses the realtime connection's RLS-scoped access token. Compatible row lookups are batched while callback delivery preserves publication order. If the row is absent or the query fails, the callback receives the lightweight From 38ce6752f45994ca6fd63394ae867409d2fdb2d9 Mon Sep 17 00:00:00 2001 From: Sean Keever <33592180+swkeever@users.noreply.github.com> Date: Fri, 18 Sep 2026 16:03:48 -0400 Subject: [PATCH 3/3] docs(realtime): remove obsolete custom-schema fallback claim --- README.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 86f41d33..b354f4f2 100644 --- a/README.md +++ b/README.md @@ -772,8 +772,8 @@ Binding a database automatically fetches the matching row for lightweight realtime connection's RLS-scoped access token. Compatible row lookups are batched while callback delivery preserves publication order. If the row is absent or the query fails, the callback receives the lightweight -notification with its `id` and `mode` intact. Non-public schemas also retain -that lightweight form. Lightweight deletes never query the database; they +notification with its `id` and `mode` intact. Custom-schema row lookups preserve +the schema from the notification. Lightweight deletes never query the database; they preserve `old_record`, or provide `{"id": change.id}` when no old row was included. Tune a channel's batching with `fetch_batch_window_ms` and `fetch_max_batch_size`; the defaults are 20 milliseconds and 50 rows. Set