Skip to content

实时消息推送

本章实现项目的核心功能——实时消息推送。Flutter 客户端通过 WebSocket 与服务端保持长连接,外部系统(如 Postman)调用 HTTP API 推送消息,服务端通过 WebSocket 实时送达客户端。

整体架构

text
┌──────────────┐     HTTP POST      ┌──────────────┐     WebSocket      ┌──────────────┐
│              │ ─────────────────→ │              │ ─────────────────→ │              │
│  Postman /   │   /api/push/send   │   Nuxt 服务   │   ws://.../_ws    │   Flutter    │
│  外部系统     │                    │              │                   │   客户端      │
│              │ ←───────────────── │              │ ←──────────────── │              │
└──────────────┘    200 OK 响应     └──────────────┘    实时消息推送    └──────────────┘

                                         │  Redis PubSub
                                         │  (多实例部署时跨进程通信)

                                   ┌──────────────┐
                                   │   实例 2      │
                                   │   (可选)      │
                                   └──────────────┘

数据流向

  1. Postman 调用 POST /api/push/send,指定目标用户和消息内容
  2. 服务端收到请求,找到目标用户的 WebSocket 连接
  3. 通过 WebSocket 推送消息给 Flutter 客户端
  4. Flutter 客户端收到消息,显示通知

数据库设计

通知表 (notifications)

ts
// server/database/schema.ts — 新增
export const notifications = pgTable('notifications', {
  id: serial('id').primaryKey(),
  userId: integer('user_id').notNull(),              // 目标用户 ID
  title: varchar('title', { length: 100 }).notNull(), // 通知标题
  content: text('content').notNull(),                 // 通知内容
  type: varchar('type', { length: 20 }).notNull(),    // 通知类型:system/order/promotion
  isRead: boolean('is_read').default(false).notNull(),// 是否已读
  createdAt: timestamp('created_at').defaultNow().notNull(),
}, (table) => ({
  userIdx: index('idx_notifications_user_id').on(table.userId),
  readIdx: index('idx_notifications_is_read').on(table.isRead),
  createdIdx: index('idx_notifications_created').on(table.createdAt),
}))

为什么要存数据库?

  • 用户离线时,消息不会丢失,上线后可以拉取未读消息
  • 提供消息历史记录
  • 支持消息已读/未读状态管理
  • 通知类型便于分类筛选

添加 shared 类型

ts
// shared/types/index.ts — 新增
export interface Notification {
  id: number
  userId: number
  title: string
  content: string
  type: 'system' | 'order' | 'promotion'
  isRead: boolean
  createdAt: string
}

// WebSocket 消息类型
export interface WsMessage {
  type: 'connected' | 'notification' | 'ping' | 'pong' | 'auth_ok' | 'error'
  data?: any
}

WebSocket 服务端

连接管理器

核心模块——管理所有 WebSocket 连接,支持按用户 ID 查找连接:

ts
// server/utils/ws-manager.ts
interface ManagedPeer {
  peer: any
  userId: number
  connectedAt: number
  lastPing: number
  pingTimer: ReturnType<typeof setInterval>
}

// 所有连接:peerId → 连接信息
const connections = new Map<string, ManagedPeer>()

// 用户连接:userId → peerId(一个用户可能有多个连接,如手机+平板)
const userConnections = new Map<number, Set<string>>()

