> ## Documentation Index
> Fetch the complete documentation index at: https://docs.metastreams.dev/llms.txt
> Use this file to discover all available pages before exploring further.

# Connecting to streams

> Open a WebSocket, subscribe to channels, and reconnect after a close.

The stream pushes trades, candles and token updates as they happen. One WebSocket connection carries up to 100 subscriptions across every channel.

## Connect

Open a WebSocket to `wss://{{API_HOST}}/v1/stream`, with your key in the `Authorization` header of the handshake.

<CodeGroup>
  ```bash wscat theme={null}
  wscat -c "wss://{{API_HOST}}/v1/stream" -H "Authorization: Bearer $API_KEY"
  ```

  ```typescript Node.js (ws) theme={null}
  import WebSocket from "ws";

  const socket = new WebSocket("wss://{{API_HOST}}/v1/stream", {
    headers: { Authorization: `Bearer ${process.env.API_KEY}` },
  });
  // A refused handshake arrives as an `error` event. Without a listener, Node exits.
  socket.on("error", (error) => console.error("stream error:", error.message));
  ```

  ```python Python (websockets 13+) theme={null}
  import asyncio
  import os

  from websockets.asyncio.client import connect

  async def main():
      async with connect(
          "wss://{{API_HOST}}/v1/stream",
          additional_headers={"Authorization": f"Bearer {os.environ['API_KEY']}"},
      ) as socket:
          ...

  asyncio.run(main())
  ```
</CodeGroup>

The server checks the handshake before it upgrades the connection:

| Response | Why |
| - | - |
| `401 UNAUTHORIZED` | The key is missing or invalid. Do not retry. |
| `429 RATE_LIMITED` | Your key already holds 100 open connections, or a rate limit applies. |
| `503 SERVICE_UNAVAILABLE` | The key could not be checked. Retry with backoff. |

A request without WebSocket upgrade headers gets a plain-text `4xx` instead of an upgrade.

Browsers cannot set the `Authorization` header on a WebSocket, so connect from your backend.

## Frames

Every frame is a JSON text message.

### From you

```jsonc theme={null}
// Open a subscription. `id` is optional; the ack, or an error about this subscribe, echoes it.
{ "action": "subscribe", "channel": "spot.trades", "filters": { "chain": "solana", "token": "DezXAZ8z7PnrnRJjz3wXBoRgixCa6xjnB7YaB1pPB263" }, "id": "r1" }

// Close one subscription.
{ "action": "unsubscribe", "channel": "spot.trades", "subscriptionId": "5f0c6d1e-8a2b-4c3d-9e4f-1a2b3c4d5e6f" }

// Close every subscription on a channel.
{ "action": "unsubscribe", "channel": "spot.trades" }

// Check the connection is alive.
{ "action": "ping" }
```

### From the server

```jsonc theme={null}
// A subscribe or unsubscribe took effect.
{ "type": "ack", "channel": "spot.trades", "subscriptionId": "5f0c6d1e-8a2b-4c3d-9e4f-1a2b3c4d5e6f", "id": "r1", "status": "success", "message": "subscribe successful" }

// One payload for one subscription.
{ "type": "update", "channel": "spot.trades", "subscriptionId": "5f0c6d1e-8a2b-4c3d-9e4f-1a2b3c4d5e6f", "data": { } }

// A frame was refused.
{ "type": "error", "id": "r1", "code": "VALIDATION_ERROR", "message": "One or more parameters are invalid.", "details": [ { "field": "chain", "message": "Name the chain." } ] }

// The answer to a ping.
{ "type": "pong" }
```

### Subscription IDs

The server assigns each subscription a `subscriptionId` and returns it on the subscribe ack. Every update carries the ID of the subscription it belongs to. Use the ID to route updates in your client and to unsubscribe exactly.

An unsubscribe without a `subscriptionId` closes every subscription you hold on that channel, and its ack carries no ID.

### Ordering

* A subscription's ack always arrives before its first update.
* An unsubscribe's ack always arrives after that subscription's last update.

## Channels

| Channel | Filter | Each update's `data` |
| - | - | - |
| [`spot.trades`](/reference/streams/spot-trades) | One of `chain` + `token`, `chain` + `wallet`, or `identity`. Optional `flags`. | One trade |
| [`spot.candles`](/reference/streams/spot-candles) | `chain` + `token`. Optional `timeframe`, `metric` and `denomination`. | The current state of one candle |
| [`spot.tokens`](/reference/streams/spot-tokens) | `chain` + 1 to 100 `tokens`. | The whole token object |

Each channel's reference page lists every filter, every update field and every error.

## Errors

A refused frame gets an `error` frame. It uses the same codes as the REST API. An error about a subscribe frame echoes that frame's `id`. A frame the server cannot parse, and a refused unsubscribe, get an error with no `id`.

