Skip to content

Websocket Client

solana.rpc.websocket_api

Confirmed Solana WebSocket subscriptions and JSON-RPC request dispatch.

ConnectionState

Lifecycle of the RPC dispatcher, independent of the wire protocol state.

Source code in src/solana/rpc/websocket_api.py
127
128
129
130
131
class ConnectionState(Enum):
    """Lifecycle of the RPC dispatcher, independent of the wire protocol state."""

    OPEN = "open"
    CLOSED = "closed"

OverflowPolicy

What a full notification queue does to the next notification.

Blocking the reader is deliberately not offered: it also stops RPC responses, so unsubscribe -- the one call that could drain the flood -- would never be confirmed, and the unread socket stalls the closing handshake. Both lossy policies count what they discard in :attr:SolanaWsClient.dropped_notifications.

Source code in src/solana/rpc/websocket_api.py
134
135
136
137
138
139
140
141
142
143
144
145
146
class OverflowPolicy(StrEnum):
    """What a full notification queue does to the next notification.

    Blocking the reader is deliberately not offered: it also stops RPC
    responses, so ``unsubscribe`` -- the one call that could drain the flood --
    would never be confirmed, and the unread socket stalls the closing
    handshake. Both lossy policies count what they discard in
    :attr:`SolanaWsClient.dropped_notifications`.
    """

    RAISE = "raise"
    DROP_OLDEST = "drop_oldest"
    DROP_NEWEST = "drop_newest"

SolanaWsClient

One reader dispatches RPC responses and delivers typed notifications.

Use :meth:connect to open and manage the connection. Subscribe helpers await confirmation and return :class:Subscription; :meth:recv returns notifications only. Fatal errors abort every subscription on this connection.