export const wsManager = {
  /** 注册连接 */
  add(peer: any, userId: number) {
    const now = Date.now()

    // 心跳定时器
    const pingTimer = setInterval(() => {
      const managed = connections.get(peer.id)
      if (!managed) return  // 连接已不存在

      if (Date.now() - managed.lastPing > 60000) {
        // 60 秒无响应,断开
        peer.close(4000, '心跳超时')
        return
      }
      try {
        peer.send({ type: 'ping' })
      } catch {
        // 发送失败,连接可能已断
        wsManager.remove(peer.id)
      }
    }, 30000)

    const managed: ManagedPeer = {
      peer,
      userId,
      connectedAt: now,
      lastPing: now,
      pingTimer,
    }

    connections.set(peer.id, managed)

    // 维护用户→连接映射
    if (!userConnections.has(userId)) {
      userConnections.set(userId, new Set())
    }
    userConnections.get(userId)!.add(peer.id)

    logger.info('WebSocket 连接', { peerId: peer.id, userId, onlineCount: connections.size })
  },

  /** 移除连接 */
  remove(peerId: string) {
    const managed = connections.get(peerId)
    if (!managed) return

    clearInterval(managed.pingTimer)
    connections.delete(peerId)

    const userSet = userConnections.get(managed.userId)
    if (userSet) {
      userSet.delete(peerId)
      if (userSet.size === 0) {
        userConnections.delete(managed.userId)
      }
    }

    logger.info('WebSocket 断开', { peerId, userId: managed.userId, onlineCount: connections.size })
  },

  /** 更新心跳时间 */
  updatePing(peerId: string) {
    const managed = connections.get(peerId)
    if (managed) {
      managed.lastPing = Date.now()
    }
  },

  /** 给指定用户推送消息 */
  sendToUser(userId: number, message: any) {
    const userSet = userConnections.get(userId)
    if (!userSet || userSet.size === 0) {
      logger.info('用户不在线,跳过推送', { userId })
      return false
    }

    const text = JSON.stringify(message)
    let sentCount = 0

    for (const peerId of userSet) {
      const managed = connections.get(peerId)
      if (managed) {
        try {
          managed.peer.send(text)
          sentCount++
        } catch (error) {
          logger.error('推送失败', { peerId, userId, error })
          wsManager.remove(peerId)
        }
      }
    }

    return sentCount > 0
  },

  /** 广播给所有在线用户 */
  broadcast(message: any) {
    const text = JSON.stringify(message)
    for (const [peerId, managed] of connections) {
      try {
        managed.peer.send(text)
      } catch {
        wsManager.remove(peerId)
      }
    }
  },

  /** 获取在线用户数 */
  getOnlineUserCount() {
    return userConnections.size
  },

  /** 获取在线连接数 */
  getOnlineConnectionCount() {
    return connections.size
  },

  /** 检查用户是否在线 */
  isUserOnline(userId: number) {
    return userConnections.has(userId) && userConnections.get(userId)!.size > 0
  },
}

为什么不用 Nitro 的 peer.subscribe/publish

  • subscribe/publish 基于主题(topic),需要从 HTTP API 中调用 publish,但 HTTP handler 没有直接的 peer 引用
  • 自己管理连接 Map 更灵活:可以按用户 ID 精确查找、统计在线人数、支持一个用户多端连接
  • 多实例部署时,Map 不能跨进程,但可以配合 Redis PubSub 解决(见后文)

WebSocket Ticket API

为浏览器/Flutter Web 客户端提供短期一次性 Ticket,避免 JWT 暴露在 URL 中。Flutter 原生客户端可直接用 Authorization 头,无需 Ticket。

ts
// server/api/ws/ticket.post.ts
import { randomUUID } from 'crypto'

// 存储 ticket(生产环境建议用 Redis)
const tickets = new Map<string, { userId: number; expiresAt: number }>()

export default defineEventHandler(async (event) => {
  const userId = event.context.userId
  if (!userId) {
    throw createError({ statusCode: 401, statusMessage: 'Unauthorized' })
  }

  // 生成一次性 Ticket,30 秒有效
  const ticket = randomUUID()
  tickets.set(ticket, { userId, expiresAt: Date.now() + 30_000 })

  // 清理过期 ticket
  for (const [key, value] of tickets) {
    if (value.expiresAt < Date.now()) tickets.delete(key)
  }

  return { ticket }
})

/** 验证并消耗 ticket(供 WebSocket handler 调用) */
export function consumeTicket(ticket: string): number | null {
  const data = tickets.get(ticket)
  if (!data) return null

  tickets.delete(ticket)  // 一次性使用

  if (data.expiresAt < Date.now()) return null

  return data.userId
}

WebSocket 路由

服务端同时支持两种认证方式:原生客户端用 Authorization 头,浏览器/Flutter Web 用 Ticket。

ts
// server/api/_ws.ts
import { consumeTicket } from './ws/ticket.post'

