Files
EchoChat/media-server/poc/server.mjs
bujinyuan 6e2e792db0 feat(phase2e-2): media-server Task 0/1/2 落地 + 代码审查修复
本次提交一次性落盘 Phase 2e-2 的 Task 0(PoC Spike)、Task 1(骨架 + Fastify 5 升级)
和 Task 2(9 个内部 REST API + code-reviewer 审查修复),覆盖 media-server 子项目
从零到可用的全部工作。

【Task 0 - PoC Spike】
- media-server/poc/:Node + mediasoup + fastify-websocket + 前端 mediasoup-client
- Playwright 双 tab 自动化验证 2 人会议:4 transports / 4 producers / 4 consumers / RSS 61MB
- media-server/docs/poc-notes.md 归档 7 项关键坑 + 启动步骤 + Task 1/2/9 复用映射
- 锁定技术栈:mediasoup + fastify + mediasoup-client,不改用 livekit-server

【Task 1 - 骨架 + Fastify 5 升级】
- 依赖版本:mediasoup@3.19.0 + fastify@5.8.5 + fastify-plugin@5.1.0
           + @fastify/sensible@6.0.4 + @fastify/websocket@11.2.0
           + pino@9.3.2 + zod@3.23.8
- src 五件套:app.ts / config.ts / utils/logger.ts / mediasoup/worker.ts
             / middlewares/internal-auth.ts
- /healthz + /readyz + /internal/info 三端点实测通过
- X-Internal-Token 鉴权用 timingSafeEqual 防侧信道
- mediasoup Worker died 指数退避自愈通过 kill -9 验证
- 多阶段 Dockerfile:非 root + curl HEALTHCHECK + 暴露 40000-40199 UDP/TCP
- Fastify 4 → 5 升级:loggerInstance: logger 消灭两个 pino 实例