Source code in src/solana/rpc/websocket_api.py
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
class SolanaWsClient:
    """One reader dispatches RPC responses and delivers typed notifications.

    Use :meth:`connect` to open and manage the connection. Subscribe helpers
    await confirmation and return :class:`Subscription`; :meth:`recv` returns
    notifications only. Fatal errors abort every subscription on this connection.
    """

    def __init__(
        self,
        uri: str = "ws://localhost:8900",
        *,
        request_timeout: float = 10.0,
        notification_queue_size: int = 10_000,
        overflow: OverflowPolicy = OverflowPolicy.RAISE,
        **kwargs: Any,
    ) -> None:
        """Create a client; the WebSocket is opened by :meth:`connect`."""
        self._uri = uri
        self._connect_kwargs = dict(kwargs)
        self._connect_kwargs["close_timeout"] = _positive_timeout(kwargs.get("close_timeout", 10.0), "close_timeout")
        self._ws: ClientConnection | None = None
        self._loop: asyncio.AbstractEventLoop | None = None
        self.request_timeout = _positive_timeout(request_timeout, "request_timeout")
        if type(notification_queue_size) is not int or notification_queue_size <= 0:
            raise ValueError("notification_queue_size must be a positive integer")
        self._notification_queue_size = notification_queue_size
        self._overflow = OverflowPolicy(overflow)
        self._dropped_notifications = 0
        self._notifications: deque[Notification] = deque()
        self._notification_ready = asyncio.Event()
        self._receiving = False
        self._connect_lock = asyncio.Lock()
        self._pending_requests: dict[int, _PendingRequest[Any]] = {}
        self._request_counter = itertools.count(1)
        self._subscriptions: dict[int, Subscription] = {}
        self._unsubscribing: set[int] = set()
        self._send_lock = asyncio.Lock()
        self._closed_exc: BaseException | None = None
        self._reader_task: asyncio.Task[None] | None = None

    async def connect(self) -> SolanaWsClient:
        """Open and own the native WebSocket, then start its sole reader."""
        async with self._connect_lock:
            if self.connection_state is ConnectionState.OPEN:
                return self
            # A failed handshake owns nothing, so only a used client is refused.
            if self._ws is not None or self._closed_exc is not None:
                raise RuntimeError("SolanaWsClient instances cannot be reused")
            # Pinned before the first loop-bound resource exists, so every task and future shares one loop.
            loop = self._loop = asyncio.get_running_loop()
            self._ws = await ws_connect(self._uri, **self._connect_kwargs)
            self._reader_task = loop.create_task(self._read_loop(), name="solana-ws-reader")
            return self

    async def __aenter__(self) -> SolanaWsClient:
        """Connect this client for use as an async context manager."""
        return await self.connect()

    @property
    def connection_state(self) -> ConnectionState:
        """Return the RPC dispatcher's lifecycle state."""
        if self._ws is None or self._closed_exc is not None:
            return ConnectionState.CLOSED
        return ConnectionState.OPEN

    @property
    def dropped_notifications(self) -> int:
        """Count notifications a lossy overflow policy discarded; always 0 under ``RAISE``."""
        return self._dropped_notifications

    def _abandon(self, exc: BaseException) -> None:
        """Record why the stream ended and wake every waiter; the first cause wins.

        A reader task cannot raise into the caller's stack, so the cause is
        stored here and re-raised by whoever is waiting.
        """
        if self._closed_exc is not None:
            return
        self._closed_exc = exc
        self._notification_ready.set()
        for pending in self._pending_requests.values():
            if not pending.future.done():
                pending.future.set_exception(exc)
        # Subscriptions are owned by one physical connection and die with it.
        self._subscriptions.clear()

    async def close(self, code: int = CloseCode.NORMAL_CLOSURE, reason: str = "") -> None:
        """Release the connection; the only cleanup path, idempotent and safe to call concurrently."""
        frame = Close(code, reason)
        frame.check()
        self._abandon(_local_close_exc(frame))
        loop = self._loop
        if self._ws is None or loop is None:
            return
        # Shielded so a cancelled caller -- a repeated Ctrl+C -- cannot abandon a half-closed socket.
        release = loop.create_task(self._release(frame), name="solana-ws-release")
        try:
            await asyncio.shield(release)
        except asyncio.CancelledError:
            await asyncio.shield(release)
            raise

    async def _release(self, frame: Close) -> None:
        ws = self._ws
        try:
            if ws is not None:
                # Idempotent, bounded by close_timeout, and aborts the transport if that elapses.
                await ws.close(frame.code, frame.reason)
        finally:
            reader = self._reader_task
            if reader is not None:
                reader.cancel()
                await asyncio.gather(reader, return_exceptions=True)

    async def __aexit__(
        self,
        exc_type: type[BaseException] | None,
        exc_value: BaseException | None,
        traceback: TracebackType | None,
    ) -> None:
        """Report the context's own failure to waiters, then always release the socket."""
        if exc_value is not None:
            self._abandon(exc_value)
        await self.close()

    async def recv(self) -> Notification:
        """Receive one notification; cancellation leaves queued notifications intact.

        Only one caller may receive at a time. Notifications already received
        are delivered before the closure is reported, as ``websockets``' own
        ``recv()`` does. Raw text/bytes decoding isn't supported.
        """
        if self._receiving:
            raise ConcurrencyError("Only one notification receiver may run at a time")
        self._receiving = True
        try:
            while True:
                if self._notifications:
                    return self._notifications.popleft()
                if self._closed_exc is not None:
                    raise self._closed_exc
                if self._ws is None:
                    raise RuntimeError("WebSocket is not connected")
                self._notification_ready.clear()
                await self._notification_ready.wait()
        finally:
            self._receiving = False

    async def __aiter__(self) -> AsyncIterator[Notification]:
        """Iterate over notifications until a normal connection closure."""
        try:
            while True:
                yield await self.recv()
        except ConnectionClosedOK:
            return

    async def _read_loop(self) -> None:
        ws = self._ws
        if ws is None:
            return
        try:
            while self._closed_exc is None:
                raw = await ws.recv()
                for envelope in _parse_frame(raw.decode() if isinstance(raw, bytes) else raw):
                    if self._closed_exc is not None:
                        return
                    if isinstance(envelope, JsonRpcResponseEnvelope):
                        self._dispatch_error(envelope)
                    elif isinstance(envelope, (SubscriptionResult, UnsubscribeResult)):
                        self._dispatch_response(envelope)
                    else:
                        self._dispatch_notification(cast(Notification, envelope))
        # The reader records every failure of this connection.
        except Exception as exc:  # noqa: BLE001
            self._abandon(exc)

    def _dispatch_error(self, envelope: JsonRpcResponseEnvelope) -> None:
        pending = self._pending_requests.pop(cast(int, envelope.id), None)
        if pending is None:
            return
        pending.future.set_exception(
            SolanaJsonRpcError.from_error_object(
                cast(JsonRpcErrorObject, envelope.error),
                request_id=envelope.id,
                method=pending.method,
            )
        )

    def _dispatch_response(self, envelope: SubscriptionResult | UnsubscribeResult) -> None:
        request_id = envelope.id
        # Popping first makes dispatch idempotent: a response for a request that
        # already timed out, was cancelled, or arrives twice has no waiter left.
        pending = self._pending_requests.pop(request_id, None)
        if pending is None:
            return
        if pending.kind is not None:
            # Registered before the caller is woken: the reader keeps draining
            # frames while that coroutine is merely scheduled, so a notification
            # following the confirmation must already find the handle.
            pending.future.set_result(self._register_subscription(pending.kind, cast(int, envelope.result)))
        else:
            # No kind means an unsubscribe, whose result is the server's boolean.
            pending.future.set_result(envelope.result)

    def _dispatch_notification(self, notification: Notification) -> None:
        subscription_id = notification.subscription
        if isinstance(notification, SignatureNotification):
            self._subscriptions.pop(subscription_id, None)
        if len(self._notifications) >= self._notification_queue_size:
            if self._overflow is OverflowPolicy.RAISE:
                raise ProtocolError("WebSocket notification queue overflow")
            self._dropped_notifications += 1
            if self._overflow is OverflowPolicy.DROP_NEWEST:
                return
            self._notifications.popleft()
        self._notifications.append(notification)
        self._notification_ready.set()

    async def _request(
        self,
        request: JsonRpcRequestSerializer,
        kind: SubscriptionKind | None = None,
    ) -> _PendingRequest[T]:
        ws, loop = self._ws, self._loop
        if self.connection_state is not ConnectionState.OPEN or ws is None or loop is None:
            raise RuntimeError("WebSocket is not connected")
        serialized = request.to_json()
        body = json.loads(serialized)
        # solders request serializers produce valid JSON-RPC envelopes.
        # The protocol serializer exposes ``id`` at runtime; the shared
        # serializer type omits that concrete request attribute.
        request_id = cast(Any, request).id
        future: asyncio.Future[T] = cast(asyncio.Future[T], loop.create_future())
        future.add_done_callback(_consume_future_exception)
        pending = _PendingRequest(request_id, future, body["method"], kind)
        self._pending_requests[request_id] = pending
        try:
            async with self._send_lock:
                if self._closed_exc is not None:
                    raise self._closed_exc
                pending.send_started = True
                try:
                    await ws.send(serialized)
                except Exception as exc:
                    self._abandon(exc)
                    raise
        except BaseException:
            self._pending_requests.pop(request_id, None)
            if not future.done():
                future.cancel()
            raise
        return pending

    async def _wait_pending(self, pending: _PendingRequest[T], timeout: float | None = None) -> T:
        """Await a registered request and remove it on caller cancellation/timeout."""
        duration = self.request_timeout if timeout is None else _positive_timeout(timeout, "timeout")
        timer = asyncio.timeout(duration)
        try:
            async with timer:
                return await asyncio.shield(pending.future)
        except asyncio.CancelledError as exc:
            if pending.send_started:
                self._abandon(exc)
            raise
        except TimeoutError as exc:
            if pending.send_started and timer.expired():
                self._abandon(exc)
            raise
        finally:
            self._pending_requests.pop(pending.request_id, None)
            if not pending.future.done():
                pending.future.cancel()

    async def _subscribe(
        self,
        kind: SubscriptionKind,
        request: JsonRpcRequestSerializer,
    ) -> Subscription:
        pending: _PendingRequest[Subscription] = await self._request(request, kind)
        return await self._wait_pending(pending)

    def _register_subscription(self, kind: SubscriptionKind, subscription_id: int) -> Subscription:
        subscription = Subscription(subscription_id, kind)
        self._subscriptions[subscription_id] = subscription
        return subscription

    def _remove_subscription(self, subscription_id: int) -> None:
        self._subscriptions.pop(subscription_id, None)

    def _unsubscribe_request(self, subscription: Subscription) -> JsonRpcRequestSerializer:
        """Build the JSON-RPC request that cancels a subscription."""
        request_type = _UNSUBSCRIBE_REQUESTS[subscription.kind]
        return request_type(subscription.subscription_id, next(self._request_counter))

    async def unsubscribe(self, subscription: Subscription) -> None:
        """Cancel a live subscription after receiving server confirmation.

        Closing the connection already cancels everything it owned, so this is a
        no-op afterwards and never masks the error that ended the stream.
        """
        if self._closed_exc is not None:
            return
        subscription_id = subscription.subscription_id
        if self._subscriptions.get(subscription_id) is not subscription:
            raise ValueError("Subscription handle is no longer active")
        if subscription_id in self._unsubscribing:
            raise ValueError("Unsubscribe is already pending for this handle")
        request = self._unsubscribe_request(subscription)
        self._unsubscribing.add(subscription_id)
        try:
            pending: _PendingRequest[bool] = await self._request(request)
            result = await self._wait_pending(pending)
            if not result:
                raise UnsubscribeError(subscription)
            self._remove_subscription(subscription_id)
        finally:
            self._unsubscribing.discard(subscription_id)

    async def account_subscribe(
        self,
        *,
        pubkey: Pubkey,
        commitment: Commitment | None = None,
        encoding: str | None = None,
        data_slice: DataSliceOpts | None = None,
        min_context_slot: int | None = None,
    ) -> Subscription:
        """Subscribe to account notifications for a public key."""
        config = None
        if any(value is not None for value in (commitment, encoding, data_slice, min_context_slot)):
            account_encoding = _ACCOUNT_ENCODING_TO_SOLDERS[encoding] if encoding is not None else None
            account_commitment = _COMMITMENT_TO_SOLDERS[commitment] if commitment is not None else None
            account_data_slice = (
                UiDataSliceConfig(offset=data_slice.offset, length=data_slice.length) if data_slice else None
            )
            config = RpcAccountInfoConfig(
                account_encoding,
                account_data_slice,
                account_commitment,
                min_context_slot,
            )
        return await self._subscribe(
            SubscriptionKind.ACCOUNT,
            AccountSubscribe(pubkey, config, next(self._request_counter)),
        )

    async def program_subscribe(
        self,
        *,
        program_id: Pubkey,
        commitment: Commitment | None = None,
        encoding: str | None = None,
        data_slice: DataSliceOpts | None = None,
        min_context_slot: int | None = None,
        filters: Sequence[int | MemcmpOpts] | None = None,
        with_context: bool | None = None,
        sort_results: bool | None = None,
    ) -> Subscription:
        """Subscribe to program account notifications."""
        config = None
        if any(
            value is not None
            for value in (
                commitment,
                encoding,
                data_slice,
                min_context_slot,
                filters,
                with_context,
                sort_results,
            )
        ):
            account = RpcAccountInfoConfig(
                encoding=(None if encoding is None else _ACCOUNT_ENCODING_TO_SOLDERS[encoding]),
                commitment=(None if commitment is None else _COMMITMENT_TO_SOLDERS[commitment]),
                min_context_slot=min_context_slot,
                data_slice=(
                    None
                    if data_slice is None
                    else UiDataSliceConfig(offset=data_slice.offset, length=data_slice.length)
                ),
            )
            parsed_filters = (
                None
                if filters is None
                else [x if isinstance(x, int) else Memcmp(offset=x.offset, bytes_=x.bytes) for x in filters]
            )
            config = cast(Any, RpcProgramAccountsConfig)(account, parsed_filters, with_context, sort_results)
        return await self._subscribe(
            SubscriptionKind.PROGRAM,
            ProgramSubscribe(program_id, config, next(self._request_counter)),
        )

    async def logs_subscribe(
        self,
        *,
        filter_: (RpcTransactionLogsFilter | RpcTransactionLogsFilterMentions) = RpcTransactionLogsFilter.All,
        commitment: Commitment | None = None,
    ) -> Subscription:
        """Subscribe to transaction log notifications."""
        logs_commitment = _COMMITMENT_TO_SOLDERS[commitment] if commitment is not None else None
        config = RpcTransactionLogsConfig(logs_commitment)
        return await self._subscribe(
            SubscriptionKind.LOGS,
            LogsSubscribe(filter_, config, next(self._request_counter)),
        )

    async def block_subscribe(
        self,
        *,
        filter_: (RpcBlockSubscribeFilter | RpcBlockSubscribeFilterMentions) = RpcBlockSubscribeFilter.All,
        commitment: Commitment | None = None,
        encoding: str | None = None,
        transaction_details: TransactionDetails | None = None,
        show_rewards: bool | None = None,
        max_supported_transaction_version: int | None = None,
    ) -> Subscription:
        """Subscribe to block notifications."""
        block_commitment = _COMMITMENT_TO_SOLDERS[commitment] if commitment is not None else None
        block_encoding = _TX_ENCODING_TO_SOLDERS[encoding] if encoding is not None else None
        config = RpcBlockSubscribeConfig(
            block_commitment,
            block_encoding,
            transaction_details,
            show_rewards,
            max_supported_transaction_version,
        )
        return await self._subscribe(
            SubscriptionKind.BLOCK,
            BlockSubscribe(filter_, config, next(self._request_counter)),
        )

    async def signature_subscribe(
        self,
        *,
        signature: Signature,
        commitment: Commitment | None = None,
        enable_received_notification: bool | None = None,
    ) -> Subscription:
        """Subscribe to signature status notifications."""
        config = None
        if commitment is not None or enable_received_notification is not None:
            signature_commitment = _COMMITMENT_TO_SOLDERS[commitment] if commitment is not None else None
            config = RpcSignatureSubscribeConfig(signature_commitment, enable_received_notification)
        return await self._subscribe(
            SubscriptionKind.SIGNATURE,
            SignatureSubscribe(signature, config, next(self._request_counter)),
        )

    async def slot_subscribe(self) -> Subscription:
        """Subscribe to slot notifications."""
        return await self._subscribe(SubscriptionKind.SLOT, SlotSubscribe(next(self._request_counter)))

    async def slots_updates_subscribe(self) -> Subscription:
        """Subscribe to slot update notifications."""
        return await self._subscribe(
            SubscriptionKind.SLOTS_UPDATES,
            SlotsUpdatesSubscribe(next(self._request_counter)),
        )

    async def root_subscribe(self) -> Subscription:
        """Subscribe to root notifications."""
        return await self._subscribe(SubscriptionKind.ROOT, RootSubscribe(next(self._request_counter)))

    async def vote_subscribe(self) -> Subscription:
        """Subscribe to vote notifications."""
        return await self._subscribe(SubscriptionKind.VOTE, VoteSubscribe(next(self._request_counter)))