export default defineWebSocketHandler({
  upgrade(request) {
    // 校验 Origin(防止跨站 WebSocket 劫持)
    const origin = request.headers.get('origin')
    const allowed = ['https://your-app.com']
    if (origin && !allowed.includes(origin)) {
      throw new Response('Forbidden origin', { status: 403 })
    }

    // 认证方式一:Authorization 头(Flutter 原生、Node.js 等原生客户端)
    const authHeader = request.headers.get('authorization')
    const token = authHeader?.replace('Bearer ', '')

    if (token) {
      try {
        const payload = verifyToken(token)
        request.context.userId = payload.userId
      } catch {
        throw new Response('Invalid token', { status: 403 })
      }
    }

    // 认证方式二:URL 参数传 Ticket(浏览器、Flutter Web)
    const url = new URL(request.url)
    const ticket = url.searchParams.get('ticket')

    if (!ticket) {
      throw new Response('Missing authentication', { status: 401 })
    }

    const userId = consumeTicket(ticket)
    if (!userId) {
      throw new Response('Invalid or expired ticket', { status: 403 })
    }

    request.context.userId = userId
  },

  open(peer) {
    // upgrade 已完成认证,peer.context.userId 可用
    wsManager.add(peer, peer.context.userId)
    peer.send({ type: 'connected', data: { userId: peer.context.userId } })
  },

  message(peer, message) {
    try {
      const data = message.json()

      switch (data.type) {
        case 'pong':
          wsManager.updatePing(peer.id)
          break

        default:
          peer.send({ type: 'error', data: { message: `未知的消息类型: ${data.type}` } })
      }
    } catch {
      // 忽略无法解析的消息
    }
  },

  close(peer) {
    wsManager.remove(peer.id)
  },

  error(peer, error) {
    logger.error('WebSocket 错误', { peerId: peer.id, error: error.message })
    wsManager.remove(peer.id)
  },
})

认证策略

  • Flutter 原生(Android/iOS/Desktop):直接用 Authorization: Bearer <jwt> 头,最简单最安全,JWT 不出现在 URL 中
  • 浏览器 / Flutter Web:先通过 HTTP API 换取一次性 Ticket,再用 Ticket 连接 WebSocket
  • 服务端 upgrade 钩子优先检查 Authorization 头,未提供则回退到 Ticket

消息推送 API

外部系统(Postman、其他服务)通过 HTTP API 推送消息,服务端再通过 WebSocket 送达客户端:

推送消息

ts
// server/api/push/send.post.ts
import { z } from 'zod'

const sendSchema = z.object({
  userId: z.number().int().positive('用户 ID 必须是正整数'),
  title: z.string().min(1, '标题不能为空').max(100),
  content: z.string().min(1, '内容不能为空'),
  type: z.enum(['system', 'order', 'promotion']).default('system'),
})

export default defineEventHandler(async (event) => {
  const body = await readValidatedBody(event, sendSchema.parse)

  // 1. 存入数据库(确保消息不丢失)
  const [notification] = await db.insert(notifications).values({
    userId: body.userId,
    title: body.title,
    content: body.content,
    type: body.type,
  }).returning()

  // 2. 实时推送(如果用户在线)
  const pushed = wsManager.sendToUser(body.userId, {
    type: 'notification',
    data: notification,
  })

  // 3. 返回结果
  return {
    success: true,
    notification,
    pushed,  // 是否成功推送到客户端
  }
})

批量推送

ts
// server/api/push/batch.post.ts
import { z } from 'zod'

const batchSchema = z.object({
  userIds: z.array(z.number().int().positive()).min(1).max(1000, '最多 1000 人'),
  title: z.string().min(1).max(100),
  content: z.string().min(1),
  type: z.enum(['system', 'order', 'promotion']).default('system'),
})

export default defineEventHandler(async (event) => {
  const body = await readValidatedBody(event, batchSchema.parse)

  // 1. 批量写入数据库
  const values = body.userIds.map(userId => ({
    userId,
    title: body.title,
    content: body.content,
    type: body.type,
  }))

  const inserted = await db.insert(notifications).values(values).returning()

  // 2. 推送给在线用户
  let pushCount = 0
  for (const userId of body.userIds) {
    if (wsManager.sendToUser(userId, {
      type: 'notification',
      data: { title: body.title, content: body.content, type: body.type },
    })) {
      pushCount++
    }
  }

  return {
    success: true,
    total: body.userIds.length,
    pushed: pushCount,
    offline: body.userIds.length - pushCount,
  }
})

全站广播

ts
// server/api/push/broadcast.post.ts
import { z } from 'zod'

