Node.js net 模块练习答案(14~27)

对应 net-exercises.md 的第 14~27 题;上一册见 第 1~13 题答案。所有示例面向 Node.js 22+ 与 TypeScript。

14. 使用 pipeline() 转发 TCP 数据

一个 TCP 代理有两条独立方向:A→B、B→A。pipeline() 会传播背压,任一流错误时销毁该方向参与的流;但仍要统一关闭另一端。

import net from "node:net";
import { pipeline } from "node:stream";

const proxy = net.createServer((client) => {
  const upstream = net.createConnection({ host: "127.0.0.1", port: 9100 });
  let settled = false;
  const stopBoth = (error?: Error) => {
    if (settled) return;
    settled = true;
    client.destroy(error);
    upstream.destroy(error);
  };

  pipeline(client, upstream, (error) => {
    if (error) stopBoth(error);
    // client 的 EOF 会调用 upstream.end(),允许其正常返回响应。
  });
  pipeline(upstream, client, (error) => {
    if (error) stopBoth(error);
  });
  client.on("error", () => undefined);
  upstream.on("error", () => undefined);
  client.on("close", () => stopBoth());
  upstream.on("close", () => stopBoth());
});
proxy.listen(9000);

上例是教学最小版:真实代理还要处理“连接上游前客户端已断开”、连接/首字节超时、访问控制、连接计数和日志。不要写 client.pipe(upstream); upstream.pipe(client) 后忽略错误;它虽然会处理正常背压,却没有统一的错误与另一端清理策略。

15. 每连接累计接收限制

const MAX_RECEIVED = 5 * 1024 * 1024;
const server = net.createServer((socket) => {
  let received = 0;
  socket.on("data", (chunk: Buffer) => {
    received += chunk.length;
    if (received > MAX_RECEIVED) {
      socket.end("ERR connection byte quota exceeded\\n");
      return;
    }
    // 处理 chunk/协议帧
  });
  socket.on("error", () => undefined);
});
server.listen(9000);

单条消息限制限制一个 frame 的大小,适用于 JSON、文件元数据等;累计限制限制整个连接会话的带宽/工作量。聊天室等长期连接通常不能用永久 5MB 总额,应改成按时间窗口限速、单帧上限和业务配额。

16. 最大连接数

server.maxConnections 是 Node 的上限提示;为了能发送繁忙提示和准确学习计数,显式管理:

const LIMIT = 100;
let connections = 0;
const server = net.createServer((socket) => {
  if (connections >= LIMIT) {
    socket.end("BUSY\\n");
    return;
  }
  connections += 1;
  let released = false;
  socket.once("close", () => {
    if (!released) {
      released = true;
      connections -= 1;
    }
  });
  socket.on("error", () => undefined);
});
server.maxConnections = LIMIT;
server.listen(9000);

计数必须在 close 释放而不是 endend 只代表读端结束,异常断连未必触发它。多进程/Cluster 下每个进程的内存计数彼此独立,需要由负载均衡器、共享限流器或进程总量策略管理全局限制。

17. 协议级心跳

const HEARTBEAT_MS = 20_000;
const server = net.createServer((socket) => {
  let lastPingAt = Date.now();
  const check = setInterval(() => {
    if (Date.now() - lastPingAt > HEARTBEAT_MS) socket.end("ERR heartbeat timeout\\n");
  }, 5_000);
  socket.once("close", () => clearInterval(check));
  socket.on("data", (chunk) => {
    // 实际代码应先按第 8/10 题完成分帧
    if (chunk.toString("utf8") === "PING\\n") {
      lastPingAt = Date.now();
      socket.write("PONG\\n");
    }
  });
  socket.on("error", () => undefined);
});

实际中只要收到完整且通过验证的协议消息就可更新活动时间,不能假定一个 chunk 等于 PING\n。TCP 连接存在不等于应用仍健康:网络设备可能保留半死连接,客户端也可能卡死。socket.setTimeout() 只看任意 I/O 空闲,不理解 PING/PONG;应用心跳能验证协议对端仍在工作,通常二者配合。

18. 有限重连与抖动

import net from "node:net";

const MAX_ATTEMPTS = 6;
let attempt = 0;
let stopped = false;
let current: net.Socket | undefined;