connection_state property

Return the RPC dispatcher's lifecycle state.

dropped_notifications property

Count notifications a lossy overflow policy discarded; always 0 under RAISE.

__init__(uri='ws://localhost:8900', *, request_timeout=10.0, notification_queue_size=10000, overflow=OverflowPolicy.RAISE, **kwargs)

Create a client; the WebSocket is opened by :meth:connect.

Source code in src/solana/rpc/websocket_api.py
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
def __init__(
    self,
    uri: str = "ws://localhost:8900",
    *,
    request_timeout: float = 10.0,
    notification_queue_size: int = 10_000,
    overflow: OverflowPolicy = OverflowPolicy.RAISE,
    **kwargs: Any,
) -> None:
    """Create a client; the WebSocket is opened by :meth:`connect`."""
    self._uri = uri
    self._connect_kwargs = dict(kwargs)
    self._connect_kwargs["close_timeout"] = _positive_timeout(kwargs.get("close_timeout", 10.0), "close_timeout")
    self._ws: ClientConnection | None = None
    self._loop: asyncio.AbstractEventLoop | None = None
    self.request_timeout = _positive_timeout(request_timeout, "request_timeout")
    if type(notification_queue_size) is not int or notification_queue_size <= 0:
        raise ValueError("notification_queue_size must be a positive integer")
    self._notification_queue_size = notification_queue_size
    self._overflow = OverflowPolicy(overflow)
    self._dropped_notifications = 0
    self._notifications: deque[Notification] = deque()
    self._notification_ready = asyncio.Event()
    self._receiving = False
    self._connect_lock = asyncio.Lock()
    self._pending_requests: dict[int, _PendingRequest[Any]] = {}
    self._request_counter = itertools.count(1)
    self._subscriptions: dict[int, Subscription] = {}
    self._unsubscribing: set[int] = set()
    self._send_lock = asyncio.Lock()
    self._closed_exc: BaseException | None = None
    self._reader_task: asyncio.Task[None] | None = None