| Code | When |
| - | - |
| `VALIDATION_ERROR` | The frame is not valid JSON, has no known `action`, or has a missing or mistyped field; the filters do not match the channel; or the connection already holds 100 subscriptions |
| `UNSUPPORTED_CHAIN` | The filter's `chain` is not [supported](/concepts/chains) |
| `NOT_FOUND` | The channel does not exist, or an unsubscribe names a subscription this connection does not hold |

An error frame never closes the connection.

## Keeping the connection alive

Send `{"action": "ping"}` every 30 seconds or so. If no `pong` arrives within a few seconds, treat the connection as dead and reconnect.

## When the server closes the connection

| Close code | Reason | What to do |
| - | - | - |
| `4008` | `slow consumer`: your client did not read fast enough. 1,024 frames were waiting unsent, or one write stalled for 10 seconds. | Reconnect and resubscribe. Read frames faster, or subscribe to less. |
| `1001` | `server restarting`: the server is shutting down, for example during a deploy. | Reconnect and resubscribe straight away. |
| Any other | The network dropped, or the connection failed. | Reconnect with backoff and resubscribe. |

Subscriptions do not survive a reconnect. After you reconnect, send your subscribe frames again. To fill the gap while you were away, read recent trades or candles from the REST API.

### Reconnect loop

These loops resubscribe after every close, reconnect straight away after `1001`, back off after anything else, and stop only on `401`.

<CodeGroup>
  ```typescript Node.js (ws) theme={null}
  import WebSocket from "ws";

  const URL = "wss://{{API_HOST}}/v1/stream";
  const subscriptions = [
    { action: "subscribe", channel: "spot.trades", filters: { chain: "solana", token: "DezXAZ8z7PnrnRJjz3wXBoRgixCa6xjnB7YaB1pPB263" } },
  ];

  function connect(attempt = 0) {
    const socket = new WebSocket(URL, { headers: { Authorization: `Bearer ${process.env.API_KEY}` } });
    let lastPong = Date.now();
    let heartbeat: NodeJS.Timeout | undefined;
    let unauthorized = false;

    // A refused handshake arrives here too. Without an `error` listener, Node exits.
    socket.on("error", (error) => {
      unauthorized = error.message.includes("Unexpected server response: 401");
      console.error("stream error:", error.message);
    });

    socket.on("open", () => {
      attempt = 0;
      for (const frame of subscriptions) socket.send(JSON.stringify(frame));
      heartbeat = setInterval(() => {
        if (Date.now() - lastPong > 45_000) return socket.terminate();
        socket.send(JSON.stringify({ action: "ping" }));
      }, 30_000);
    });

    socket.on("message", (raw) => {
      const frame = JSON.parse(raw.toString());
      if (frame.type === "pong") lastPong = Date.now();
      if (frame.type === "update") handleUpdate(frame);
      if (frame.type === "error") console.error(frame.code, frame.message);
    });

    socket.on("close", (code) => {
      clearInterval(heartbeat);
      if (unauthorized) return console.error("The API key was refused. Not reconnecting.");
      const delay = code === 1001 ? 0 : Math.min(30_000, 1000 * 2 ** attempt) + Math.random() * 1000;
      setTimeout(() => connect(attempt + 1), delay);
    });
  }

  function handleUpdate(frame: { channel: string; data: unknown }) {
    console.log(frame.channel, frame.data);
  }

  connect();
  ```

  ```python Python (websockets 13+) theme={null}
  import asyncio
  import json
  import os
  import random

  from websockets.asyncio.client import connect
  from websockets.exceptions import ConnectionClosed, InvalidStatus

  URL = "wss://{{API_HOST}}/v1/stream"
  HEADERS = {"Authorization": f"Bearer {os.environ['API_KEY']}"}
  SUBSCRIPTIONS = [
      {"action": "subscribe", "channel": "spot.trades", "filters": {"chain": "solana", "token": "DezXAZ8z7PnrnRJjz3wXBoRgixCa6xjnB7YaB1pPB263"}},
  ]

  async def run():
      attempt = 0
      while True:
          close_code = None
          try:
              # websockets pings the server every 20 seconds by default.
              async with connect(URL, additional_headers=HEADERS) as socket:
                  attempt = 0
                  for frame in SUBSCRIPTIONS:
                      await socket.send(json.dumps(frame))
                  async for raw in socket:
                      frame = json.loads(raw)
                      if frame["type"] == "update":
                          print(frame["channel"], frame["data"])
                      elif frame["type"] == "error":
                          print("error", frame["code"], frame["message"])
                  close_code = socket.close_code
          except InvalidStatus as error:
              if error.response.status_code == 401:
                  raise
          except ConnectionClosed as error:
              close_code = error.rcvd.code if error.rcvd else None
          except OSError:
              pass
          if close_code == 1001:
              continue
          await asyncio.sleep(min(30, 2 ** attempt) + random.random())
          attempt += 1

  asyncio.run(run())
  ```
</CodeGroup>


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.