Skip to content

Repository files navigation

WebSocket 插件化客户端

一个轻量、可扩展的 WebSocket 客户端库,采用插件化架构,核心极简,功能按需组合。

特性

  • 极简核心,仅管理连接生命周期和插件管道
  • 插件化设计,功能按需加载
  • 类型安全的事件系统(TypeScript)
  • 链式调用 API
  • 内置 4 个常用插件:自动重连、心跳检测、离线消息队列、消息路由

快速开始

import {
  WebSocketClient,
  ReconnectPlugin,
  HeartbeatPlugin,
  MessageQueuePlugin,
  MessageRouterPlugin,
} from './websocket';

const ws = new WebSocketClient('wss://example.com/stream', {
  binaryType: 'arraybuffer',
  timeout: 10000,
})
  .use(new ReconnectPlugin({ maxAttempts: 5, delay: 3000 }))
  .use(new HeartbeatPlugin({ interval: 5000, timeout: 3000 }))
  .use(new MessageQueuePlugin({ maxSize: 50 }))
  .use(new MessageRouterPlugin());

ws.connect();

API 参考

WebSocketClient

构造函数

new WebSocketClient(url: string, options?: WebSocketClientOptions)
参数 类型 默认值 说明
url string WebSocket 服务器地址
options.binaryType 'arraybuffer' | 'blob' 'arraybuffer' 接收二进制数据格式
options.timeout number 10000 连接超时时间(ms)

方法

方法 返回值 说明
use(plugin) this 注册插件,支持链式调用
connect() void 连接到服务器
close(code?, reason?) void 关闭连接
reconnect() void 重新连接(先清理旧连接)
send(data) boolean 发送数据,返回是否成功
isConnected() boolean 是否已连接
getState() number 获取 WebSocket readyState
destroy() void 销毁实例,清理所有资源

事件监听

// 订阅事件
const unsubscribe = ws.on('message', (data, event) => { ... });

// 取消订阅
unsubscribe();
// 或
ws.off('message', listener);

// 一次性监听
ws.once('open', (event) => { ... });

支持的事件

事件名 回调签名 说明
open (event: Event) => void 连接建立
close (event: CloseEvent) => void 连接关闭
error (event: Event, error?: Error) => void 连接出错
message (data: MessageData, event: MessageEvent) => void 收到消息(原始)
reconnect (attempt: number, maxAttempts: number) => void 触发重连(需 ReconnectPlugin)
message:json (data: any, event: MessageEvent) => void JSON 消息(需 MessageRouterPlugin)
message:text (data: string, event: MessageEvent) => void 文本消息(需 MessageRouterPlugin)
message:binary (data: ArrayBuffer | Blob, event: MessageEvent) => void 二进制消息(需 MessageRouterPlugin)

发送数据

send() 支持多种数据类型,普通对象会自动 JSON 序列化:

// 发送 JSON 对象(自动序列化)
ws.send({ action: 'play', id: 123 });

// 发送字符串
ws.send('hello');

// 发送二进制
ws.send(new Uint8Array([1, 2, 3]));
ws.send(someArrayBuffer);

插件详解

ReconnectPlugin — 自动重连

非正常断开时自动重连,采用指数退避策略。

import { ReconnectPlugin } from './websocket';

ws.use(new ReconnectPlugin({
  maxAttempts: 5,         // 最大重连次数,默认 5
  delay: 3000,            // 初始重连延迟(ms),默认 3000
  backoffMultiplier: 1.5, // 退避倍数,默认 1.5
  maxDelay: 30000,        // 最大延迟上限(ms),默认 30000
}));

// 监听重连事件
ws.on('reconnect', (attempt, maxAttempts) => {
  console.log(`正在第 ${attempt}/${maxAttempts} 次重连...`);
});

行为说明:

  • 手动调用 close() 不触发重连
  • 正常关闭码 (1000) 不触发重连
  • 重连延迟按 delay * backoffMultiplier^(attempt-1) 递增,不超过 maxDelay

HeartbeatPlugin — 心跳检测

定时发送 ping 消息,检测 pong 响应超时则关闭连接(配合 ReconnectPlugin 触发重连)。

import { HeartbeatPlugin } from './websocket';

ws.use(new HeartbeatPlugin({
  interval: 5000,    // 心跳间隔(ms),默认 5000
  timeout: 3000,     // pong 超时(ms),默认 3000
  pingMessage: { type: 'heartbeat', cmd: 'ping' }, // 自定义 ping 消息
  isPong: (data) => {
    // 自定义 pong 判断逻辑
    if (typeof data !== 'string') return false;
    try {
      const parsed = JSON.parse(data);
      return parsed.type === 'heartbeat' || parsed.cmd === 'pong';
    } catch {
      return data === 'pong';
    }
  },
}));

默认 pong 判断逻辑:

  • JSON 中 type === 'heartbeat'cmd === 'pong'
  • 纯文本 'pong'

MessageQueuePlugin — 离线消息队列

连接未就绪时自动缓存待发送消息,连接恢复后自动 flush。

import { MessageQueuePlugin } from './websocket';

const queuePlugin = new MessageQueuePlugin({
  maxSize: 100, // 队列最大长度,默认 100
});