account_subscribe(*, pubkey, commitment=None, encoding=None, data_slice=None, min_context_slot=None) async

Subscribe to account notifications for a public key.

Source code in src/solana/rpc/websocket_api.py
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
async def account_subscribe(
    self,
    *,
    pubkey: Pubkey,
    commitment: Commitment | None = None,
    encoding: str | None = None,
    data_slice: DataSliceOpts | None = None,
    min_context_slot: int | None = None,
) -> Subscription:
    """Subscribe to account notifications for a public key."""
    config = None
    if any(value is not None for value in (commitment, encoding, data_slice, min_context_slot)):
        account_encoding = _ACCOUNT_ENCODING_TO_SOLDERS[encoding] if encoding is not None else None
        account_commitment = _COMMITMENT_TO_SOLDERS[commitment] if commitment is not None else None
        account_data_slice = (
            UiDataSliceConfig(offset=data_slice.offset, length=data_slice.length) if data_slice else None
        )
        config = RpcAccountInfoConfig(
            account_encoding,
            account_data_slice,
            account_commitment,
            min_context_slot,
        )
    return await self._subscribe(
        SubscriptionKind.ACCOUNT,
        AccountSubscribe(pubkey, config, next(self._request_counter)),
    )

block_subscribe(*, filter_=RpcBlockSubscribeFilter.All, commitment=None, encoding=None, transaction_details=None, show_rewards=None, max_supported_transaction_version=None) async

