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
3 changes: 3 additions & 0 deletions docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -124,3 +124,6 @@ Read the [database guide](./database.md) for projection, filters, ordered pagina
For uploads, visibility, and resumable sessions, see [Storage](./storage.md).

See [Logs](./logs.md) for project-token authentication, search, pagination, and activity.


See [Realtime](./realtime.md) for broadcasts, presence, database changes, and shutdown.
123 changes: 123 additions & 0 deletions docs/realtime.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
---
title: Realtime
description: Send broadcasts, observe presence, and receive database changes with Python async channels.
order: 8
---

Subscribe to a broadcast channel and send a JSON publication using Python's async lifecycle.
Sign in with the [quickstart credentials](./README.md) and enable the project's realtime capabilities and access policies.

```python
import asyncio
import os

from volcano_sdk import VolcanoClient

client = VolcanoClient(
anon_key=os.environ["VOLCANO_ANON_KEY"],
api_url=os.environ.get("VOLCANO_API_URL", "https://api.volcano.dev"),
)
client.auth.sign_in(
email=os.environ["VOLCANO_USER_EMAIL"],
password=os.environ["VOLCANO_USER_PASSWORD"],
)


async def main():
channel = client.realtime.channel("updates")
channel.on("message", lambda message: print(message))
try:
await channel.subscribe()
await channel.send({"event": "message", "value": "hello"})
await asyncio.sleep(2)
finally:
await client.realtime.disconnect()


asyncio.run(main())
```

`subscribe()` waits for the server acknowledgement; a presence subscription also waits for its initial roster.
Channel names receive their type prefix, so `updates` becomes `broadcast:updates`.
Use the same event loop for a client's realtime operations.
The following channel examples belong inside an async function with an authenticated `client`.

## Pause or remove channels

```python
channel = client.realtime.channel("updates")
channel.on("message", print)
await channel.subscribe()
await channel.unsubscribe()
await channel.subscribe()
await client.realtime.remove_channel("updates")
```

Unsubscribe pauses delivery while retaining handlers and the in-memory broadcast recovery position.
Subscribe resumes and requests missed broadcasts when retained history is available.
Pausing discards queued callbacks; a callback already running may finish.
Removal forgets the channel and its recovery position.
`remove_all_channels()` removes all channels while leaving the shared connection available; `disconnect()` closes it.
Application code owns any work its callbacks start and should await that work during shutdown.

## Observe presence

```python
lobby = client.realtime.channel("lobby", channel_type="presence")
stop_sync = lobby.on_presence_sync(lambda state: print("Online", len(state)))
lobby.on("join", lambda info: print("Joined", info.user))
lobby.on("leave", lambda info: print("Left", info.user))
await lobby.subscribe()
await lobby.track({"status": "online"})
print(lobby.tracked_state)
print(lobby.get_presence_state())
await client.realtime.remove_channel("lobby", channel_type="presence")
stop_sync()
```

Presence identity and metadata come from the authenticated user.
`track()` stores optional local application state; it does not replace server-managed presence metadata.
The roster maps connection IDs to immutable client identity and user metadata snapshots.
One user can have several connections. The original sync callback observes connections joining and leaving.
The roster is refreshed after reconnection; a failed roster query reports an error and clears the snapshot.

## Subscribe to database changes

Use an existing `app` database and `public.messages` table with realtime and suitable Row-Level Security policies configured:

```python
client.realtime.set_database_name("app")
changes = client.realtime.channel("public:messages", channel_type="postgres")
stop_changes = changes.on_postgres_changes(
"INSERT",
schema="public",
table="messages",
callback=lambda change: print(change.record),
)
await changes.subscribe()
await asyncio.sleep(30)
stop_changes()
await client.realtime.remove_channel("public:messages", channel_type="postgres")
```

Insert a row from another client during the listening period.
Matching insert and update notifications in the `public` schema 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.
Use `auto_fetch=False` or `set_database_name(None)` to retain lightweight notifications without row lookups.
Missing rows, failed lookups, and non-public schemas retain the lightweight notification.
Deletes use `old_record` or the row ID and do not query the database.

## Handle connection changes

```python
stop_errors = client.realtime.on_error(lambda context: print(context.message))
stop_connected = client.realtime.on_connect(lambda context: print(context.client))
stop_disconnected = client.realtime.on_disconnect(lambda context: print(context.reason))
```

Each registration returns an idempotent unsubscribe function.
Callbacks receive immutable contexts and run outside connection processing.
Broadcast recovery stays within one client lifetime and one authenticated session lineage; it is not persisted across processes.
After signing in again or changing users, disconnect before subscribing for the new session.
Do not use realtime delivery as a durable record of every database change.
80 changes: 80 additions & 0 deletions features/presence_membership.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
from __future__ import annotations

import asyncio
from typing import TYPE_CHECKING

if TYPE_CHECKING:
from collections.abc import Callable, Mapping

from contract_support import ContractWorld

from volcano_sdk.realtime import Channel, RealtimePresenceInfo


