Python Client API
Full reference for useceleris-client: create_client, Channel, Segment, payload helpers, types and error codes.
pip install --pre useceleris-client
Python 3.10 to 3.14, on asyncio. Runtime dependencies: pydantic, websockets and typing-extensions. Type information ships with the package (py.typed).
create_client()
def create_client(
*,
credential_provider: CredentialProvider,
base_url: str = "wss://realtime.useceleris.com",
allow_insecure_loopback: bool = False,
connect_timeout_ms: int = 15_000,
presence_query_timeout_ms: int = 10_000,
) -> Client: ...
| Option | Type | Default | Description |
|---|---|---|---|
credential_provider | CredentialProvider | — | Required. Called for every connection attempt. |
base_url | str | production endpoint | Only for a local or self-hosted stack |
allow_insecure_loopback | bool | False | Permits ws:// for loopback hosts. Never set in production. |
connect_timeout_ms | int | 15000 | Covers credential acquisition and handshake together |
presence_query_timeout_ms | int | 10000 | Deadline for one presence_list() |
Performs no network work. Options are validated strictly: both timeouts are positive int milliseconds, and a bool, float or str is refused rather than converted. base_url is the realtime socket endpoint, not your HTTP credential endpoint; it must be wss://, or ws:// to a loopback host with allow_insecure_loopback=True, and carry no credentials, query string or fragment. Client(...) takes the same keyword arguments.
Client
class Client:
def channel(self, reference: str) -> Channel: ...
channel() is side-effect free and returns a new connection handle every call. A reference is 1 to 255 characters of [A-Za-z0-9_-].
Channel
class Channel:
@property
def state(self) -> ChannelState: ...
async def connect(self) -> None: ...
async def close(self) -> None: ...
def segment(self, segment_id: str) -> Segment: ...
def default_segment(self) -> Segment: ... # segment("default")
def events(self) -> ChannelEventHandler: ...
ChannelState = Literal[
"idle", "connecting", "connected", "reconnecting", "failed", "closing", "closed"
]
| Member | Notes |
|---|---|
state | The current state |
connect() | Raises OperationInProgress if one is already running, NotConnected once closed. An initial failure both raises and sets failed. Cancelling the awaiting task abandons the attempt and sets failed. |
close() | Idempotent and terminal, with a 5 s graceful budget. Every call waits for the same close. Ends every segment object with the channel. |
segment(id) | Side-effect free. The id is required: non-empty, no CR or LF. |
default_segment() | The "default" segment every connection joins automatically. |
failed is recoverable by calling connect() again; closed is not. Initial connection failures are not retried automatically. See Reconnection and Recovery for triggers, deadlines, retry budgets and restored subscriptions.
A successful connect() joins default; default_segment().on_message(listener) receives its messages without subscribe(). Read and write permissions still apply, and default membership includes no other segment and no presence subscription. The automatic join also happens on reconnect.
Use a channel from the event loop thread that connected it.
Segment
class Segment:
@property
def segment_id(self) -> str: ...
def subscribe(self) -> Subscription: ...
def on_message(self, listener: MessageListener) -> Callable[[], None]: ...
def on_presence(
self, listener: Callable[[PresenceEvent], None]
) -> Callable[[], None]: ...
async def publish(self, payload: bytes, *, message_id: str | None = None) -> None: ...
def subscribe_presence(self) -> Subscription: ...
async def presence_list(self, *, page: int, per_page: int) -> PresencePage: ...
| Method | Notes |
|---|---|
subscribe() | A real server operation. Never sent for "default", which is joined automatically. |
on_message() / on_presence() | Return a function that removes that listener. Segment objects for the same segment share one listener set. async def listeners raise ConfigurationError. |
publish() | Returns on local acceptance only. payload must be bytes; an empty payload is valid. Raises ConfigurationError for a command over 2 MiB encoded; your plan's payload cap is enforced by the server afterwards. Joins the segment server-side. Cancelling the awaiting task withdraws a publish not yet sent. |
subscribe_presence() | Also joins the segment for messages; cancelling presence does not leave it. |
presence_list() | One in flight per channel; a second raises OperationInProgress. page is 1 to 2147483647, per_page 1 to 100, both validated rather than clamped. A query the server refuses raises a ServerError with sub_type "PRES_LIST" at once. A timeout or cancellation frees the slot and leaves the connection up. |
Payload helpers
def text_payload(value: str) -> bytes: ...
def json_payload(value: object) -> bytes: ...
def read_text(payload: bytes) -> str: ...
def read_json(payload: bytes) -> Any: ... # returns the JSON value; does not validate
def create_payload_codec(
*, encode: Callable[[T], bytes], decode: Callable[[bytes], T]
) -> PayloadCodec[T]: ...
@dataclass(frozen=True)
class PayloadCodec(Generic[T]):
encode_payload: Callable[[T], bytes]
read_payload: Callable[[bytes], T]
text_payload raises ConfigurationError for text with unpaired surrogates. json_payload writes compact JSON (as JSON.stringify does) and raises ConfigurationError for NaN, infinities, circular structures and values JSON cannot represent. read_text and read_json raise ConfigurationError for invalid UTF-8 or JSON. Errors raised by your own encode or decode propagate unchanged. No size check happens here: publish() enforces the encoded-command bound.
Types
@dataclass(frozen=True)
class MessageMetadata:
token_reference: str
segment_id: str
message_id: str # always present; the SDK generates one if you omit it
timestamp: int # Unix milliseconds
MessageListener = Callable[[bytes, MessageMetadata], None]
@dataclass(frozen=True)
class PresenceEvent:
segment_id: str
token_reference: str
connection_id: str
joined: bool # False = left
timestamp: int
@dataclass(frozen=True)
class PresencePage:
segment_id: str
total: int
per_page: int
current_page: int
from_: int # from_ > to past the last page
to: int
connections: tuple[PresenceConnection, ...]
@dataclass(frozen=True)
class PresenceConnection:
token_reference: str
connection_id: str
timestamp: int
@dataclass(frozen=True)
class ServerNotice:
timestamp: int
payload: bytes
@dataclass(frozen=True)
class RecoveryEvent:
retry_index: int
possible_gaps: Literal[True] = True
possible_duplicates: Literal[True] = True
@dataclass(frozen=True)
class CredentialRequest:
channel_reference: str
reason: Literal["initial", "reconnect"]
disconnected_at: int | None = None # Unix ms at the start of the outage
replay_lookback_ms: int | None = None # at most 4294967295
@dataclass(frozen=True)
class Credentials:
payload: str # both left out of repr()
signature: str
CredentialProvider = Callable[[CredentialRequest], Awaitable[Credentials]]
class Subscription:
def cancel(self) -> None: ... # idempotent
Every number is an int: presence figures are signed 32-bit values as received, and timestamps signed 64-bit Unix milliseconds. from_ carries a trailing underscore because from is a Python keyword. A credential provider must return a Credentials instance with non-empty, well-formed strings.
ChannelEventHandler
class ChannelEventHandler:
def on_state_change(self, listener: Callable[[ChannelState], None]) -> Callable[[], None]: ...
def on_recovery(self, listener: Callable[[RecoveryEvent], None]) -> Callable[[], None]: ...
def on_notice(self, listener: Callable[[ServerNotice], None]) -> Callable[[], None]: ...
def on_error(self, listener: Callable[[ChannelError], None]) -> Callable[[], None]: ...
ChannelError = ConfigurationError | CelerisConnectionError | ProtocolError | ServerError
Each registration returns a function that removes it. Listeners run synchronously in registration order with no queue: a slow listener delays delivery rather than growing a backlog. A listener that raises, or that returns a coroutine, is contained and reported through on_error; failures of error listeners themselves are swallowed. A recovery event follows the connected state change; it does not signal replay completion.
on_notice delivers raw server prose, such as subscription acknowledgements and refusals. It carries no segment id and no correlation id. Do not parse it.
Errors
Every error is a CelerisError with a code and a message. Three kinds are raised by the SDK itself: ConfigurationError (code "Configuration"), CelerisConnectionError (a code from the table below) and ProtocolError (code "ProtocolError", with field and offset). Their messages say what failed and which rule or limit it broke, such as Invalid channel reference. Must not be empty., but never carry your input, credentials or server text, and never chain an exception that holds them. The fourth, ServerError, carries an error the server sent, described below. Match on code, and on type for server errors, never on message text.
code | Raised when |
|---|---|
Configuration | The call is wrong, or a credential provider returned invalid credentials. Fix the code; do not retry. |
Timeout | The connect deadline elapsed, or a presence query timed out |
Cancelled | The channel was closed during a connect, a queued publish or a presence query |
Transport | Handshake or socket failure, or a credential provider that raised |
NotConnected | Publishing or querying while not connected, or acting on a closed channel |
Backpressure | 64 publishes are already waiting, or a presence query while the writer is full or paused |
OperationInProgress | A concurrent connect(), or a second presence_list() |
DeliveryUnknown | The send failed after hand-off: acceptance genuinely unknown |
ProtocolError | One received frame could not be decoded; the error also has field and offset |
Cancelling a task you own raises asyncio.CancelledError, not an SDK error. ProtocolError means one message could not be decoded: that message is dropped, the error is reported, and the connection stays up. A command this SDK version does not recognise is skipped silently, so a newer server cannot break a deployed client.
There is no authentication code: the WebSocket handshake status is not exposed, so a refused handshake is reported as Transport rather than guessed at.
ServerError
ServerErrorType = Literal[
"ParserError",
"SendError",
"PermissionDeniedError",
"RateLimitError",
"MessageSizeLimitError",
"InternalError",
]
ServerErrorResource = str | int | tuple["ServerErrorResource", ...] | None
class ServerError(CelerisError): # code == "Server"
type: str # one of ServerErrorType, or a type a newer server adds
sub_type: str | None # the command it answers, e.g. "PUB"
message: str # the server's own text, unchanged
resource: ServerErrorResource # what that command names, e.g. the segment
An error the server sent, every field as sent, and the connection stays up. It arrives through on_error, after the call that caused it, except for a presence query error, which raises from presence_list() instead. Keep a default branch for types a newer server adds. Check resource's type before using it.
| Error | sub_type | resource |
|---|---|---|
| A permission denial for a publish, subscribe or presence subscription | "PUB", "SUB", "UNSUB", "PRES_SUB" or "PRES_UNSUB" | The segment id |
| A refused or failed presence query | "PRES_LIST" | The query's request id, which the SDK sets |
| Anything else | None | None |
type | Sent when |
|---|---|
PermissionDeniedError | The token lacks access to what the command tried |
MessageSizeLimitError | A publish exceeded your plan's payload cap |
RateLimitError | A per-second, per-hour, per-month or per-connection message limit was exceeded; the SDK pauses, then resends recent commands |
ParserError | The server could not parse a command |
SendError | The server failed to send |
InternalError | A fault inside the server, such as a failed presence read |
Limits
| Limit | Value |
|---|---|
| Outbound command | 2 MiB encoded, rejected before any write |
| Publish payload | Per plan: 64 KiB free, 128 KiB standard, 512 KiB pro, 1024 KiB prime; enforced by the server |
| Writer bounds | 64 pending commands, 2 MiB of unsent data; up to 64 more publishes wait in a queue, behind subscription changes |
| Rate-limit recovery | Pause of 1 s plus growing jitter (at most 31 s), then a resend of the last 2 s of commands (at most 64 publishes, each once); after 8 limits in a row, dropped subscriptions are retried after 1 min, doubling to 1 h |
| Duplicate filtering | 1024 message ids per channel |
| Reconnect | 10 failed attempts per budget, full-jitter delays up to 30 s; the budget resets at a drop after at least 60 s connected |
| Connect deadline | 15 s, covering credentials and handshake |
| Presence query | 10 s, one in flight per channel |
| Close | 5 s graceful, then forced |
Received messages are never size-checked: they are already in memory when they arrive, so the SDK processes whatever the server sends. A publish rejected for its plan payload limit still counts toward your usage; a local encoded-command rejection sends nothing.
What this SDK will not do
- Confirm delivery. There are no receipts anywhere in the protocol.
- Queue while offline. Publishes are resent only after a rate limit, once each, with their original ID.
- Provide durable history or global ordering. Replay is a bounded window, not a cursor.
- Sign anything. Signing lives in useceleris-server.