Subscribe to block notifications.

Source code in src/solana/rpc/websocket_api.py
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
async def block_subscribe(
    self,
    *,
    filter_: (RpcBlockSubscribeFilter | RpcBlockSubscribeFilterMentions) = RpcBlockSubscribeFilter.All,
    commitment: Commitment | None = None,
    encoding: str | None = None,
    transaction_details: TransactionDetails | None = None,
    show_rewards: bool | None = None,
    max_supported_transaction_version: int | None = None,
) -> Subscription:
    """Subscribe to block notifications."""
    block_commitment = _COMMITMENT_TO_SOLDERS[commitment] if commitment is not None else None
    block_encoding = _TX_ENCODING_TO_SOLDERS[encoding] if encoding is not None else None
    config = RpcBlockSubscribeConfig(
        block_commitment,
        block_encoding,
        transaction_details,
        show_rewards,
        max_supported_transaction_version,
    )
    return await self._subscribe(
        SubscriptionKind.BLOCK,
        BlockSubscribe(filter_, config, next(self._request_counter)),
    )

close(code=CloseCode.NORMAL_CLOSURE, reason='') async

Release the connection; the only cleanup path, idempotent and safe to call concurrently.

Source code in src/solana/rpc/websocket_api.py
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
async def close(self, code: int = CloseCode.NORMAL_CLOSURE, reason: str = "") -> None:
    """Release the connection; the only cleanup path, idempotent and safe to call concurrently."""
    frame = Close(code, reason)
    frame.check()
    self._abandon(_local_close_exc(frame))
    loop = self._loop
    if self._ws is None or loop is None:
        return
    # Shielded so a cancelled caller -- a repeated Ctrl+C -- cannot abandon a half-closed socket.
    release = loop.create_task(self._release(frame), name="solana-ws-release")
    try:
        await asyncio.shield(release)
    except asyncio.CancelledError:
        await asyncio.shield(release)
        raise