function connect(): void {
  if (stopped) return;
  const socket = net.createConnection({ host: "127.0.0.1", port: 9000 });
  current = socket;
  let connected = false;
  socket.once("connect", () => { connected = true; attempt = 0; });
  socket.on("error", (error) => console.warn("socket", error.code));
  socket.once("close", () => {
    if (stopped || connected && attempt >= MAX_ATTEMPTS) return;
    attempt += 1;
    if (attempt > MAX_ATTEMPTS) return void console.error("giving up");
    const base = Math.min(1_000 * 2 ** (attempt - 1), 30_000);
    const delay = Math.round(base * (0.5 + Math.random())); // full-ish jitter range
    setTimeout(connect, delay).unref();
  });
}
connect();
process.on("SIGINT", () => { stopped = true; current?.destroy(); });

示例中“成功连接后又断开”重新从短退避开始;若要限制整个进程的重连总时长,单独记录 firstFailureAt。抖动让大量客户端不在 1、2、4 秒整齐地同时冲击刚恢复的服务。重连前要确保旧 Socket 已 close,并区分用户主动关闭与意外断线。

19. 监听地址

server.listen(9000, "127.0.0.1"); // 仅本机 IPv4
// server.listen(9000, "0.0.0.0"); // 所有 IPv4 网卡,可能对公网可达
// server.listen(9000, "::1");      // 仅本机 IPv6

127.0.0.1 适合仅供 Nginx、同机应用或本地开发访问;0.0.0.0 是暴露到所有 IPv4 接口,是否能从公网访问还受云安全组、防火墙、路由影响。生产默认应最小暴露:无须外部访问的管理/内部 TCP 服务监听 loopback 或 Unix Socket,不要仅依赖“端口不常见”。

20. 管理 TCP 接口不能裸露公网

CLEAR_ALL 如果没有认证,任何能连端口的人都能清缓存,造成可用性事故。最低设计是:

公网不监听该端口
→ Unix Socket 或 localhost
→ 防火墙/安全组只允许运维网段
→ mTLS 或短期、可轮换的管理凭据
→ 明确命令白名单、审计日志、速率限制和双人确认

不要自己发明“明文 token + TCP”后认为已安全:明文可被窃听/重放,=== 比较也不能解决传输安全。管理操作更适合受认证的 HTTPS 管理 API、VPN/堡垒机或云平台控制面,并遵循最小权限。

21. Unix Domain Socket

Linux/macOS 示例:

import net from "node:net";
import { rm } from "node:fs/promises";

const path = "/run/my-app/service.sock";
await rm(path, { force: true }); // 仅确认是本服务专属、陈旧的 socket 文件时
const server = net.createServer((socket) => socket.end("ok\\n"));
server.listen(path);

UDS 不经过 TCP 端口,适合同机 Nginx↔应用、进程间服务;其文件系统权限决定谁可连接,仍需正确 chown/chmod 与协议认证。它不适合跨主机。Windows 对应概念是 Named Pipe(例如 \\\\.\\pipe\\my-app);不要把 POSIX 路径假设为跨平台。程序退出后 socket 文件可能残留,启动时清理必须限定为已验证的专属路径。

22. 优雅关闭 TCP 服务

import net from "node:net";

const sockets = new Set<net.Socket>();
const server = net.createServer((socket) => {
  sockets.add(socket);
  socket.once("close", () => sockets.delete(socket));
  socket.on("error", () => undefined);
});
server.listen(9000);

function shutdown(signal: string): void {
  console.info(`${signal}: stop accepting`);
  server.close((error) => {
    if (error) console.error(error);
    process.exitCode = error ? 1 : 0;
  });
  for (const socket of sockets) socket.end("SERVER_SHUTTING_DOWN\\n");
  const forceTimer = setTimeout(() => {
    for (const socket of sockets) socket.destroy();
  }, 10_000);
  forceTimer.unref();
}
process.once("SIGINT", () => shutdown("SIGINT"));
process.once("SIGTERM", () => shutdown("SIGTERM"));

server.close() 只停止接受新连接,等待现有连接关闭;不主动结束长期连接。socket.end() 是优雅请求,超时后的 destroy() 才是强制兜底。真实代码应防止两个 signal 重复进入 shutdown,并在退出前关闭数据库、队列等资源;不要在 signal 回调开头 process.exit(),那会截断在途写入。

23. 连接日志

import { randomUUID } from "node:crypto";

