From 65d96894207719784bd81fe8b93880774f12277b Mon Sep 17 00:00:00 2001 From: Sean Keever <33592180+swkeever@users.noreply.github.com> Date: Fri, 18 Sep 2026 15:01:43 -0400 Subject: [PATCH] test(realtime): stage presence membership acceptance --- docs/README.md | 3 + docs/realtime.md | 123 ++++++++++++++++++++++ features/presence_membership.py | 80 ++++++++++++++ features/staged/realtime-presence.feature | 8 ++ features/steps/sdk_contract_steps.py | 19 ++++ tests/unit/test_contract_bindings.py | 31 ++++++ 6 files changed, 264 insertions(+) create mode 100644 docs/realtime.md create mode 100644 features/presence_membership.py create mode 100644 features/staged/realtime-presence.feature diff --git a/docs/README.md b/docs/README.md index e7a9dad5..05ee39b9 100644 --- a/docs/README.md +++ b/docs/README.md @@ -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. diff --git a/docs/realtime.md b/docs/realtime.md new file mode 100644 index 00000000..739d0849 --- /dev/null +++ b/docs/realtime.md @@ -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. diff --git a/features/presence_membership.py b/features/presence_membership.py new file mode 100644 index 00000000..1e59f535 --- /dev/null +++ b/features/presence_membership.py @@ -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 diff --git a/features/staged/realtime-presence.feature b/features/staged/realtime-presence.feature new file mode 100644 index 00000000..17dc2b3c --- /dev/null +++ b/features/staged/realtime-presence.feature @@ -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 diff --git a/features/steps/sdk_contract_steps.py b/features/steps/sdk_contract_steps.py index ab8093bc..86089a5a 100644 --- a/features/steps/sdk_contract_steps.py +++ b/features/steps/sdk_contract_steps.py @@ -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 @@ -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] diff --git a/tests/unit/test_contract_bindings.py b/tests/unit/test_contract_bindings.py index e2234cf5..284885be 100644 --- a/tests/unit/test_contract_bindings.py +++ b/tests/unit/test_contract_bindings.py @@ -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" + ), } @@ -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 + )