class PresenceObserver:
def __init__(self, channel: Channel, user_id: str) -> None:
self.channel = channel
self.user_id = user_id
self.changed = asyncio.Event()
self.snapshots: list[set[str]] = []
self.unsubscribe = channel.on_presence_sync(self.record)

def record(self, state: Mapping[str, RealtimePresenceInfo]) -> None:
self.snapshots.append(set(state))
self.changed.set()

async def wait(self, predicate: Callable[[], bool]) -> None:
async with asyncio.timeout(10):
while True:
self.changed.clear()
if predicate():
return
await self.changed.wait()

async def roster(self, count: int) -> set[str]:
await self.wait(lambda: len(self.channel.get_presence_state()) == count)
state = self.channel.get_presence_state()
assert all(info.user == self.user_id for info in state.values())
assert all(key == info.client and key for key, info in state.items())
return set(state)


async def verify_presence_membership(world: ContractWorld) -> list[int]:
first, second = [
client.realtime.channel(
world.realtime_channel + "-presence", channel_type="presence"
)
for client in world.realtime_clients
]
first_observer = PresenceObserver(first, world.fixture["user_id"])
second_observer = PresenceObserver(second, world.fixture["user_id"])
try:
await first.subscribe()
initial = await first_observer.roster(1)
await second.subscribe()
joined = await first_observer.roster(2)
assert joined == await second_observer.roster(2)
assert initial < joined
await second.unsubscribe()
assert await first_observer.roster(1) == initial
await first_observer.wait(
lambda: _observed_membership(first_observer.snapshots, initial, joined)
)
return [1, 2, 1]
finally:
first_observer.unsubscribe()
second_observer.unsubscribe()
await asyncio.gather(first.unsubscribe(), second.unsubscribe())


def _observed_membership(
snapshots: list[set[str]], initial: set[str], joined: set[str]
) -> bool:
expected = iter([initial, joined, initial])
target: set[str] | None = next(expected)
for state in snapshots:
if state == target:
target = next(expected, None)
if target is None:
return True
return False
8 changes: 8 additions & 0 deletions features/staged/realtime-presence.feature
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
Feature: SDK presence membership contract

@realtime @SDK-REALTIME-003
Scenario: The original presence handler observes another connection joining and leaving
Given two authenticated realtime clients
When one presence client joins and leaves while the other remains subscribed
Then the SDK operation succeeds
And both rosters identify the contract user and the original handler observes membership changes
19 changes: 19 additions & 0 deletions features/steps/sdk_contract_steps.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
classify_error,
)
from logs_contract import LogContract
from presence_membership import verify_presence_membership

from volcano_sdk import NotFoundError, Session, VolcanoClient

Expand Down Expand Up @@ -1159,3 +1160,21 @@ def verify_contract_logs(context: Any) -> None:
@then("activity counts exactly that event in its function and level buckets")
def verify_contract_log_activity(context: Any) -> None:
context.logs_contract.verify_activity(_world(context).last_outcome.value)


@when("one presence client joins and leaves while the other remains subscribed")
def observe_presence_membership(context: Any) -> None:
world = _world(context)
if world.last_outcome is not None and not world.last_outcome.ok:
return
world.record(lambda: world.run(verify_presence_membership(world)))


@then(
"both rosters identify the contract user "
"and the original handler observes membership changes"
)
def verify_presence_rosters(context: Any) -> None:
world = _world(context)
assert world.last_outcome is not None
assert world.last_outcome.value == [1, 2, 1]
31 changes: 31 additions & 0 deletions tests/unit/test_contract_bindings.py
Original file line number Diff line number Diff line change
Expand Up @@ -241,6 +241,11 @@ def test_every_contract_phrase_is_bound_verbatim() -> None:
"the client reads matching log activity within 120 seconds",
"all three structured events retain their metadata without duplicates",
"activity counts exactly that event in its function and level buckets",
"one presence client joins and leaves while the other remains subscribed",
(
"both rosters identify the contract user "
"and the original handler observes membership changes"
),
}


Expand Down Expand Up @@ -525,3 +530,29 @@ def test_log_bounds_allow_server_clock_skew(server_skew_seconds: int) -> None:
contract.emit(1)
assert datetime.fromisoformat(contract.request["start_time"]) < server_time
assert server_time < datetime.fromisoformat(contract.request["end_time"])


def test_staged_presence_feature_matches_proposed_shared_source() -> None:
staged = ROOT / "features" / "staged" / "realtime-presence.feature"
expected = "b4429f6e3df60a6a98be4daf1d8517e2cd7cee651f9eb6463a1090ab49a102b5"
assert hashlib.sha256(staged.read_bytes()).hexdigest() == expected


@pytest.mark.parametrize(
("snapshots", "expected"),
[
([{"first"}, {"first", "second"}, {"first"}], True),
([{"first"}, {"first", "second"}], False),
([{"first", "second"}, {"first"}], False),
],
)
def test_presence_requires_original_handler_membership_sequence(
snapshots: list[set[str]], *, expected: bool
) -> None:
module = _load_module(
"presence_membership", ROOT / "features" / "presence_membership.py"
)
assert (
module._observed_membership(snapshots, {"first"}, {"first", "second"})
is expected
)
Loading