diff --git a/README.md b/README.md index 9d271d28..b354f4f2 100644 --- a/README.md +++ b/README.md @@ -768,12 +768,12 @@ 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 -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 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, },