connect() async

Open and own the native WebSocket, then start its sole reader.

Source code in src/solana/rpc/websocket_api.py
260
261
262
263
264
265
266
267
268
269
270
271
272
async def connect(self) -> SolanaWsClient:
    """Open and own the native WebSocket, then start its sole reader."""
    async with self._connect_lock:
        if self.connection_state is ConnectionState.OPEN:
            return self
        # A failed handshake owns nothing, so only a used client is refused.
        if self._ws is not None or self._closed_exc is not None:
            raise RuntimeError("SolanaWsClient instances cannot be reused")
        # Pinned before the first loop-bound resource exists, so every task and future shares one loop.
        loop = self._loop = asyncio.get_running_loop()
        self._ws = await ws_connect(self._uri, **self._connect_kwargs)
        self._reader_task = loop.create_task(self._read_loop(), name="solana-ws-reader")
        return self

logs_subscribe(*, filter_=RpcTransactionLogsFilter.All, commitment=None) async

Subscribe to transaction log notifications.

Source code in src/solana/rpc/websocket_api.py
613
614
615
616
617
618
619
620
621
622
623
624
625
async def logs_subscribe(
    self,
    *,
    filter_: (RpcTransactionLogsFilter | RpcTransactionLogsFilterMentions) = RpcTransactionLogsFilter.All,
    commitment: Commitment | None = None,
) -> Subscription:
    """Subscribe to transaction log notifications."""
    logs_commitment = _COMMITMENT_TO_SOLDERS[commitment] if commitment is not None else None
    config = RpcTransactionLogsConfig(logs_commitment)
    return await self._subscribe(
        SubscriptionKind.LOGS,
        LogsSubscribe(filter_, config, next(self._request_counter)),
    )

