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.
| State | What your program should do |
|---|---|
idle | Call connect() when ready |
connecting | Wait; avoid another connect() |
connected | Publish, subscribe and query |
reconnecting | Pause sending and report connection status |
failed | Report the failure; call connect() again to retry |
closing | Wait for cleanup |
closed | Create 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.