ws.use(queuePlugin);

// 即使离线也可以调用 send,消息会被缓存
ws.send({ action: 'play' });

// 查看队列状态
console.log(queuePlugin.getQueueSize());

// 手动清空队列
queuePlugin.clearQueue();

行为说明:

  • 离线时 send() 返回 false,但消息已入队
  • 队列满时新消息被丢弃,控制台输出警告
  • 连接恢复后按入队顺序依次发送

MessageRouterPlugin — 消息路由

根据消息数据类型自动分发到不同事件通道。

import { MessageRouterPlugin } from './websocket';

ws.use(new MessageRouterPlugin());

// 只关心 JSON 消息
ws.on('message:json', (data, event) => {
  console.log('收到 JSON:', data);
});

// 只关心文本消息
ws.on('message:text', (text, event) => {
  console.log('收到文本:', text);
});

// 只关心二进制消息
ws.on('message:binary', (buf, event) => {
  console.log('收到二进制:', buf);
});

// 原始 message 事件仍然会触发
ws.on('message', (data, event) => {
  console.log('收到原始消息:', data);
});

路由规则:

  • string → 触发 message:text,若为合法 JSON 还额外触发 message:json
  • ArrayBuffer / Blob → 触发 message:binary
  • 不拦截原始 message 事件

插件执行顺序

插件按 use() 注册顺序执行。推荐顺序:

ws.use(new ReconnectPlugin())      // 1. 重连(最先感知断开)
  .use(new HeartbeatPlugin())      // 2. 心跳(依赖重连恢复连接)
  .use(new MessageQueuePlugin())   // 3. 队列(拦截离线发送)
  .use(new MessageRouterPlugin()); // 4. 路由(最后分发消息)

自定义插件

实现 WebSocketPlugin 接口即可:

import type { WebSocketPlugin, IWebSocketClient, SendData, MessageData } from './websocket';

class LogPlugin implements WebSocketPlugin {
  readonly name = 'logger';

  install(client: IWebSocketClient): void {
    console.log('[LogPlugin] 已安装');
  }

  onOpen(client: IWebSocketClient, event: Event): void {
    console.log('[LogPlugin] 连接已建立');
  }

  onMessage(client: IWebSocketClient, data: MessageData, event: MessageEvent): void {
    console.log('[LogPlugin] 收到消息:', data);
    // 不返回 false,消息继续传递
  }

  onBeforeSend(client: IWebSocketClient, data: SendData): SendData | false {
    console.log('[LogPlugin] 即将发送:', data);
    return data; // 返回数据继续发送,返回 false 拦截
  }

  onClose(client: IWebSocketClient, event: CloseEvent): void {
    console.log('[LogPlugin] 连接已关闭', event.code, event.reason);
  }

  onError(client: IWebSocketClient, event: Event, error?: Error): void {
    console.error('[LogPlugin] 错误:', error?.message);
  }

  destroy(): void {
    console.log('[LogPlugin] 已销毁');
  }
}

插件钩子说明

钩子 调用时机 返回值作用
install(client) 注册时调用一次 可劫持 client 方法
onOpen(client, event) 连接建立后
onMessage(client, data, event) 收到消息时 返回 false 拦截消息
onBeforeSend(client, data) 发送前 返回修改后的数据,或 false 阻止发送
onClose(client, event) 连接关闭后
onError(client, event, error) 连接出错时
destroy() 实例销毁时 清理资源

完整示例

import {
  WebSocketClient,
  ReconnectPlugin,
  HeartbeatPlugin,
  MessageQueuePlugin,
  MessageRouterPlugin,
} from './websocket';

// 创建实例
const ws = new WebSocketClient('wss://api.example.com/ws', {
  binaryType: 'arraybuffer',
  timeout: 8000,
})
  .use(new ReconnectPlugin({ maxAttempts: 10, delay: 2000 }))
  .use(new HeartbeatPlugin({ interval: 10000, timeout: 5000 }))
  .use(new MessageQueuePlugin({ maxSize: 200 }))
  .use(new MessageRouterPlugin());

// 事件监听
ws.on('open', () => console.log('已连接'));
ws.on('close', (e) => console.log('已断开', e.code));
ws.on('error', (_, err) => console.error('错误', err?.message));
ws.on('reconnect', (n, max) => console.log(`重连中 ${n}/${max}`));

ws.on('message:json', (data) => {
  switch (data.type) {
    case 'notification':
      showNotification(data.payload);
      break;
    case 'update':
      updateState(data.payload);
      break;
  }
});

ws.on('message:binary', (buf) => {
  processAudioChunk(buf as ArrayBuffer);
});

// 连接
ws.connect();

// 发送(离线时自动入队)
ws.send({ type: 'subscribe', channel: 'live' });

// 页面卸载时销毁
window.addEventListener('beforeunload', () => {
  ws.destroy();
});

类型导出

export type {
  WebSocketPlugin,
  WebSocketClientOptions,
  WebSocketEventMap,
  CoreEventMap,
  PluginEventMap,
  SendData,
  MessageData,
  IWebSocketClient,
  ReconnectConfig,
  HeartbeatConfig,
  MessageQueueConfig,
};

About

轻量、可扩展的 WebSocket 插件化 js

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages