@crossfox/ws
v0.4.0
Published
Reliable WebSocket protocol: client (reconnect, offline queue, ack) + server (Redis Streams seq/replay/dedup) + Elysia adapter. Browser and Bun/Node.
Maintainers
Readme
@crossfox/ws
Надёжный WebSocket целиком: клиент (реконнект с backoff, офлайн-очередь) + сервер (движок на Redis Streams: seq / replay / дедупликация / pub-sub между инстансами) + общий протокол. Обе стороны в одном пакете — версия пакета и есть версия протокола.
Клиент — браузер / любой runtime с WebSocket. Сервер — framework-agnostic (@crossfox/ws/server) или готовый плагин Elysia (@crossfox/ws/elysia, обычно на Bun). Для сервера нужен Redis ≥5 (peer).
Установка
bun add @crossfox/ws
# or: npm i @crossfox/wsEntry points:
| импорт | что внутри | зависимости |
|---|---|---|
| @crossfox/ws | всё ниже, кроме react | нет |
| @crossfox/ws/protocol | shared-типы фреймов, DEFAULTS, REDIS_KEYS | нет |
| @crossfox/ws/codec | jsonCodec, msgpackCodec, serverJsonCodec, serverMsgpackCodec | нет |
| @crossfox/ws/client | WSReliableClient, connectWS | нет |
| @crossfox/ws/react | useWS | react ≥18 (peer) |
| @crossfox/ws/server | createWsEngine + redis-хелперы (framework-agnostic) | redis ≥5 (peer) |
| @crossfox/ws/elysia | websocket() — готовый Elysia-плагин | elysia + redis (peer) |
Быстрый старт
import { connectWS } from "@crossfox/ws";
const client = connectWS({
url: "wss://api.example.com/ws",
refreshToken: async () => (await fetch("/auth/ws-token", { method: "POST" })
.then(r => r.json())).access_token,
});
client.onMessage((msg) => console.log(msg.channel, msg.payload));
client.subscribe(["public:items"]);
const seq = await client.send("chat:1", { text: "hi" }); // резолв после server-ackReact:
import { useWS } from "@crossfox/ws/react";
const channels = useMemo(() => ["public:items"], []);
const { status, send } = useWS({
url: WS_URL,
refreshToken: getWsToken,
channels,
onMessage: () => queryClient.invalidateQueries({ queryKey: ["items"] }),
});Настройки (все — опциональные, кроме url)
connectWS({
url: "wss://…",
// ── Аутентификация ──
token: "…", // статический токен (subprotocol ['token', <v>])
refreshToken: async () => "…", // свежий токен перед КАЖДЫМ (ре)коннектом
// ── Кодек / URL ──
codec: "json" | "msgpack", // msgpack: −30–50% на числовых payload
queryParams: { v: "2" }, // доп. query к URL
// ── Жизненный цикл ──
autoConnect: true, // false — client.connect() вручную
connectTimeoutMs: 10_000, // open+connected дольше → закрыть и реконнект
// ── Реконнект ──
reconnect: {
enabled: true,
maxAttempts: Infinity, // исчерпано → статус "failed" (connect() оживляет)
baseDelayMs: 100, // 100 → 200 → 400 → … экспонента
maxDelayMs: 30_000, // потолок
factor: 2,
jitter: 0.2, // ±20% — против thundering herd
},
// ── Heartbeat / watchdog ──
hbIntervalMs: 15_000, // клиентский hb (+отчёт о пропусках)
idleTimeoutMs: 45_000, // сервер молчит дольше → форс-реконнект; 0 = off
// ── Отправка ──
ackTimeoutMs: 0, // >0 → send() реджект без ack за это время
queue: {
enabled: true, // копить send() в офлайне, флаш после реконнекта
maxSize: 100, // переполнение → реджект старейшего
},
// ── Браузер ──
listenBrowserEvents: true, // online/visibilitychange → реконнект немедленно
// ── Колбэки / логи ──
onMessage, onConnect, onDisconnect, onError,
logger: console | null, // null — полная тишина
});Статусы: idle → connecting → connected ⇄ reconnecting, терминальные failed (лимит попыток) и closed (после disconnect()).
Гарантии доставки
- Порядок и пропуски: у каждого сообщения
seq(по каналу); клиент двигает highWaterMark по контигуальной последовательности, пропуски репортит в heartbeat — сервер доотправляет. - Replay: при реконнекте подписка уходит с
reconnect: true, resumeFrom: hwm— сервер повторяет пропущенное из Redis Stream (окно 30 мин / 10k сообщений). - Дедупликация: входящие — по
seq; исходящие — поmsgId(офлайн-очередь шлёт с тем жеmsgId, сервер не задублирует). - Идемпотентная отправка:
send(ch, data, { msgId })— повтор с тем же msgId безопасен, придёт тот жеseq.
Серверная сторона
Движок независим от HTTP-фреймворка; Redis-клиенты, ACL каналов и приём аналитики инжектятся:
import { websocket } from "@crossfox/ws/elysia";
app.use(websocket({
redis, redisPub, redisSub, // node-redis клиенты (sub — выделенный)
canAccess: (state, channel, kind) => { // единственная точка ACL: subscribe И publish
if (kind === "publish") return state.role === "admin";
if (channel.startsWith("public:")) return true;
if (channel.startsWith("user:")) return channel === `user:${state.customerUid}`;
return state.role === "admin";
},
onTrack: (state, frame) => analytics.ingest(state, frame), // опционально
// instanceId, maxChannels (50), hbIntervalMs, isDev, logger | null
}));auth соединения берётся из ws.data.auth (положи derive'ом/authPlugin'ом: { userId, role, taxId | tax_id, customerUid }).
Без Elysia — createWsEngine из @crossfox/ws/server: три метода openConnection / handleMessage / closeConnection вешаются на любой WS-сервер (Bun.serve, ws, uWS). Серверная публикация из HTTP-кода: engine.publish(channel, payload) или publishMessage(...) напрямую.
Горизонтальное масштабирование из коробки: инстансы делят Redis (streams + pub/sub) — publish на одном доставляется подписчикам всех.
Версия пакета = версия протокола: мажорный бамп при несовместимых изменениях фреймов.
Разработка
bun install
bun test # 33 теста: backoff, кодеки, клиент на фейковом WebSocket,
# e2e клиент↔сервер через локальный Redis (нет Redis — скип)
bun run check # typecheck + test
bun run build # dist/ (tsc, d.ts)License
PolyForm Noncommercial 1.0.0 —
free for non-commercial use; copyright notice required
(Required Notice: Copyright Oleksii Fursov).
Commercial use (for-profit products, SaaS, client work, etc.) needs a separate written license from the author: [email protected].