program_subscribe(*, program_id, commitment=None, encoding=None, data_slice=None, min_context_slot=None, filters=None, with_context=None, sort_results=None) async

Subscribe to program account notifications.

Source code in src/solana/rpc/websocket_api.py
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
async def program_subscribe(
    self,
    *,
    program_id: Pubkey,
    commitment: Commitment | None = None,
    encoding: str | None = None,
    data_slice: DataSliceOpts | None = None,
    min_context_slot: int | None = None,
    filters: Sequence[int | MemcmpOpts] | None = None,
    with_context: bool | None = None,
    sort_results: bool | None = None,
) -> Subscription:
    """Subscribe to program account notifications."""
    config = None
    if any(
        value is not None
        for value in (
            commitment,
            encoding,
            data_slice,
            min_context_slot,
            filters,
            with_context,
            sort_results,
        )
    ):
        account = RpcAccountInfoConfig(
            encoding=(None if encoding is None else _ACCOUNT_ENCODING_TO_SOLDERS[encoding]),
            commitment=(None if commitment is None else _COMMITMENT_TO_SOLDERS[commitment]),
            min_context_slot=min_context_slot,
            data_slice=(
                None
                if data_slice is None
                else UiDataSliceConfig(offset=data_slice.offset, length=data_slice.length)
            ),
        )
        parsed_filters = (
            None
            if filters is None
            else [x if isinstance(x, int) else Memcmp(offset=x.offset, bytes_=x.bytes) for x in filters]
        )
        config = cast(Any, RpcProgramAccountsConfig)(account, parsed_filters, with_context, sort_results)
    return await self._subscribe(
        SubscriptionKind.PROGRAM,
        ProgramSubscribe(program_id, config, next(self._request_counter)),
    )

recv() async

Receive one notification; cancellation leaves queued notifications intact.

Only one caller may receive at a time. Notifications already received are delivered before the closure is reported, as websockets' own recv() does. Raw text/bytes decoding isn't supported.

Source code in src/solana/rpc/websocket_api.py
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
async def recv(self) -> Notification:
    """Receive one notification; cancellation leaves queued notifications intact.

    Only one caller may receive at a time. Notifications already received
    are delivered before the closure is reported, as ``websockets``' own
    ``recv()`` does. Raw text/bytes decoding isn't supported.
    """
    if self._receiving:
        raise ConcurrencyError("Only one notification receiver may run at a time")
    self._receiving = True
    try:
        while True:
            if self._notifications:
                return self._notifications.popleft()
            if self._closed_exc is not None:
                raise self._closed_exc
            if self._ws is None:
                raise RuntimeError("WebSocket is not connected")
            self._notification_ready.clear()
            await self._notification_ready.wait()
    finally:
        self._receiving = False

root_subscribe() async

Subscribe to root notifications.

Source code in src/solana/rpc/websocket_api.py
680
681
682
async def root_subscribe(self) -> Subscription:
    """Subscribe to root notifications."""
    return await self._subscribe(SubscriptionKind.ROOT, RootSubscribe(next(self._request_counter)))

signature_subscribe(*, signature, commitment=None, enable_received_notification=None) async

Subscribe to signature status notifications.

Source code in src/solana/rpc/websocket_api.py
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
async def signature_subscribe(
    self,
    *,
    signature: Signature,
    commitment: Commitment | None = None,
    enable_received_notification: bool | None = None,
) -> Subscription:
    """Subscribe to signature status notifications."""
    config = None
    if commitment is not None or enable_received_notification is not None:
        signature_commitment = _COMMITMENT_TO_SOLDERS[commitment] if commitment is not None else None
        config = RpcSignatureSubscribeConfig(signature_commitment, enable_received_notification)
    return await self._subscribe(
        SubscriptionKind.SIGNATURE,
        SignatureSubscribe(signature, config, next(self._request_counter)),
    )