const server = net.createServer((socket) => {
  const id = randomUUID();
  const openedAt = Date.now();
  let received = 0;
  let sent = 0;
  let reason = "normal";
  socket.on("data", (chunk) => { received += chunk.length; });
  socket.on("error", (error) => { reason = `error:${error.code ?? error.name}`; });
  const send = (text: string): boolean => {
    sent += Buffer.byteLength(text, "utf8");
    return socket.write(text);
  };
  send("WELCOME\\n");
  socket.once("close", (hadError) => console.info({ id, remote: socket.remoteAddress, received, sent, ms: Date.now() - openedAt, reason, hadError }));
  console.info({ id, event: "connected" });
});

工程中更好的方式是在自己所有 send() 包装函数统计 sent,而不是覆盖 Socket 方法。TCP connection ID 覆盖连接生命周期;HTTP request ID 覆盖一条请求。HTTP keep-alive 下一个 TCP ID 可关联多个 request ID。日志勿记录原始敏感 payload。

24. 原始 TCP 发送 HTTP

import net from "node:net";

const socket = net.createConnection({ host: "127.0.0.1", port: 8080 });
const parts: Buffer[] = [];
socket.on("connect", () => socket.write(
  "GET /health HTTP/1.1\\r\\nHost: localhost:8080\\r\\nConnection: close\\r\\n\\r\\n",
));
socket.on("data", (chunk) => parts.push(chunk));
socket.on("end", () => console.log(Buffer.concat(parts).toString("utf8")));
socket.on("error", console.error);

响应是 HTTP/1.1 200 OK 状态行、若干 Name: value 响应头、空行 \r\n\r\n、响应体。HTTP/1.1 语法规定行结束为 CRLF(\r\n);宽容服务端未必能掩盖不合规客户端。这里 Connection: close 令 EOF 可作为响应结束标记;chunked、Content-Length、压缩时解析更复杂。

25. 分块发送 HTTP 请求

socket.on("connect", () => {
  socket.write("GET /hea");
  setTimeout(() => socket.write("lth HTTP/1.1\\r\\nHost: localhost:8080\\r\\n"), 10);
  setTimeout(() => socket.end("Connection: close\\r\\n\\r\\n"), 20);
});

Node http 解析器会把 TCP 字节流重组成请求行、headers 和 body;应用收到的是已解析的 IncomingMessage。若自己实现 HTTP,至少要处理 CRLF、header 大小/数量、Content-Length、chunked 编码、多个 keep-alive 请求、请求 body 流、超时、畸形语法、升级协议与安全限制。不要用 data.toString().split("\\r\\n\\r\\n") 假装实现 HTTP。

26. HTTP keep-alive 复用

import http from "node:http";

let connections = 0;
const server = http.createServer((_req, res) => res.end("ok"));
server.on("connection", (socket) => {
  connections += 1;
  console.log("connection +", connections);
  socket.once("close", () => console.log("connection -", --connections));
});
server.listen(8080);

同一个 net.Socket 上可顺序发送两个完整 HTTP/1.1 请求(本示例须正确解析每个响应边界后再发下一条,避免把响应混在一起):

GET /health HTTP/1.1\r\nHost: localhost:8080\r\n\r\n
GET /health HTTP/1.1\r\nHost: localhost:8080\r\nConnection: close\r\n\r\n

keep-alive 减少 TCP/TLS 建连延迟和短连接端口压力,但闲置 Socket 占 FD/内存,必须有合理 keepAliveTimeout、最大连接数和负载均衡策略。短连接释放快但握手成本高。Node HTTP 默认也会依据请求/响应 header 决定是否保持连接。

27. HTTP 响应背压

ServerResponse 是可写流,规则与 net.Socket.write() 相同:

import http from "node:http";

const server = http.createServer((_req, res) => {
  res.writeHead(200, { "content-type": "text/plain; charset=utf-8" });
  let index = 0;
  const writeMore = (): void => {
    while (index < 100_000) {
      if (!res.write(`line ${index++}\\n`)) {
        res.once("drain", writeMore);
        return;
      }
    }
    res.end();
  };
  writeMore();
});
server.listen(8080);

false 不代表客户端已经断开;它说明可写缓冲达到阈值。drain 才能恢复生产,同时仍应监听请求/响应的 close,客户端断开后停止生成数据。对文件、压缩、代理等流式任务优先 pipeline(readable, transform, res, callback),它会把下游背压传给上游,并让错误在一个回调中收束。