const broadcastSchema = z.object({
  title: z.string().min(1).max(100),
  content: z.string().min(1),
})

export default defineEventHandler(async (event) => {
  // 仅管理员可广播
  if (event.context.userRole !== 'admin') {
    throw createError({ statusCode: 403, statusMessage: 'Forbidden', message: '仅管理员可发送广播' })
  }

  const body = await readValidatedBody(event, broadcastSchema.parse)

  // 广播给所有在线用户
  wsManager.broadcast({
    type: 'notification',
    data: {
      title: body.title,
      content: body.content,
      type: 'system',
    },
  })

  // 不存数据库(广播量大,按需决定是否存储)
  return {
    success: true,
    onlineUsers: wsManager.getOnlineUserCount(),
  }
})

通知历史查询

ts
// server/api/notifications.get.ts
export default defineEventHandler(async (event) => {
  const userId = event.context.userId
  if (!userId) {
    throw createError({ statusCode: 401, statusMessage: 'Unauthorized', message: '未认证' })
  }

  const query = getQuery(event)
  const page = Number(query.page) || 1
  const pageSize = Math.min(Number(query.pageSize) || 20, 50)
  const offset = (page - 1) * pageSize

  const [list, [{ count }]] = await Promise.all([
    db.select()
      .from(notifications)
      .where(eq(notifications.userId, userId))
      .orderBy(desc(notifications.createdAt))
      .limit(pageSize)
      .offset(offset),
    db.select({ count: sql`count(*)` })
      .from(notifications)
      .where(eq(notifications.userId, userId)),
  ])

  return { list, total: count, page, pageSize }
})

标记已读

ts
// server/api/notifications/read.put.ts
import { z } from 'zod'

const schema = z.object({
  ids: z.array(z.number().int().positive()),  // 通知 ID 列表
})

export default defineEventHandler(async (event) => {
  const userId = event.context.userId
  if (!userId) {
    throw createError({ statusCode: 401, statusMessage: 'Unauthorized', message: '未认证' })
  }

  const { ids } = await readValidatedBody(event, schema.parse)

  await db.update(notifications)
    .set({ isRead: true })
    .where(and(
      inArray(notifications.id, ids),
      eq(notifications.userId, userId),  // 只能标记自己的通知
    ))

  return { success: true }
})

未读消息拉取(上线补发)

用户上线后,先拉取离线期间的未读消息,确保不遗漏:

ts
// server/api/_ws.ts — 在 open 中加入拉取逻辑(upgrade 钩子已完成认证)
open(peer) {
  const userId = peer.context.userId  // upgrade 中通过 request.context 已设置
  wsManager.add(peer, userId)
  peer.send({ type: 'connected', data: { userId } })

  // 拉取未读消息并推送
  fetchUnreadNotifications(userId).then(notifications => {
    if (notifications.length > 0) {
      for (const notification of notifications) {
        peer.send({ type: 'notification', data: notification })
      }
      logger.info('补发未读消息', { userId, count: notifications.length })
    }
  })
},
ts
// server/utils/notification.ts
export const fetchUnreadNotifications = async (userId: number) => {
  return await db.select()
    .from(notifications)
    .where(and(
      eq(notifications.userId, userId),
      eq(notifications.isRead, false),
    ))
    .orderBy(desc(notifications.createdAt))
    .limit(50)  // 最多补发 50 条
}

上线补发逻辑

用户连接 WebSocket 后,立即查询数据库中的未读消息并推送。这保证了即使推送时用户不在线,上线后也能收到消息。

用 Postman 测试

1. 先登录获取 Token

text
POST http://localhost:3000/api/auth/login
Content-Type: application/json

{
  "phone": "13800138000",
  "password": "your-password"
}

记下返回的 token

2. 推送消息

text
POST http://localhost:3000/api/push/send
Content-Type: application/json
Authorization: Bearer <token>

{
  "userId": 1,
  "title": "订单更新",
  "content": "您的订单 #10086 已发货",
  "type": "order"
}

返回:

json
{
  "success": true,
  "notification": {
    "id": 1,
    "userId": 1,
    "title": "订单更新",
    "content": "您的订单 #10086 已发货",
    "type": "order",
    "isRead": false,
    "createdAt": "2025-01-15T10:30:00.000Z"
  },
  "pushed": true
}

pushed: true

