Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 3 additions & 2 deletions docs/realtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Update the README for custom-schema fetching

Update the main README alongside this behavior change: its Postgres section still states that automatic fetching applies only to the public schema and that non-public schemas retain lightweight notifications (README.md lines 770–776). Users following that contract may bind a database without setting auto_fetch=False and unexpectedly issue row queries for custom-schema notifications.

Useful? React with 👍 / 👎.

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
Expand Down
7 changes: 5 additions & 2 deletions src/volcano_sdk/realtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -768,15 +768,18 @@ 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
):
return None
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,
)

Expand Down
33 changes: 18 additions & 15 deletions tests/unit/test_realtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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(
Expand All @@ -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",
)

Expand All @@ -1108,7 +1111,7 @@ def on_insert(change: Any) -> None:

channel.on_postgres_changes(
"INSERT",
schema="public",
schema=schema,
table="messages",
callback=on_insert,
)
Expand All @@ -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",
Expand All @@ -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,
},
Expand Down
Loading