Files
EchoChat/media-server/poc/public/client.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

348 lines
11 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.

// PoC 浏览器端mediasoup-client + 原生 WebSocket
// 通过 esm.sh CDN 加载 mediasoup-client保持 PoC 零构建步骤。
// 正式 media-server 落地时会改为 frontend 侧 npm 依赖 + vite 打包。
// esm.sh 自动解析到最新 3.x本 PoC 验证时为 3.19.0);正式 media-server 落地会用 npm + vite 打包
import { Device } from 'https://esm.sh/mediasoup-client@3';
const logEl = document.getElementById('log');
const statusEl = document.getElementById('status');
const videosEl = document.getElementById('videos');
const btnJoin = document.getElementById('btnJoin');
const btnLeave = document.getElementById('btnLeave');
function log(msg, level = 'info') {
const ts = new Date().toLocaleTimeString();
const line = document.createElement('div');
line.className = `log-${level}`;
line.textContent = `[${ts}] ${msg}`;
logEl.appendChild(line);
logEl.scrollTop = logEl.scrollHeight;
console[level === 'err' ? 'error' : 'log'](msg);
}
// ---------------- WS 信令包装Promise-based ----------------
class SignalClient {
constructor(url) {
this.url = url;
this.ws = null;
this.pending = new Map(); // reqId → {resolve, reject}
this.listeners = new Map();
}
connect() {
return new Promise((resolve, reject) => {
this.ws = new WebSocket(this.url);
this.ws.onopen = () => resolve();
this.ws.onerror = (e) => reject(e);
this.ws.onclose = () => {
this.emit('__close__', {});
};
this.ws.onmessage = (ev) => this._onMessage(ev);
});
}
_onMessage(ev) {
let msg;
try { msg = JSON.parse(ev.data); } catch { return; }
const { type, data, reqId } = msg;
if (reqId && this.pending.has(reqId)) {
const { resolve, reject } = this.pending.get(reqId);
this.pending.delete(reqId);
if (type === 'error') reject(new Error(data?.message || 'signal error'));
else resolve(data);
return;
}
this.emit(type, data);
}
request(type, data = {}) {
const reqId = Math.random().toString(36).slice(2);
return new Promise((resolve, reject) => {
this.pending.set(reqId, { resolve, reject });
this.ws.send(JSON.stringify({ type, data, reqId }));
setTimeout(() => {
if (this.pending.has(reqId)) {
this.pending.delete(reqId);
reject(new Error(`request ${type} timeout`));
}
}, 10_000);
});
}
on(type, fn) {
if (!this.listeners.has(type)) this.listeners.set(type, new Set());
this.listeners.get(type).add(fn);
}
emit(type, data) {
const set = this.listeners.get(type);
if (set) for (const fn of set) fn(data);
}
close() {
if (this.ws) this.ws.close();
}
}
// ---------------- App 状态 ----------------
const state = {
signal: null,
device: null,
sendTransport: null,
recvTransport: null,
myPeerId: null,
localStream: null,
producers: { audio: null, video: null },
consumers: new Map(), // producerId → {consumer, peerId}
};
function setStatus(text) {
statusEl.textContent = text;
}
function addVideoTile(id, stream, label, isSelf = false) {
let box = document.getElementById(`tile-${id}`);
if (!box) {
box = document.createElement('div');
box.id = `tile-${id}`;
box.className = 'video-box' + (isSelf ? ' self' : '');
const v = document.createElement('video');
v.autoplay = true;
v.playsInline = true;
if (isSelf) v.muted = true;
v.srcObject = stream;
box.appendChild(v);
const lab = document.createElement('div');
lab.className = 'label';
lab.textContent = label;
box.appendChild(lab);
videosEl.appendChild(box);
} else {
const v = box.querySelector('video');
// 将新的 track 合并进已有 stream
for (const t of stream.getTracks()) {
if (!v.srcObject) v.srcObject = new MediaStream();
v.srcObject.addTrack(t);
}
}
}
function removeVideoTile(id) {
const box = document.getElementById(`tile-${id}`);
if (box) box.remove();
}
// ---------------- 关键流程 ----------------
async function join() {
btnJoin.disabled = true;
setStatus('connecting WS...');
const proto = location.protocol === 'https:' ? 'wss' : 'ws';
const signal = new SignalClient(`${proto}://${location.host}/ws`);
state.signal = signal;
signal.on('welcome', ({ peerId }) => {
state.myPeerId = peerId;
log(`WS connected as peer ${peerId.slice(0, 8)}`, 'ok');
});
signal.on('newProducer', async ({ peerId, producerId, kind }) => {
log(`remote producer: peer=${peerId.slice(0, 8)} kind=${kind}`, 'info');
await consumeRemote(peerId, producerId, kind);
});
signal.on('peerLeft', ({ peerId }) => {
log(`peer left: ${peerId.slice(0, 8)}`, 'info');
removeVideoTile(peerId);
// 清理该 peer 的 consumers
for (const [pid, entry] of state.consumers) {
if (entry.peerId === peerId) {
entry.consumer.close();
state.consumers.delete(pid);
}
}
});
signal.on('consumerClosed', ({ consumerId }) => {
for (const [pid, entry] of state.consumers) {
if (entry.consumer.id === consumerId) {
entry.consumer.close();
state.consumers.delete(pid);
}
}
});
signal.on('__close__', () => {
setStatus('disconnected');
log('WS closed', 'err');
});
await signal.connect();
// 1. Device
setStatus('loading device...');
const rtpCapabilities = await signal.request('getRtpCapabilities');
const device = new Device();
await device.load({ routerRtpCapabilities: rtpCapabilities });
state.device = device;
log('mediasoup Device loaded', 'ok');
// 2. getUserMedia
setStatus('requesting media...');
let localStream;
try {
localStream = await navigator.mediaDevices.getUserMedia({
audio: true,
video: { width: { ideal: 640 }, height: { ideal: 360 } },
});
} catch (e) {
log(`getUserMedia failed: ${e.message}`, 'err');
setStatus('media denied');
btnJoin.disabled = false;
return;
}
state.localStream = localStream;
addVideoTile('local', localStream, '本地(你自己)', true);
// 3. SendTransport
setStatus('creating send transport...');
const sendInfo = await signal.request('createTransport', { direction: 'send' });
const sendTransport = device.createSendTransport({
id: sendInfo.id,
iceParameters: sendInfo.iceParameters,
iceCandidates: sendInfo.iceCandidates,
dtlsParameters: sendInfo.dtlsParameters,
});
state.sendTransport = sendTransport;
sendTransport.on('connect', ({ dtlsParameters }, cb, errb) => {
signal
.request('connectTransport', { transportId: sendTransport.id, dtlsParameters })
.then(() => cb())
.catch(errb);
});
sendTransport.on('produce', ({ kind, rtpParameters }, cb, errb) => {
signal
.request('produce', { transportId: sendTransport.id, kind, rtpParameters })
.then(({ producerId }) => cb({ id: producerId }))
.catch(errb);
});
sendTransport.on('connectionstatechange', (s) => {
log(`sendTransport ICE: ${s}`, s === 'failed' ? 'err' : 'info');
});
// 4. produce audio + video
const audioTrack = localStream.getAudioTracks()[0];
const videoTrack = localStream.getVideoTracks()[0];
if (audioTrack) {
state.producers.audio = await sendTransport.produce({ track: audioTrack });
log(`local audio producer: ${state.producers.audio.id.slice(0, 8)}`, 'ok');
}
if (videoTrack) {
state.producers.video = await sendTransport.produce({
track: videoTrack,
encodings: [
{ maxBitrate: 150_000, scaleResolutionDownBy: 4 },
{ maxBitrate: 400_000, scaleResolutionDownBy: 2 },
{ maxBitrate: 1_000_000 },
],
codecOptions: { videoGoogleStartBitrate: 1000 },
});
log(`local video producer: ${state.producers.video.id.slice(0, 8)}`, 'ok');
}
// 5. RecvTransport
setStatus('creating recv transport...');
const recvInfo = await signal.request('createTransport', { direction: 'recv' });
const recvTransport = device.createRecvTransport({
id: recvInfo.id,
iceParameters: recvInfo.iceParameters,
iceCandidates: recvInfo.iceCandidates,
dtlsParameters: recvInfo.dtlsParameters,
});
state.recvTransport = recvTransport;
recvTransport.on('connect', ({ dtlsParameters }, cb, errb) => {
signal
.request('connectTransport', { transportId: recvTransport.id, dtlsParameters })
.then(() => cb())
.catch(errb);
});
recvTransport.on('connectionstatechange', (s) => {
log(`recvTransport ICE: ${s}`, s === 'failed' ? 'err' : 'info');
});
// 6. 拉取当前已在房间的其他 peers
const { items } = await signal.request('getPeers');
log(`existing remote producers: ${items.length}`, 'info');
for (const it of items) {
await consumeRemote(it.peerId, it.producerId, it.kind);
}
setStatus('in room');
btnLeave.disabled = false;
}
async function consumeRemote(peerId, producerId, kind) {
if (!state.recvTransport) return;
try {
const data = await state.signal.request('consume', {
transportId: state.recvTransport.id,
producerId,
rtpCapabilities: state.device.rtpCapabilities,
});
const consumer = await state.recvTransport.consume({
id: data.id,
producerId: data.producerId,
kind: data.kind,
rtpParameters: data.rtpParameters,
});
await state.signal.request('resumeConsumer', { consumerId: consumer.id });
state.consumers.set(producerId, { consumer, peerId });
const stream = new MediaStream([consumer.track]);
addVideoTile(peerId, stream, `远端 ${peerId.slice(0, 8)} · ${kind}`);
log(`consuming ${kind} from ${peerId.slice(0, 8)}`, 'ok');
consumer.on('transportclose', () => state.consumers.delete(producerId));
} catch (e) {
log(`consume failed: ${e.message}`, 'err');
}
}
async function leave() {
btnLeave.disabled = true;
setStatus('leaving...');
try {
for (const { consumer } of state.consumers.values()) consumer.close();
state.consumers.clear();
if (state.producers.audio) state.producers.audio.close();
if (state.producers.video) state.producers.video.close();
if (state.sendTransport) state.sendTransport.close();
if (state.recvTransport) state.recvTransport.close();
if (state.localStream) state.localStream.getTracks().forEach((t) => t.stop());
if (state.signal) state.signal.close();
} catch (e) {
log(`leave error: ${e.message}`, 'err');
}
videosEl.innerHTML = '';
setStatus('idle');
btnJoin.disabled = false;
}
btnJoin.addEventListener('click', () => {
join().catch((e) => {
log(`join error: ${e.message}`, 'err');
setStatus('error');
btnJoin.disabled = false;
});
});
btnLeave.addEventListener('click', () => leave());
log('PoC client ready. 点击「加入并推流」,授权摄像头/麦克风后可与其他窗口互通。', 'info');