表示消息已通过 WebSocket 实时送达客户端。如果为 false,说明用户当前不在线,但消息已存入数据库,用户上线后会补发。

3. 批量推送

http
POST http://localhost:3000/api/push/batch
Content-Type: application/json
Authorization: Bearer <token>

{
  "userIds": [1, 2, 3],
  "title": "系统维护通知",
  "content": "系统将于今晚 22:00 进行维护",
  "type": "system"
}

4. 测试 WebSocket 连接(Postman)

Postman 支持 WebSocket 测试,两种认证方式都可以:

方式一:Authorization 头(推荐,模拟原生客户端)

  1. 新建 WebSocket Request
  2. URL 填写:ws://localhost:3000/api/_ws
  3. 在 Headers 中添加:Authorization: Bearer <your-jwt-token>
  4. 点击 Connect
  5. 发送心跳消息:{"type":"pong"}
  6. 观察:当通过 HTTP API 推送消息时,这里会实时收到

方式二:Ticket(模拟浏览器/Flutter Web 客户端)

  1. 先调用 POST /api/ws/ticket(带 Authorization: Bearer <token>)获取 Ticket
  2. 新建 WebSocket Request
  3. URL 填写:ws://localhost:3000/api/_ws?ticket=<上一步获取的ticket>
  4. 点击 Connect
  5. 发送心跳消息:{"type":"pong"}
  6. 观察:当通过 HTTP API 推送消息时,这里会实时收到

Flutter 客户端接入

以下是 Flutter 端连接 WebSocket 的核心代码。

Flutter 原生端(推荐:Authorization 头)

Flutter 原生平台(Android/iOS/Desktop)使用 IOWebSocketChannel,通过 Authorization 头传递 JWT——最简单最安全:

dart
// lib/services/websocket_service.dart
import 'dart:convert';
import 'dart:async';
import 'dart:io' show WebSocket;              // 原生平台
import 'package:web_socket_channel/io.dart';    // 原生平台
import 'package:web_socket_channel/web_socket_channel.dart';

class WebSocketService {
  WebSocketChannel? _channel;
  String? _jwt;
  Timer? _reconnectTimer;
  Timer? _pingTimer;
  int _retryCount = 0;
  static const maxRetries = 5;

  final String wsUrl;  // wss://your-server.com

  final void Function(Map<String, dynamic>)? onNotification;

  WebSocketService({
    required this.wsUrl,
    this.onNotification,
  });

  /// 连接 WebSocket(使用 Authorization 头)
  void connect(String jwt) {
    _jwt = jwt;
    _retryCount = 0;

    try {
      // ✅ 通过 Authorization 头传递 JWT,不会出现在 URL 中
      final channel = IOWebSocketChannel.connect(
        Uri.parse('$wsUrl/api/_ws'),
        headers: {
          'Authorization': 'Bearer $jwt',
        },
      );
      _channel = channel;

      _channel!.stream.listen(
        (message) {
          _handleMessage(message);
          _retryCount = 0;
        },
        onDone: () {
          print('WebSocket 连接关闭');
          _scheduleReconnect();
        },
        onError: (error) {
          print('WebSocket 错误: $error');
          _scheduleReconnect();
        },
      );

      _startPing();
    } catch (e) {
      print('WebSocket 连接失败: $e');
      _scheduleReconnect();
    }
  }

  /// 处理服务端消息
  void _handleMessage(dynamic raw) {
    final data = jsonDecode(raw as String) as Map<String, dynamic>;

    switch (data['type']) {
      case 'ping':
        send({'type': 'pong'});
        break;
      case 'notification':
        onNotification?.call(data['data'] as Map<String, dynamic>);
        _showLocalNotification(data['data']);
        break;
      case 'connected':
        print('WebSocket 已连接,userId: ${data['data']['userId']}');
        break;
    }
  }

  /// 发送消息
  void send(Map<String, dynamic> message) {
    _channel?.sink.add(jsonEncode(message));
  }

  /// 心跳
  void _startPing() {
    _pingTimer?.cancel();
    _pingTimer = Timer.periodic(
      const Duration(seconds: 30),
      (_) => send({'type': 'pong'}),
    );
  }

