Documentation

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: ...
OptionTypeDefaultDescription
credential_providerCredentialProvider—Required. Called for every connection attempt.
base_urlstrproduction endpointOnly for a local or self-hosted stack
allow_insecure_loopbackboolFalsePermits ws:// for loopback hosts. Never set in production.
connect_timeout_msint15000Covers credential acquisition and handshake together
presence_query_timeout_msint10000Deadline 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"
]
MemberNotes
stateThe 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: ...
MethodNotes
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.

codeRaised when
ConfigurationThe call is wrong, or a credential provider returned invalid credentials. Fix the code; do not retry.
TimeoutThe connect deadline elapsed, or a presence query timed out
CancelledThe channel was closed during a connect, a queued publish or a presence query
TransportHandshake or socket failure, or a credential provider that raised
NotConnectedPublishing or querying while not connected, or acting on a closed channel
Backpressure64 publishes are already waiting, or a presence query while the writer is full or paused
OperationInProgressA concurrent connect(), or a second presence_list()
DeliveryUnknownThe send failed after hand-off: acceptance genuinely unknown
ProtocolErrorOne 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.

Errorsub_typeresource
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 elseNoneNone
typeSent when
PermissionDeniedErrorThe token lacks access to what the command tried
MessageSizeLimitErrorA publish exceeded your plan's payload cap
RateLimitErrorA per-second, per-hour, per-month or per-connection message limit was exceeded; the SDK pauses, then resends recent commands
ParserErrorThe server could not parse a command
SendErrorThe server failed to send
InternalErrorA fault inside the server, such as a failed presence read

Limits

LimitValue
Outbound command2 MiB encoded, rejected before any write
Publish payloadPer plan: 64 KiB free, 128 KiB standard, 512 KiB pro, 1024 KiB prime; enforced by the server
Writer bounds64 pending commands, 2 MiB of unsent data; up to 64 more publishes wait in a queue, behind subscription changes
Rate-limit recoveryPause 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 filtering1024 message ids per channel
Reconnect10 failed attempts per budget, full-jitter delays up to 30 s; the budget resets at a drop after at least 60 s connected
Connect deadline15 s, covering credentials and handshake
Presence query10 s, one in flight per channel
Close5 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.

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