Documentation

Connecting and Lifecycle

Supply credentials, connect a channel, observe recovery, and release resources.

Before you start: install useceleris-client (see the overview) and have an endpoint that signs credentials, as described in Authentication. A trusted backend can sign its own instead; see Server-side usage.

1. Supply credentials

The SDK has a built-in production endpoint. Give it a credential provider: a coroutine that obtains signed credentials from your application for each connection attempt.

import asyncio
import json
import urllib.request

from useceleris_client import CredentialRequest, Credentials, create_client


def post_json(url: str, body: dict[str, object]) -> dict[str, str]:
    request = urllib.request.Request(
        url,
        data=json.dumps(body).encode(),
        headers={"content-type": "application/json", "authorization": "Bearer ..."},
        method="POST",
    )

    with urllib.request.urlopen(request, timeout=10) as response:
        result: dict[str, str] = json.load(response)

    return result


async def provide_credentials(request: CredentialRequest) -> Credentials:
    body = await asyncio.to_thread(
        post_json,
        "https://your-app.example/api/realtime-credentials",
        {
            "channel_reference": request.channel_reference,
            "replay_lookback_ms": request.replay_lookback_ms,
        },
    )
    return Credentials(payload=body["payload"], signature=body["signature"])


client = create_client(credential_provider=provide_credentials)
channel = client.channel("demo-room")

The provider receives a CredentialRequest with channel_reference and reason ("initial" or "reconnect"). On recovery it also carries disconnected_at, Unix milliseconds, and replay_lookback_ms, a duration in milliseconds; both are None on the first attempt. Any HTTP client works; this one uses only the standard library.

Return the payload and signature exactly as your backend signed them. Do not generate, cache or inspect them. The provider runs on every connection attempt, including recovery. If it raises, the SDK reports a safe Transport error; your exception's details are not passed on.

Constructing a client, channel or segment opens no connection. Use base_url only for a local or self-hosted Celeris, not for your credential endpoint. See the Python client API for the options.

2. Observe state, then connect

from useceleris_client import CelerisConnectionError, ChannelError, ChannelState


def show_state(state: ChannelState) -> None:
    print("Connection:", state)


def show_error(error: ChannelError) -> None:
    print("Error:", error.code, error)


events = channel.events()
stop_state = events.on_state_change(show_state)
stop_errors = events.on_error(show_error)

try:
    await channel.connect()
except CelerisConnectionError as error:
    print("Initial connection failed:", error.code)

Expected: connecting, then connected. A failed first attempt raises from connect() and sets failed, without also reporting to on_error. It is not retried automatically.

Once connected, the channel is a member of default. Use channel.default_segment() for that stream; do not open another connection or call subscribe() just to join it.

StateWhat your program should do
idleCall connect() when ready
connectingWait; avoid another connect()
connectedPublish, subscribe and query
reconnectingPause sending and report connection status
failedReport the failure; call connect() again to retry
closingWait for cleanup
closedCreate a new channel if you need another connection

Calling connect() while connected, connecting or reconnecting raises OperationInProgress.

To give up on a connection attempt, cancel the task awaiting it. connect() then raises CancelledError and the channel is failed:

connecting = asyncio.create_task(channel.connect())
await asyncio.sleep(0)  # the attempt is under way
connecting.cancel()

try:
    await connecting
except asyncio.CancelledError:
    print(channel.state)  # failed

A task cancelled before it first runs never starts the attempt, and the channel stays idle. asyncio.wait_for(channel.connect(), timeout=5) works the same way. The connect_timeout_ms option, 15 seconds by default, bounds credentials and handshake together.

3. Respond to recovery

Register a recovery listener once, when setting up the channel:

from useceleris_client import RecoveryEvent


def recovered(event: RecoveryEvent) -> None:
    print("Recovered after retry", event.retry_index)
    # Reload your application's authoritative state here.


stop_recovery = events.on_recovery(recovered)

retry_index is the number of failed reconnect attempts already used in the current budget: 0 when recovery succeeds at the first try. possible_gaps and possible_duplicates are always True. Reload what matters from your own API, and reconcile live updates with an application-owned version.

Listeners survive reconnects; do not register them again on every connected state. Catch errors inside any task you start from a listener: the SDK contains failures of the listener itself, not of work it starts.

See Reconnection and Recovery for triggers, timing, restored interests and replay. The recovery event is not a replay-complete signal.

4. Release resources

Each registration returns a function that removes it; subscriptions have cancel(). The channel owns the socket.

await channel.close()
stop_recovery()
stop_errors()
stop_state()

close() is idempotent and permanent, with a five-second graceful budget. It cancels recovery and ends the channel's segment objects too. Do not reuse a closed channel; obtain a new one from the client.

Next: Messaging.

We use Google Analytics cookies to understand how people use Celeris, only if you allow it. See our Cookie Policy.