  /// 自动重连(指数退避)
  void _scheduleReconnect() {
    _pingTimer?.cancel();

    if (_retryCount >= maxRetries) {
      print('达到最大重试次数,停止重连');
      return;
    }
    _retryCount++;
    final delay = Duration(seconds: _retryCount * 2);
    print('${delay.inSeconds} 秒后重连(第 $_retryCount 次)');

    _reconnectTimer = Timer(delay, () {
      if (_jwt != null) connect(_jwt!);
    });
  }

  /// 显示本地通知
  void _showLocalNotification(Map<String, dynamic> data) {
    // 使用 flutter_local_notifications 等包显示系统通知
    print('收到通知: ${data['title']} - ${data['content']}');
  }

  /// 断开连接
  void disconnect() {
    _pingTimer?.cancel();
    _reconnectTimer?.cancel();
    _channel?.sink.close();
    _channel = null;
    _jwt = null;
    _retryCount = 0;
  }
}

Flutter Web 端(Ticket 方案回退)

Flutter Web 不支持 IOWebSocketChannel(浏览器 WebSocket API 不支持自定义头),需要改用 Ticket 方案:

dart
import 'package:web_socket_channel/web_socket_channel.dart';
import 'package:http/http.dart' as http;

/// Flutter Web 连接(通过 Ticket)
Future<WebSocketChannel> connectForWeb(String jwt, String baseUrl, String wsUrl) async {
  // 1. 通过 HTTP API 换取一次性 Ticket
  final response = await http.post(
    Uri.parse('$baseUrl/api/ws/ticket'),
    headers: {'Authorization': 'Bearer $jwt'},
  );

  if (response.statusCode != 200) {
    throw Exception('获取 WebSocket ticket 失败');
  }

  final ticket = jsonDecode(response.body)['ticket'] as String;

  // 2. 用 Ticket 连接 WebSocket
  return WebSocketChannel.connect(
    Uri.parse('$wsUrl/api/_ws?ticket=$ticket'),
  );
}

Flutter 端要点

  • 原生端首选 IOWebSocketChannel + Authorization:JWT 通过请求头传递,不出现在 URL 中,一次连接直接完成
  • Flutter Web 用 Ticket 方案:浏览器不支持自定义头,需先通过 HTTP API 换取短期 Ticket
  • 自动重连:网络断开后指数退避重试(2→4→6→8→10 秒)
  • 心跳保活:每 30 秒发一次 pong,保持连接活跃
  • 跨平台统一:可用 kIsWeb 判断平台,自动选择连接方式

在 Flutter 中使用

dart
// main.dart
final wsService = WebSocketService(
  wsUrl: 'wss://your-server.com',
  onNotification: (data) {
    print('收到推送: $data');
  },
);

// 登录成功后连接(传 JWT)
wsService.connect(loginResponse.token);

// 退出登录时断开
wsService.disconnect();

多实例部署(Redis PubSub)

单实例部署时,wsManager 的内存 Map 足够。但多实例(如 Docker 集群)部署时,连接分布在不同实例上,需要 Redis PubSub 跨实例通信:

text
实例 1 ─── Redis PubSub ─── 实例 2
  │                           │
  ├─ 用户 A 连接              ├─ 用户 C 连接
  └─ 用户 B 连接              └─ 用户 D 连接

Postman → POST /api/push/send (userId=C)
→ 实例 1 收到请求
→ 发现用户 C 不在本实例
→ 发布消息到 Redis PubSub
→ 实例 2 收到 Redis 消息
→ 推送给用户 C

实现

ts
// server/utils/ws-redis.ts
export const setupRedisPubSub = () => {
  const storage = useStorage('redis')

  // 订阅推送频道
  storage.watch?.('ws:push', async (event, key, value) => {
    if (event === 'update' && value) {
      const message = JSON.parse(value)
      wsManager.sendToUser(message.userId, message.payload)
    }
  })
}

/** 跨实例推送 */
export const crossInstancePush = async (userId: number, payload: any) => {
  // 先尝试本实例推送
  if (wsManager.sendToUser(userId, payload)) {
    return true // 本实例推送成功
  }

  // 本实例没有该用户,通过 Redis 通知其他实例
  const storage = useStorage('redis')
  await storage.setItem('ws:push', JSON.stringify({ userId, payload }), { ttl: 10 })
  return false // 不确定其他实例是否推送成功
}

INFO