【Task 2 - 9 个内部 REST API + 审查修复】
接口全部挂 /internal/v1/* 前缀,覆盖 Router/Transport/Producer/Consumer 完整生命周期:
- POST /routers, DELETE /routers/:id
- POST /transports, POST /transports/:id/connect
- POST /producers, DELETE /producers/:id
- POST /consumers, POST /consumers/:id/resume, DELETE /consumers/:id

工程特性:
- zod 手动 parse + 全局 errorHandler(不引入 fastify-type-provider-zod 避免
  zod v4 依赖冲突)
- AppError 统一错误码(NOT_FOUND/CONFLICT/CAN_NOT_CONSUME/ROUTER_LIMIT_EXCEEDED
  /MEDIASOUP_ERROR/VALIDATION_ERROR/UNAUTHORIZED/INTERNAL_ERROR)
- 所有资源用 Map + observer.once('close') 自清理,Consumer 监听 producerclose 级联关闭
- Consumer 强制 paused:true 创建,/resume 独立接口
- Transport direction 强约束:recv 拒 produce、send 拒 consume

code-reviewer 子代理审查"有条件通过",同步修复:
- M1 _clearXxxMap 新增 src/utils/test-guard.ts#assertTestOnly 守卫(生产误调用抛错)
- M2 新增 src/schemas/rtp.ts 对 rtpParameters / rtpCapabilities 做 codecs 浅层校验
  (mimeType/clockRate/payloadType 必填、codecs ≥1),消除 as unknown as 双跳断言
- m1 connectTransport 改乐观锁:先置位再 await,失败回退
- m2 改读 consumer.producerPaused(更符合 mediasoup 语义)
- m3 producerclose 改为 once(风格一致)
- m5 internal-auth 改为反向白名单 PRIVATE_PATH_PREFIXES = ['/internal/']

验证:
- typecheck / lint 0 错误
- vitest:65 passed / 8 spec 文件
- 覆盖率 stmts 82.87% / branches 75.83% / funcs 91.3% / lines 82.87%
- 9 接口 happy path + 6 类错误路径 curl 手测全部按预期返回

【文档同步】
- CURRENT_STATUS.md +232 行:新增 Task 0/1/2 完整记录 + 代码审查修复章节 + 延后清单
- project-context.mdc:Task 0/1/2 状态同步
- phase2e-2-design.md +20 行:Fastify 5 升级相关决策记录
- phase2e-2-implementation.plan.md +117 行:Task 0/1/2 实际产出 + 修复记录

余下 Minor/Nits(m4/m6~m10/n1~n10)登记至 Task 16 收尾清单一次性清扫。

Made-with: Cursor
2026-04-21 15:31:01 +08:00

349 lines
10 KiB
JavaScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// Phase 2e-2 Task 0: mediasoup PoC Server
// 架构Fastify 提供静态资源 + WS 信令 + mediasoup 单 Worker 单 Router
// 仅用于验证技术栈PoC 结束后可删除或演化为正式 media-server
import Fastify from 'fastify';
import fastifyStatic from '@fastify/static';
import fastifyWebsocket from '@fastify/websocket';
import * as mediasoup from 'mediasoup';
import pino from 'pino';
import path from 'node:path';
import { fileURLToPath } from 'node:url';
import { randomUUID } from 'node:crypto';
import os from 'node:os';
const __dirname = path.dirname(fileURLToPath(import.meta.url));
const logger = pino({
transport: {
target: 'pino-pretty',
options: { colorize: true, translateTime: 'SYS:HH:MM:ss.l' },
},
});
// ---------------- Config ----------------
const HTTP_PORT = Number(process.env.HTTP_PORT ?? 3300);
const LISTEN_IP = process.env.MEDIASOUP_LISTEN_IP ?? '0.0.0.0';
const ANNOUNCED_IP = process.env.MEDIASOUP_ANNOUNCED_IP ?? ''; // 本机留空走 iceCandidates 自动
const RTC_MIN_PORT = Number(process.env.MEDIASOUP_RTC_MIN_PORT ?? 40000);
const RTC_MAX_PORT = Number(process.env.MEDIASOUP_RTC_MAX_PORT ?? 40099);
// 媒体编解码Router 创建时注入)
const mediaCodecs = [
{
kind: 'audio',
mimeType: 'audio/opus',
clockRate: 48000,
channels: 2,
},
{
kind: 'video',
mimeType: 'video/VP8',
clockRate: 90000,
parameters: { 'x-google-start-bitrate': 1000 },
},
{
kind: 'video',
mimeType: 'video/H264',
clockRate: 90000,
parameters: {
'packetization-mode': 1,
'profile-level-id': '42e01f',
'level-asymmetry-allowed': 1,
},
},
];
// ---------------- Mediasoup 全局资源 ----------------
let worker;
let router;
// peerspeerId → { ws, transports: Map, producers: Map, consumers: Map }
const peers = new Map();
async function initMediasoup() {
worker = await mediasoup.createWorker({
rtcMinPort: RTC_MIN_PORT,
rtcMaxPort: RTC_MAX_PORT,
logLevel: 'warn',
});
worker.on('died', (err) => {
logger.error({ err }, 'mediasoup worker died, exiting');
process.exit(1);
});
router = await worker.createRouter({ mediaCodecs });
logger.info(
{
workerPid: worker.pid,
routerId: router.id,
rtcPortRange: `${RTC_MIN_PORT}-${RTC_MAX_PORT}`,
listenIp: LISTEN_IP,
announcedIp: ANNOUNCED_IP || '(unset, use local LAN ip)',
},
'mediasoup worker + router ready',
);
}
// ---------------- Transport 创建辅助 ----------------
async function createWebRtcTransport() {
const listenIps = [{ ip: LISTEN_IP, announcedIp: ANNOUNCED_IP || undefined }];
const transport = await router.createWebRtcTransport({
listenIps,
enableUdp: true,
enableTcp: true,
preferUdp: true,
initialAvailableOutgoingBitrate: 1_000_000,
});
return transport;
}
// ---------------- WS 信令 ----------------
function send(ws, type, data = {}, reqId) {
if (ws.readyState !== 1) return;
const msg = { type, data };
if (reqId) msg.reqId = reqId;
ws.send(JSON.stringify(msg));
}
function broadcast(excludeId, type, data) {
for (const [pid, peer] of peers) {
if (pid === excludeId) continue;
send(peer.ws, type, data);
}
}
async function handleMessage(peerId, peer, raw) {
let msg;
try {
msg = JSON.parse(raw);
} catch {
logger.warn({ peerId, raw: String(raw).slice(0, 80) }, 'invalid json');
return;
}
const { type, data = {}, reqId } = msg;
try {
switch (type) {
case 'getRtpCapabilities': {
send(peer.ws, 'rtpCapabilities', router.rtpCapabilities, reqId);
break;
}
case 'createTransport': {
const transport = await createWebRtcTransport();
peer.transports.set(transport.id, transport);
send(
peer.ws,
'transportCreated',
{
direction: data.direction,
id: transport.id,
iceParameters: transport.iceParameters,
iceCandidates: transport.iceCandidates,
dtlsParameters: transport.dtlsParameters,
},
reqId,
);
break;
}
case 'connectTransport': {
const { transportId, dtlsParameters } = data;
const transport = peer.transports.get(transportId);
if (!transport) throw new Error(`transport ${transportId} not found`);
await transport.connect({ dtlsParameters });
send(peer.ws, 'transportConnected', { transportId }, reqId);
break;
}
case 'produce': {
const { transportId, kind, rtpParameters } = data;
const transport = peer.transports.get(transportId);
if (!transport) throw new Error(`transport ${transportId} not found`);
const producer = await transport.produce({ kind, rtpParameters });
peer.producers.set(producer.id, producer);
producer.on('transportclose', () => {
peer.producers.delete(producer.id);
});
send(peer.ws, 'produced', { producerId: producer.id, kind }, reqId);
// 广播给其他 peer
broadcast(peerId, 'newProducer', {
peerId,
producerId: producer.id,
kind,
});
logger.info({ peerId, producerId: producer.id, kind }, 'peer produced');
break;
}
case 'consume': {
const { transportId, producerId, rtpCapabilities } = data;
const transport = peer.transports.get(transportId);
if (!transport) throw new Error(`transport ${transportId} not found`);
if (!router.canConsume({ producerId, rtpCapabilities })) {
throw new Error(`router cannot consume producer ${producerId}`);
}
const consumer = await transport.consume({
producerId,
rtpCapabilities,
paused: true, // 先暂停,客户端确认后再 resume
});
peer.consumers.set(consumer.id, consumer);
consumer.on('transportclose', () => {
peer.consumers.delete(consumer.id);
});
consumer.on('producerclose', () => {
peer.consumers.delete(consumer.id);
send(peer.ws, 'consumerClosed', { consumerId: consumer.id });
});
send(
peer.ws,
'consumed',
{
id: consumer.id,
producerId,
kind: consumer.kind,
rtpParameters: consumer.rtpParameters,
},
reqId,
);
break;
}
case 'resumeConsumer': {
const { consumerId } = data;
const consumer = peer.consumers.get(consumerId);
if (!consumer) throw new Error(`consumer ${consumerId} not found`);
await consumer.resume();
send(peer.ws, 'consumerResumed', { consumerId }, reqId);
break;
}
case 'getPeers': {
// 返回当前其他 peers 正在 produce 的 producerId 列表
const others = [];
for (const [pid, p] of peers) {
if (pid === peerId) continue;
for (const producer of p.producers.values()) {
others.push({ peerId: pid, producerId: producer.id, kind: producer.kind });
}
}
send(peer.ws, 'peers', { items: others }, reqId);
break;
}
default:
send(peer.ws, 'error', { message: `unknown type: ${type}` }, reqId);
}
} catch (err) {
logger.error({ peerId, type, err: err.message }, 'handle message failed');
send(peer.ws, 'error', { message: err.message }, reqId);
}
}
function cleanupPeer(peerId) {
const peer = peers.get(peerId);
if (!peer) return;
for (const p of peer.producers.values()) p.close();
for (const c of peer.consumers.values()) c.close();
for (const t of peer.transports.values()) t.close();
peers.delete(peerId);
broadcast(peerId, 'peerLeft', { peerId });
logger.info(
{
peerId,
remaining: peers.size,
routerStats: { producers: [...peers.values()].reduce((n, p) => n + p.producers.size, 0) },
},
'peer cleaned up',
);
}
// ---------------- Fastify 启动 ----------------
async function main() {
await initMediasoup();
const app = Fastify({ logger: false });
await app.register(fastifyWebsocket);
await app.register(fastifyStatic, {
root: path.join(__dirname, 'public'),
prefix: '/',
});
app.get('/healthz', async () => ({
ok: true,
workerPid: worker.pid,
routerId: router.id,
peers: peers.size,
}));
app.get('/stats', async () => {
// 基础压测观测peer 数 + producer/consumer 总数 + 内存
let totalProducers = 0;
let totalConsumers = 0;
let totalTransports = 0;
for (const p of peers.values()) {
totalProducers += p.producers.size;
totalConsumers += p.consumers.size;
totalTransports += p.transports.size;
}
const mem = process.memoryUsage();
return {
peers: peers.size,
transports: totalTransports,
producers: totalProducers,
consumers: totalConsumers,
memoryMB: {
rss: Math.round(mem.rss / 1024 / 1024),
heapUsed: Math.round(mem.heapUsed / 1024 / 1024),
},
loadavg: os.loadavg(),
uptime: Math.round(process.uptime()),
};
});
app.register(async function (f) {
f.get('/ws', { websocket: true }, (socket, req) => {
const peerId = randomUUID();
const peer = {
ws: socket,
transports: new Map(),
producers: new Map(),
consumers: new Map(),
};
peers.set(peerId, peer);
logger.info({ peerId, total: peers.size, ua: req.headers['user-agent'] }, 'peer connected');
send(socket, 'welcome', { peerId });
socket.on('message', (raw) => {
handleMessage(peerId, peer, raw.toString());
});
socket.on('close', () => {
logger.info({ peerId }, 'ws closed');
cleanupPeer(peerId);
});
socket.on('error', (err) => {
logger.warn({ peerId, err: err.message }, 'ws error');
});
});
});
await app.listen({ port: HTTP_PORT, host: '0.0.0.0' });
logger.info(`PoC ready: http://localhost:${HTTP_PORT} (open in 2 browser windows)`);
}
main().catch((err) => {
logger.error({ err }, 'poc server failed to start');
process.exit(1);
});