slot_subscribe() async

Subscribe to slot notifications.

Source code in src/solana/rpc/websocket_api.py
669
670
671
async def slot_subscribe(self) -> Subscription:
    """Subscribe to slot notifications."""
    return await self._subscribe(SubscriptionKind.SLOT, SlotSubscribe(next(self._request_counter)))

slots_updates_subscribe() async

Subscribe to slot update notifications.

Source code in src/solana/rpc/websocket_api.py
673
674
675
676
677
678
async def slots_updates_subscribe(self) -> Subscription:
    """Subscribe to slot update notifications."""
    return await self._subscribe(
        SubscriptionKind.SLOTS_UPDATES,
        SlotsUpdatesSubscribe(next(self._request_counter)),
    )

unsubscribe(subscription) async

Cancel a live subscription after receiving server confirmation.

Closing the connection already cancels everything it owned, so this is a no-op afterwards and never masks the error that ended the stream.

Source code in src/solana/rpc/websocket_api.py
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
async def unsubscribe(self, subscription: Subscription) -> None:
    """Cancel a live subscription after receiving server confirmation.

    Closing the connection already cancels everything it owned, so this is a
    no-op afterwards and never masks the error that ended the stream.
    """
    if self._closed_exc is not None:
        return
    subscription_id = subscription.subscription_id
    if self._subscriptions.get(subscription_id) is not subscription:
        raise ValueError("Subscription handle is no longer active")
    if subscription_id in self._unsubscribing:
        raise ValueError("Unsubscribe is already pending for this handle")
    request = self._unsubscribe_request(subscription)
    self._unsubscribing.add(subscription_id)
    try:
        pending: _PendingRequest[bool] = await self._request(request)
        result = await self._wait_pending(pending)
        if not result:
            raise UnsubscribeError(subscription)
        self._remove_subscription(subscription_id)
    finally:
        self._unsubscribing.discard(subscription_id)

vote_subscribe() async

Subscribe to vote notifications.

Source code in src/solana/rpc/websocket_api.py
684
685
686
async def vote_subscribe(self) -> Subscription:
    """Subscribe to vote notifications."""
    return await self._subscribe(SubscriptionKind.VOTE, VoteSubscribe(next(self._request_counter)))

Subscription dataclass

A server-confirmed subscription owned by one physical connection.

Handles are created by subscribe helpers. Copying or reconstructing a handle doesn't grant ownership; unsubscribe checks its object identity.

Source code in src/solana/rpc/websocket_api.py
115
116
117
118
119
120
121
122
123
124
@dataclass(frozen=True, slots=True)
class Subscription:
    """A server-confirmed subscription owned by one physical connection.

    Handles are created by subscribe helpers. Copying or reconstructing a handle
    doesn't grant ownership; unsubscribe checks its object identity.
    """

    subscription_id: int
    kind: SubscriptionKind

SubscriptionKind

Solana subscription methods supported by the typed helpers.

Source code in src/solana/rpc/websocket_api.py
88
89
90
91
92
93
94
95
96
97
98
99
class SubscriptionKind(StrEnum):
    """Solana subscription methods supported by the typed helpers."""

    ACCOUNT = "account"
    BLOCK = "block"
    LOGS = "logs"
    PROGRAM = "program"
    SIGNATURE = "signature"
    SLOT = "slot"
    SLOTS_UPDATES = "slotsUpdates"
    ROOT = "root"
    VOTE = "vote"

UnsubscribeError

The server explicitly refused to cancel a subscription.

Source code in src/solana/rpc/websocket_api.py
149
150
151
152
153
154
155
class UnsubscribeError(Exception):
    """The server explicitly refused to cancel a subscription."""

    def __init__(self, subscription: Subscription) -> None:
        """Retain the handle whose unsubscribe request returned false."""
        self.subscription = subscription
        super().__init__(f"Unsubscribe returned false for {subscription.kind.value} {subscription.subscription_id}")

__init__(subscription)

Retain the handle whose unsubscribe request returned false.

Source code in src/solana/rpc/websocket_api.py
152
153
154
155
def __init__(self, subscription: Subscription) -> None:
    """Retain the handle whose unsubscribe request returned false."""
    self.subscription = subscription
    super().__init__(f"Unsubscribe returned false for {subscription.kind.value} {subscription.subscription_id}")