Redis PubSub 的局限:Nitro 的 useStorage('redis') 不直接支持 PubSub 上面的代码是概念演示。实际项目中推荐使用 ioredis 的 PubSub 功能

ts
// server/utils/redis-pubsub.ts
import Redis from 'ioredis'

const publisher = new Redis(redisConfig)
const subscriber = new Redis(redisConfig)

// 订阅频道
subscriber.subscribe('ws:push')
subscriber.on('message', (channel, message) => {
if (channel === 'ws:push') {
const { userId, payload } = JSON.parse(message)
wsManager.sendToUser(userId, payload)
}
})

// 发布消息
export const publishPush = (userId: number, payload: any) => {
publisher.publish('ws:push', JSON.stringify({ userId, payload }))
}

修改推送 API

ts
// server/api/push/send.post.ts — 使用跨实例推送
export default defineEventHandler(async (event) => {
  const body = await readValidatedBody(event, sendSchema.parse)

  const [notification] = await db.insert(notifications).values({
    userId: body.userId,
    title: body.title,
    content: body.content,
    type: body.type,
  }).returning()

  // 使用跨实例推送替代 wsManager.sendToUser
  const pushed = await crossInstancePush(body.userId, {
    type: 'notification',
    data: notification,
  })

  return { success: true, notification, pushed }
})

Nginx 配置

生产环境 Nginx 需要额外配置 WebSocket 代理:

nginx
# nginx.conf — 在 server 块中添加
location /api/_ws {
    proxy_pass http://127.0.0.1:3000;
    proxy_http_version 1.1;
    proxy_set_header Upgrade $http_upgrade;
    proxy_set_header Connection "upgrade";
    proxy_set_header Host $host;
    proxy_set_header X-Real-IP $remote_addr;
    proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
    proxy_set_header X-Forwarded-Proto $scheme;
    proxy_read_timeout 86400;      # WebSocket 长连接 24 小时超时
    proxy_send_timeout 86400;
}

关键配置解释

  • Upgrade + Connection: "upgrade":告诉 Nginx 这是 WebSocket 升级请求
  • proxy_read_timeout 86400:Nginx 默认 60 秒无数据就断开,WebSocket 需要更长时间
  • 这个 location 块要放在通用 /api/ 代理之前,确保 WebSocket 请求优先匹配

完整文件清单

目录文件说明
server/utils/ws-manager.tsWebSocket 连接管理器
notification.ts通知查询工具函数
ws-redis.tsRedis PubSub 跨实例推送(可选)
server/api/_ws.tsWebSocket 端点
ws/ticket.post.tsWebSocket 认证 Ticket API
push/send.post.ts单人推送 API
push/batch.post.ts批量推送 API
push/broadcast.post.ts全站广播 API
notifications.get.ts通知列表 API
notifications/read.put.ts标记已读 API

常见问题

Flutter 客户端连接后马上断开?

检查 Token 是否有效。连接时 Token 通过 URL 传递,确保没有编码问题。如果 Token 包含特殊字符,需要 URL 编码:

dart
final encodedToken = Uri.encodeComponent(token);
final uri = Uri.parse('ws://your-server.com/api/_ws?token=$encodedToken');

消息推送 API 返回 pushed: false

说明目标用户当前不在线。消息已存入数据库,用户上线连接时会自动补发。这是正常行为。

如何限制推送频率?

在推送 API 中加速率限制,防止被恶意刷接口:

ts
// server/middleware/rate-limit.ts 中添加
const pushRateLimits = {
  '/api/push/send': { max: 60, window: 60 },     // 每分钟 60 次
  '/api/push/batch': { max: 10, window: 60 },     // 每分钟 10 次
  '/api/push/broadcast': { max: 5, window: 3600 }, // 每小时 5 次
}

多个 Flutter 设备同时登录?

wsManager 支持一个用户多个连接(手机+平板)。所有设备都会收到推送消息。如果需要"只推送到最新活跃设备",可以在 add 时关闭该用户的旧连接。

WebSocket 连接数上限?

Node.js 单进程默认最大约 64000 个 WebSocket 连接。如果需要更多,可以:

  1. 增加系统文件描述符限制:ulimit -n 65535
  2. 使用多实例部署
  3. 使用 Redis PubSub 跨实例通信

基于 Nuxt 4 官方文档整理编写