第 0 章

全景:一次语音会话的完整旅程

先看全貌,再看代码。

voice_chat 很小——后端 9 个项目模块、前端 3 个页面文件和 1 个音频资源——但它把 Elixir 里最重要的几件事串了起来:进程、消息传递、OTP 监督树、Plug 管线、WebSocket 协议。本章用一张会动的图,带你走一遍「你对着麦克风说话」到「服务端落盘 + 约每 3 秒推送一次提示文本与音频」的全程。

Bandit 服务器 · BEAM 运行时(Elixir 一侧) 浏览器 · 麦克风 getUserMedia 采集 PCM16 / Opus 分块 (第 9 章) WSHandler 会话进程 init / handle_in / handle_info 每个连接一个(第 5 章) 运行在 Bandit 为连接开的进程里 Recorder 进程 收 chunk → 写入文件(第 6 章) Prompter 进程 每 3 秒推一次提示音(第 7 章) 录音文件 priv/recordings WAV · WebM · Ogg 浏览器 · 日志 / 喇叭 handleText + decodeAudioData 显示提示并播放音频(第 9 章) ① 连接 + config ② 音频块(二进制帧) ③ 持续写入文件 ④a 文本 + 音频消息 ④b WebSocket 下行 spawn_link ×2
图 0-1 · 一次语音会话的旅程:① 浏览器连接并发送 config;② 音频块不断上行;③ Recorder 持续落盘;④ Prompter 约每 3 秒把文本与音频消息发给 ws 进程,再由该连接推到服务器边界外的浏览器显示、播放。

一张表记住「需求 → 代码 → 章节」

你的需求落在哪对应章节
启动 WebSocket 接口endpoint.ex 的 /ws 路由 + ws_handler.ex第 4、5 章
连接 = 新建一个 session每个连接一个 WSHandler 进程,注册进 Registry第 5 章
持续接收麦克风音频(PCM16 / Opus 容器分块)二进制帧 → handle_in → Recorder 进程第 5、6 章
轻量进程把音频流保存为文件recorder.ex(spawn_link + receive 循环)第 6 章
另一个轻量进程约每 3 秒推送提示prompter.ex + 现有 haode.wav(缺失时由 wav.ex 合成双音)第 7、8 章
前端调试页(消息日志 + 麦克风采集/回放)priv/static/index.html / app.js / app.css第 9 章
为什么“每连接两个进程”会是这个项目的灵魂

Bandit 在 BEAM 进程中管理连接,本项目再为每个 WebSocket 会话显式 spawn 两个帮手进程,用消息把连接处理、落盘和定时提示分开。各会话进程的堆和业务状态不直接共享;Registry 则封装了共享查找状态。业务代码也没有自行管理线程池;真正的 socket、调度与协议细节由 Bandit、Thousand Island 和 BEAM 运行时承担。进程和消息是 Elixir 的核心并发原语,voice_chat 把它们放进了一个直观的小场景。

学完这一章你应该能做到

  • 在 30 秒内说出「你说话 → 落盘 → 收到提示音」经过哪几个进程、什么顺序
  • 看着任意一行后端代码,说出它运行在哪个进程里
  • 解释为什么「每 3 秒发一次」不需要引入任何定时器库
  • 指出 4 个 OTP 概念在这个项目里的具体化身:Application / Supervisor / 进程 / 消息
  • 说清 Plug 管线和 WebSocket 升级在 /ws 这条路上各做了什么
自测:给这段旅程排个序

把下面 6 步按真实发生顺序排列:
a) Prompter 进程向 ws 进程发 {:send_binary, 音频} b) 浏览器 getUserMedia 拿到麦克风  c) Bandit 为连接开进程并调用 WSHandler.init/1 d) Recorder 进程把 chunk 追加进文件  e) 浏览器 decodeAudioData 播放“好的” f) terminate 停掉两个子进程

答案:HTTP Upgrade 成功后,服务端会切换协议并执行 c,浏览器的 open 回调则开始 b;这两边属于并发的握手后工作,不能要求 c 一定先完成。服务端在处理后续帧前会先完成 init。随后音频块上行并触发 d;Prompter 从 init 创建后约 3 秒执行 a,浏览器收到音频后执行 e。断开时进入 f。麦克风授权、持续落盘和周期提示会交错进行,因此不存在把 a、d 排成唯一全局顺序的答案。
第 1 章

mix new:一个 Elixir 项目从哪来

先认识你脚下的脚手架,再看每一块地砖。

整个项目从 mix new voice_chat --sup 提供的基础骨架长出来。本章先区分哪些文件由命令生成、哪些是项目后续加入,再搞清楚 lib/、priv/、config/、test/、mix.exs、mix.lock 各由谁在什么时候使用。

基础骨架 vs 后续新增

# 当前项目的教学摘录(省略 README、scripts、部分测试与构建产物) voice_chat/ ├── mix.exs # mix new 生成,后续补入依赖和应用配置 ├── mix.lock # 首次解析依赖后生成,不是 mix new 的初始文件 ├── .formatter.exs # mix format 的规则 ├── .gitignore # 忽略 _build/ deps/ 等 ├── config/ # 本项目后续加入 │ ├── config.exs # 共享配置入口 │ └── {dev,prod,test}.exs # 分环境配置 ├── lib/ │ ├── voice_chat.ex # 顶层模块(我们把它改成文档 + 两个辅助函数) │ └── voice_chat/ # Application + 7 个业务模块 ├── priv/static/ # 本项目后续加入:页面、脚本、样式、提示音 ├── priv/recordings/ # 运行时创建的录音输出目录 └── test/ ├── test_helper.exs └── *_test.exs # ExUnit 单元/生命周期测试

mix new --sup 提供 mix.exs、顶层模块、Application 骨架和测试等基础文件;本项目随后新增了 config/、priv/ 以及七个业务模块,并改写了 Application、顶层模块和测试。mix.lock 则在依赖被解析后生成。

六类文件,各用一句人话说明

目录 / 文件一句话谁在消费它
mix.exs项目的「出生证明」:名字、版本、依赖、应用入口mix 本身(deps.get / compile / run / test…)
lib/按惯例放应用的 Elixir 源码编译器 → _build/ 里的 .beam 字节码
priv/随应用打包、运行期按路径读取的资源(前端页、音频)Application.app_dir/2、Plug.Static
config/编译或启动应用前加载的配置(端口、录音目录)Mix 配置系统 → Application.get_env/3
test/ExUnit 测试mix test
mix.lock记录已解析依赖的精确版本与校验信息mix deps.get(让团队/CI 复用同一解析结果)
类比:mix 之于 Elixir,约等于 cargo / npm / rebar3

mix 身兼三职:项目脚手架(new)、构建工具(compile / test / run)、包管理器入口(deps.get 去 Hex 拉依赖,deps 锁进 mix.lock)。它采用强约定:默认会从 lib/ 编译源码、从 test/ 找测试;这些路径可配置,但多数项目沿用默认布局。

一个关键区分:Mix 项目 ≠ OTP 应用

这两个词经常被混用,但拆开讲更清楚:

  • Mix 项目(project):构建工具看到的工程,由 mix.exs 描述——怎么编译、测试、解析依赖和打包。
  • OTP 应用(application):mix.exs 里 app: :voice_chat 声明的那个东西,服务于运行期——一组模块 + 一个可选的启动入口(mod:),由 BEAM 的应用控制器负责启动/停止。

你写 mix run --no-halt 时,mix 先编译项目,再让应用控制器把 :voice_chat 应用连同它的依赖(logger、bandit…)一起启动——启动时调用 mod: 指向的 VoiceChat.Application.start/2。这条链到第 3 章会完整展开。

自测:lib/ 和 priv/ 有什么区别

为什么本项目把前端页面放在 priv/static/,而不是 lib/?

答案:lib/ 默认用于会编译成 BEAM 字节码的 Elixir 源码;priv/ 用于随应用打包、运行期按路径读取的资源。前端页面要原样交给浏览器,所以放进惯用的 priv/static/,再由 Plug.Static 或 Application.app_dir/2 定位。
第 2 章

mix.exs 逐行:项目的出生证明

全文只有 29 行,但每一行都决定了后面所有章节的写法。

这是本书第一个代码走读。走读组件左侧是带行号的源码,右侧是讲解卡片;点代码里的任意一行,或点卡片标题,或按 ▶ 逐步播放,讲解会跟着行号走。代码块可用 python3 gen.py --check 对照 voice_chat/ 源码校验。

mix.exs· 29 行 · 整个项目的元信息 0/0
defmodule VoiceChat.MixProject do
  use Mix.Project

  def project do
    [
      app: :voice_chat,
      version: "0.1.0",
      elixir: "~> 1.15",
      start_permanent: Mix.env() == :prod,
      deps: deps()
    ]
  end

  def application do
    [
      extra_applications: [:logger],
      mod: {VoiceChat.Application, []}
    ]
  end

  defp deps do
    [
      {:bandit, "~> 1.6"},
      {:websock_adapter, "~> 0.5"},
      {:plug, "~> 1.16"},
      {:jason, "~> 1.4"}
    ]
  end
end
L1–3 模块声明 + use Mix.Project

VoiceChat.MixProject 是生成器采用的命名约定;真正把模块声明为 Mix 项目的是 use Mix.Project。它设置项目行为和编译钩子,项目再实现 project/0,并可选实现 application/0。

Mix 在求值 mix.exs 时装载这个模块并读取返回的配置;关键不是模块名被“搜索到”,而是 use Mix.Project 完成注册和约定接线。

关联:第 1 章「Mix 项目」、第 3 章「OTP 应用」——这一行同时定义了两者。

L4–5 project/0:一段关键字列表

project/0 返回一个关键字列表(keyword list),是 mix 读取项目信息的主入口。第 5 行的 [ 开始这个列表。

这里体现「配置就是数据」:本项目的函数只返回关键字列表。Elixir 函数本身并不会被语言强制为纯函数,因此“无副作用”是这段实现的选择,不是 Mix 的硬保证。

L6 app: :voice_chat —— OTP 应用名

声明这个项目的 OTP 应用名(一个 atom)。它出现在三处:编译产物 voice_chat.app、Application.get_env(:voice_chat, …)、以及 Application.app_dir(:voice_chat, …)。

atom 是 Erlang VM 里的符号常量,适合应用名这类有限、可信的标识;atom 表不会回收,所以不能把任意用户输入转成新 atom。

关联:第 3 章 Application.get_env,以及第 4、8 章的 Application.app_dir——同一枚应用 atom 贯穿配置与资源定位。

L7 version: "0.1.0"

当前项目自身的版本号,会进入应用元数据并用于发布/打包;发布到 Hex 时通常遵循语义化版本。它不决定依赖锁定,mix.lock 记录的是依赖解析结果。

L8 elixir: "~> 1.15" —— 版本约束

要求 Elixir 版本。重点是 ~> 的语义:~> 1.15 允许 1.15.x 及以上、但不到 2.0 的版本(即允许 1.19,不允许 2.0)。

用宽松下限 + 语义化版本,既不锁死环境,又不放进来不兼容的大版本——这是 Hex 生态的惯例。

关联:第 1 章 mix.lock(依赖的版本锁得比这严格得多)。

L9 start_permanent: Mix.env() == :prod

当 Mix 启动应用且 MIX_ENV=prod 时,把该应用标为 :permanent;如果整个 OTP 应用最终停止,BEAM 节点会退出。子进程的普通崩溃仍先按 Supervisor 策略重启,并不是任意异常都立即关 VM。

它让“应用已无法维持监督树”成为节点级失败信号,便于外部服务管理器重启实例;开发环境不启用这个永久启动标记。

L10–13 deps: deps() —— 依赖

把私有函数 deps/0 返回的依赖列表挂到 project 上。deps/0 定义在第 21 行。

拆成私有函数是惯例:project 保持可读,依赖单独成块方便增删。

L14–16 application/0:运行期怎么启动

extra_applications: [:logger] 声明「我的应用启动前,请先把 :logger 应用启动好」——因为代码里要 Logger.info。

BEAM 的应用有依赖顺序:Logger 是 Elixir 自带的库应用,也要显式声明才保证可用。这也是为什么依赖图是「应用」,不只是「包」。

L17–20 mod: {VoiceChat.Application, []} —— 应用入口

告诉应用控制器:启动本应用时调用 VoiceChat.Application.start/2,参数为 []。若回调模块实现可选的 stop/1,应用停止时才会调用;本项目没有实现它。

没有 mod: 的应用不会启动自己的 Application 回调(仍可打包模块、配置和依赖);有 mod: 时应用控制器会调用入口,通常由它拉起进程树。voice_chat 属于后一种。

关联:第 3 章整章——start/2 里就是监督树。

L21–22 deps/0:声明式依赖

返回依赖元组列表。每个元组 {名字, 版本约束} 只是声明——真正下载、编译发生在 mix deps.get / mix compile 时。

声明式的好处是机器可读、可锁定(mix.lock)并容易在团队和 CI 中复现;Mix/Hex 负责常规依赖获取与校验。

L23–28 四个依赖,各司其职

bandit:HTTP/1、HTTP/2 与 WebSocket 服务器;websock_adapter:把 Plug 连接升级接到 WebSock 回调;plug:请求管线/路由抽象;jason:JSON 编解码。

Bandit 提供底层 socket 与协议,WebSock 定义「回调契约」,Plug 定义「请求管线」——三者互相独立、靠 adapter 拼起来,这正是 BEAM 生态「小库组合」的风格。

关联:第 4 章 Plug 管线、第 5 章 WebSock/Jason 入站处理、第 7 章提示 JSON 编码。

L29 end

模块结束。这里没有服务运行期的请求处理逻辑;这些函数会在构建/启动准备阶段被 Mix 调用,返回项目与应用元数据。

这提醒我们一个事实:Elixir 的「代码即数据」让项目配置本身也是一段可以被任何工具读取的 Elixir 表达式。

什么时候会动这个文件

加依赖(往 deps/0 加一行)、改应用入口(mod:)、加编译期依赖或 escript 配置。其他时候基本不动——就像出生证明,偶尔才拿出来改。

第 3 章

监督树:进程从哪来,谁负责谁

应用自己的顶层 Supervisor 只有两个直接 child spec,但每个组件内部还会管理更多进程。

依赖应用启动后,应用控制器按第 2 章的 mod: 调用 VoiceChat.Application.start/2。它准备资源、声明 Registry 与 Bandit 两个直接孩子,再启动顶层 Supervisor;Bandit/Thousand Island 继续管理监听器与连接生命周期。

lib/voice_chat/application.ex· 23 行 · 应用入口 + 监督树 0/0
defmodule VoiceChat.Application do
  @moduledoc false
  use Application

  @impl true
  def start(_type, _args) do
    # 确保录音目录存在
    VoiceChat.recordings_dir() |> File.mkdir_p!()

    # 确保提示音(“好的”)存在,不存在则生成
    VoiceChat.PromptAudio.ensure!()

    children = [
      # 会话注册表:session_id -> ws 进程,进程退出自动清理
      {Registry, keys: :unique, name: VoiceChat.SessionRegistry},
      # HTTP + WebSocket 服务器
      {Bandit, plug: VoiceChat.Endpoint, port: VoiceChat.port()}
    ]

    opts = [strategy: :one_for_one, name: VoiceChat.Supervisor]
    Supervisor.start_link(children, opts)
  end
end
L1–4 模块声明 + use Application

use Application 要求实现 start/2(和可选的 stop/1、prep_stop/1)。@moduledoc false 表示这个模块是内部实现细节,不生成文档。

OTP 应用的启动/停止是「框架回调」:VM 的应用控制器在正确时机调用 start/2,你只负责返回 {:ok, pid} 或 {:ok, pid, state}。

关联:第 2 章 L17 mod: {VoiceChat.Application, []}——这行就是接线。

L5–6 @impl true + start/2

@impl true 声明「下面是行为回调的实现」,编译器会校验签名是否匹配。_type 是启动类型(常见为 :normal,也可能是 takeover/failover 元组),_args 是 mod: 里传的参数 [],这里都用不上。

下划线前缀的变量名是「我故意不用你」的约定,避免未使用变量的编译警告。

L7–9 确保录音目录存在

VoiceChat.recordings_dir() 从配置读目录(priv/recordings),File.mkdir_p!/1 递归建目录。|> 是管道:左边结果作为右边函数第一个参数。

结尾 ! 是 Elixir 约定:失败直接抛异常(而非返回 {:error, …})。这一步发生在顶层 Supervisor 启动之前,因此失败会让 Application 的启动回调失败,错误会在启动期暴露。

关联:第 6 章 Recorder 的 build_path 写文件时假设这个目录已存在。

L10–12 启动期准备提示音

VoiceChat.PromptAudio.ensure!() 检查 priv/static/audio/haode.wav 是否存在,不存在就现场合成一段(第 8 章)。

启动期能尽早发现“文件缺失且兜底无法写入”这类错误;但 ensure! 对已存在文件只检查存在性,不验证可读性或 WAV 内容,损坏素材仍可能到 Prompter 读取或浏览器解码时才暴露。

关联:第 7 章 Prompter 启动时 PromptAudio.bytes() 直接读文件。

L13–15 children:被监督的孩子

声明两个直接 child spec:Registry(唯一键会话注册表,条目属于注册它的连接进程,进程退出时自动清理)和 Bandit(HTTP/WS 服务器组件)。它们本身都可能在内部启动多个进程。

放进监督树后,顶层组件异常退出会按策略重启;若超过重启强度,Supervisor 自身仍会终止,因此监督提供的是明确的恢复策略,不是“永不失败”。

关联:第 5 章 Registry.register;第 4 章 Bandit, plug: Endpoint 的 plug 选项。

L16–19 Bandit 子进程规格

{Bandit, plug: VoiceChat.Endpoint, port: VoiceChat.port()} 是「子进程规格」简写:模块 + 启动参数。Supervisor 会用 Bandit.child_spec/1 把它展开成完整的启动指令。

plug: 告诉 Bandit「HTTP 请求交给这个模块处理」,port: 从配置读端口(默认 4000;本项目的 config.exs 在被求值时读取 PORT)。若打 release 后仍要在每次启动时读取环境变量,通常应把这类逻辑放进 runtime.exs。

关联:第 4 章 Endpoint 整章;第 1 章 config/ 目录。

L20–23 把孩子们交给 Supervisor

Supervisor.start_link(children, opts) 启动监督者并 link 到调用进程。strategy: :one_for_one 表示「一个孩子挂了只重启那一个」;name: 把监督者注册为 VoiceChat.Supervisor。

link 把 Application 启动进程与顶层 Supervisor 的生命周期接起来。one_for_one 表示某个直接孩子退出时只重启该 child spec;但 Registry 重启会丢失旧注册信息,已有连接与调试视图仍可能需要额外恢复设计。

关联:第 5 章每个连接的 spawn_link 是同一机制的最简形态;第 6 章 trap_exit 处理的就是 link 产生的信号。

VoiceChat.Supervisor one_for_one Registry keys: :unique session_id → ws 进程 Bandit HTTP + WebSocket port 4000 管理连接生命周期 动态区域:连接进程由 Bandit / Thousand Island 管理 ws 进程 WSHandler Recorder spawn_link Prompter
图 3-1 · 简化后的生命周期关系:应用顶层 Supervisor 直接管理 Registry 与 Bandit;Bandit/Thousand Island 动态管理连接进程,连接进程再分别 link Recorder 与 Prompter。Registry、Bandit/Thousand Island 自己的内部进程树没有展开。
为什么应用代码不把每条连接列成顶层 child spec

WebSocket 连接短命且动态,断开后也不应按固定 child spec 重启。它们的生命周期由 Bandit/Thousand Island 的服务器层管理;本应用只把长期的服务器组件 Bandit 放进自己的顶层监督树。WSHandler 初始化后,再分别 link Recorder 与 Prompter 这两个会话辅助进程。

类比:监督树 = 值班经理

Supervisor 像值班经理:直接孩子(Registry、Bandit)退出时按策略处理。连接则像由 Bandit 业务线内部调度的短期任务——仍有明确的生命周期管理,只是不逐条出现在本应用的顶层 children 列表里。

第 4 章

Endpoint:Plug 管线和 101 升级

从 HTTP 请求切换到 WebSocket 帧协议,入口是一轮成功的 101 Upgrade 握手。

endpoint.ex 是流量的第一站:静态页面、健康检查、会话列表和 WebSocket 升级都在这个模块里。理解它需要两个核心概念:Plug 管线(请求穿过一串 conn 变换)和 WebSocket 升级(从 HTTP 握手切换到帧协议)。

lib/voice_chat/endpoint.ex· 源码全文 · 流量第一站 0/0
defmodule VoiceChat.Endpoint do
  @moduledoc """
  HTTP 路由:静态调试页面 + WebSocket 升级接口 + 调试辅助接口。
  """

  use Plug.Router

  plug(Plug.Logger)

  plug(Plug.Static,
    at: "/",
    from: :voice_chat,
    only: ~w(index.html app.js app.css)
  )

  plug(:match)
  plug(:dispatch)

  get "/" do
    send_file(conn, 200, Application.app_dir(:voice_chat, "priv/static/index.html"))
  end

  get "/health" do
    send_resp(conn, 200, "ok")
  end

  get "/sessions" do
    sessions =
      VoiceChat.SessionRegistry
      |> Registry.select([{{:"$1", :"$2", :"$3"}, [], [{{:"$1", :"$2", :"$3"}}]}])
      |> Enum.map(fn {id, pid, meta} ->
        %{
          "session" => id,
          "ws_pid" => inspect(pid),
          # PID 不是 JSON 类型,编码前统一转成字符串,否则会话进行中该接口会 500
          "meta" =>
            Map.new(meta, fn
              {k, v} when is_pid(v) -> {k, inspect(v)}
              {k, v} -> {k, v}
            end)
        }
      end)

    conn
    |> put_resp_content_type("application/json")
    |> send_resp(200, Jason.encode!(sessions))
  end

  get "/ws" do
    conn = Plug.Conn.fetch_query_params(conn)

    try do
      conn
      |> WebSockAdapter.upgrade(
        VoiceChat.WSHandler,
        %{query_params: conn.query_params},
        timeout: 60_000,
        max_frame_size: 10_000_000
      )
      |> halt()
    rescue
      e in WebSockAdapter.UpgradeError ->
        send_resp(conn, 400, "websocket upgrade failed: #{Exception.message(e)}")
    end
  end

  match _ do
    send_resp(conn, 404, "not found")
  end
end
L1–7 use Plug.Router:路由即模块

use Plug.Router 生成路由所需代码,并让模块实现 Plug 约定的 init/1、call/2:请求以 %Plug.Conn{} 进入,也以更新后的 conn 返回。

任何实现 init/1 + call/2 的模块都可以作为 plug,既能单独调用,也能串进管线;Plug.Router 再在这个基础上提供匹配与分派宏。

关联:第 3 章 Bandit 的 plug: VoiceChat.Endpoint 选项——服务器把请求交给这个模块。

L8–9 plug(Plug.Logger)

把日志中间件挂进管线,记录请求方法/路径以及响应状态与耗时。

这是本管线的第一个「关卡」:请求先经过它,再到路由。中间件可以更新 conn;需要提前结束后续管线时,还要发送响应并 halt。

L10–15 Plug.Static:静态文件服务

把 priv/static/ 挂到根路径:at: "/"(URL 前缀)、from: :voice_chat(:voice_chat 是 priv/static 的简写)、only: ~w(index.html app.js app.css)(只放行这三个文件)。

only 是白名单:调试页面的静态资源就三个,白名单既快又安全(audio/haode.wav 故意不放行——它只走 WebSocket 推给客户端,不需要被 HTTP 直接下载)。~w(...) 是字符串列表字面量。

关联:第 9 章前端页面;第 8 章提示音文件在 priv/static/audio/ 但不走 HTTP。

L16–18 :match 与 :dispatch

Plug.Router 的两段式::match 找到匹配的路由,:dispatch 执行它。本模块把 Logger、Static 排在两者之前,所以请求会先按声明顺序经过它们。

显式顺序让日志覆盖本模块收到的请求,也让 Static 有机会在路由匹配前直接服务白名单资源。

L19–22 get "/":首页

根路径显式用 send_file 返回 index.html;Plug.Static 负责的是按 URL 文件名请求的 /app.js、/app.css 等资源,并不会自动把 / 当作目录首页。

把首页路由和静态资源白名单分开后,根路径行为清楚,也不依赖目录索引约定。

L23–26 /health

健康检查:send_resp(conn, 200, "ok")。部署时负载均衡/探针打这个接口确认服务活着。

L27–43 /sessions:读取并序列化注册项

Registry.select/2 的第一个参数是 Registry 名,第二个才是 match spec(ETS 匹配规格)。{{:"$1", :"$2", :"$3"}, [], [...]} 把条目投影成键、注册进程 pid 和元数据;随后把 pid 转成字符串,避免 Jason 直接编码 pid 失败。

match spec 是 Registry/ETS 的查询语言。这里没有筛选条件,只是一次取出调试接口需要的三个字段;元数据里的 recorder pid 也在 Map.new 中转成可编码字符串。

关联:第 5 章 Registry.register(写入)——这里是同一张表的读取。

L44–48 把列表变成 JSON

put_resp_content_type("application/json") 设置 Content-Type,Jason.encode! 把 Elixir 数据结构编码成 JSON 字符串返回。

Jason.encode! 的 ! 表示编码失败会抛异常;前一段显式字符串化已知 pid,正是为了让当前元数据结构可编码。

L49–51 /ws:先取查询参数

fetch_query_params(conn) 解析 URL 里的查询串(?session=xxx)放进 conn.query_params。注意这里手动取——因为 WebSockAdapter.upgrade 之后 conn 就不再是普通 HTTP 请求了,参数得在升级前拿到。

关联:第 5 章 init/1 收到的 %{query_params: …} 就是从这里传进去的。

L52–60 WebSocket 升级:一次 101 握手

WebSockAdapter.upgrade 校验升级请求,在 conn 上登记 WebSock 回调模块和初始参数;Bandit 随后发送 101 并让同一个 HTTP/1 连接进程切换到 WebSocket handler。第三个参数会传给 WSHandler.init/1,halt(conn) 则停止后续 Plug 管线。

WebSocket 有独立的帧协议,但通常先通过 HTTP/1 Upgrade 建立:请求带升级头,服务器回 101,之后同一 TCP 连接切换为 WebSocket 帧。timeout 控制空闲断开,max_frame_size 限制单帧大小。

关联:第 5 章整章——升级完成后,本项目实现的 WebSock 业务回调进入 WSHandler;连接协议和控制帧细节仍由 Bandit/WebSockAdapter 承接。

L61–66 rescue:升级失败要友好

升级请求不合法(缺 Sec-WebSocket-Key 等)时 upgrade 会抛 WebSockAdapter.UpgradeError,这里捕获并回 400。

「让用户看到 400 而不是 500」:协议错误是客户端的问题,不该算服务器内部错误。

L67–70 match _:兜底 404

任何没被上面路由接住、也没被 Static 接住的请求,返回 404。

显式兜底把未匹配路径稳定地变成 404;若没有匹配子句,生成的路由匹配函数会因找不到可用子句而失败,而不是自动替本模块构造这条响应。

HTTP 请求 Plug.Logger Plug.Static :match :dispatch /ws WebSockAdapter.upgrade 101 → 切换 WSHandler 静态文件 / 首页 直接返回 200 match _ → 404
图 4-1 · Plug 管线:请求依次穿过 Logger、Static,再由 match/dispatch 分派。/ws 登记升级后,Bandit 让该连接切换到 WebSocket/WSHandler;未匹配路径进入 404 兜底。
最容易踩的坑:升级之后 conn 就「死」了

WebSockAdapter.upgrade 返回的 conn 状态是 :upgraded,不能再 send_resp。所以任何「升级前要做的检查」(比如鉴权、取参数)都必须在 upgrade 调用之前完成——这正是我们先把 fetch_query_params 放在前面的原因。

第 5 章

WSHandler:一个 session 的一生

升级完成之后,这段代码就活在一个属于这个连接的进程里。

第 4 章末尾,Bandit 完成 101 握手;在 HTTP/1 路径上,原连接进程切换到 WebSocket handler,并调用 WSHandler.init/1,不是升级后再开一个 ws 进程。本模块实现 init、handle_in、handle_info 三个必需回调,以及可选的 terminate;WebSock 行为另一个可选回调 handle_control/2 未实现。

先记住一个事实:这段代码跑在谁的进程里

Bandit 的 WebSocket 实现(Bandit.WebSocket.Connection)在连接进程里调用我们的回调。所以 self() 就是「这个连接的进程」;别人 send(ws_pid, msg) 时,handle_info 会被调用——这正是第 7 章 Prompter 推音频给客户端用的通道。

lib/voice_chat/ws_handler.ex· 源码全文 · 会话生命周期 0/0
defmodule VoiceChat.WSHandler do
  @moduledoc """
  WebSocket 会话处理器(WebSock 回调模块,运行在 Bandit 的 ws 进程中)。

  每个连接即一个新 session:
  - `init/1` 启动录音进程与提示音进程,并把 session 注册到 Registry;
  - 二进制帧(麦克风音频块)转发给录音进程落盘;
  - 文本帧为 JSON 控制消息(config / text 等);
  - `handle_info` 响应提示音进程发来的下发请求(文本 + 音频)。
  """

  @behaviour WebSock

  require Logger

  @impl true
  def init(%{query_params: query_params}) do
    session_id = VoiceChat.SessionId.sanitize(query_params["session"])
    ws_pid = self()

    # 先占用最终 session id,再用它启动录音进程,避免重名连接写入旧 id 的文件。
    session_id = register_session(session_id, ws_pid)
    recorder = VoiceChat.Recorder.start(session_id, ws_pid)

    {_, _} =
      Registry.update_value(
        VoiceChat.SessionRegistry,
        session_id,
        &Map.put(&1, :recorder, recorder)
      )

    prompter = VoiceChat.Prompter.start(ws_pid)

    Logger.info("session opened id=#{session_id} ws_pid=#{inspect(ws_pid)}")

    session_message = Jason.encode!(%{type: "session", session: session_id})
    state = %{session_id: session_id, ws_pid: ws_pid, recorder: recorder, prompter: prompter}

    {:push, {:text, session_message}, state}
  end

  defp register_session(session_id, ws_pid) when is_pid(ws_pid) do
    meta = %{ws_pid: ws_pid, started_at: DateTime.utc_now() |> DateTime.to_iso8601()}
    register_session(session_id, meta)
  end

  defp register_session(session_id, meta) do
    case Registry.register(VoiceChat.SessionRegistry, session_id, meta) do
      {:ok, _} ->
        session_id

      {:error, {:already_registered, _pid}} ->
        register_session(VoiceChat.SessionId.generate(), meta)
    end
  end

  @impl true
  def handle_in({data, opcode: :binary}, state) do
    VoiceChat.Recorder.chunk(state.recorder, data)
    {:ok, state}
  end

  def handle_in({text, opcode: :text}, state) do
    case Jason.decode(text) do
      {:ok, %{"type" => "config"} = config} ->
        VoiceChat.Recorder.config(state.recorder, config)
        {:ok, state}

      {:ok, %{"type" => "text", "text" => t}} when is_binary(t) ->
        # 模拟聊天:回一句 echo
        reply = Jason.encode!(%{type: "text", message: "收到: #{t}", session: state.session_id})
        {:push, {:text, reply}, state}

      {:ok, _other} ->
        {:ok, state}

      {:error, _} ->
        Logger.warning("bad json frame session=#{state.session_id}")
        {:ok, state}
    end
  end

  @impl true
  def handle_info({:send_binary, data}, state), do: {:push, {:binary, data}, state}
  def handle_info({:send_text, text}, state), do: {:push, {:text, text}, state}
  def handle_info(_msg, state), do: {:ok, state}

  @impl true
  def terminate(_reason, state) do
    VoiceChat.Prompter.stop(state.prompter)

    if VoiceChat.Recorder.stop(state.recorder) == {:error, :timeout} do
      Logger.warning("recorder stop timed out session=#{state.session_id}")
    end

    Logger.info("session closed id=#{state.session_id}")
    :ok
  end
end
L1–11 模块文档:先写清「谁在什么时候调用我」

@moduledoc 不只是注释——mix docs 会生成文档,而更重要的是一开始就想清楚回调的职责边界。

WebSock 回调的「触发者」各不相同:init 由连接建立触发、handle_in 由客户端帧触发、handle_info 由别的进程发消息触发、terminate 由连接关闭触发。文档把这些写清楚,后面每行代码都有的放矢。

L12–13 @behaviour WebSock

声明「我实现了 WebSock 这个行为」。编译器会按行为定义检查回调并对缺失或不匹配给出警告;本项目实现必需的 init/1、handle_in/2、handle_info/2,以及 terminate/2。

行为(behaviour)类似一份回调接口:适配层把协议事件翻译成函数调用,所以本项目业务模块无需自行实现 RFC 6455 的握手与帧解析。

关联:websock 是间接依赖里的行为定义方,websock_adapter 负责把 Plug/Bandit 的升级接进来。

L14–15 require Logger

引入 Logger 宏(Logger.info/1 等其实是宏,需要 require)。

细节:为什么是 require 不是 import/alias?因为 Logger 的 API 是宏(延迟求值 + 编译期裁剪日志级别),宏必须 require。

L16–17 init/1:session 的出生

连接建立后调用一次。参数是第 4 章 upgrade 时传的 %{query_params: …},用模式匹配直接解包。@impl true 标记这是行为回调。

「每个连接 = 一个新 session」就发生在这里:每连接进程各调一次 init,天然隔离,不需要 session 池、不需要锁。

L18 session id:来自查询参数或自动生成

VoiceChat.SessionId.sanitize(query_params["session"]):缺少或清洗后为空时自动生成;其他输入把非字母数字/下划线/连字符替换为 _,并截到最多 64 个字符。

用户输入永远不可信:id 会拼进文件名(第 6 章 build_path),不清洗就可能写出路径穿越。这是安全习惯,不是功能。

关联:lib/voice_chat/session_id.ex;第 6 章文件名拼接。

L19–20 ws_pid = self()

记下当前进程 pid——就是「这个连接的进程」。它会被交给 Recorder/Prompter,让它们能发消息回来(走 handle_info)。

pid 是定位进程的句柄;send/2 把异步消息放进它的邮箱。它和同步“方法调用”并不等价:没有返回值,也不会替调用方提供背压。

关联:第 6 章 trap_exit 里对比的 from 就是它;第 7 章 Process.alive?(ws_pid)。

L21–24 先确定最终 id,再启动 Recorder

先用 register_session/2 占用最终 session id,再把这个 id 交给 Recorder.start/2。这样重名连接换到新 id 时,Registry、WebSock state 和录音文件名前缀仍然一致。

Recorder 内部用 spawn_link。ws 异常退出时,Recorder 因 trap_exit 收到 EXIT 消息并收尾;反方向若 Recorder 异常退出,link 也会把故障传播给连接进程。

关联:第 6、7 章整章;第 0 章旅程图的 spawn_link 虚线。

L25–33 补齐 Registry 元数据,再启动 Prompter

Registry.update_value/3 把刚得到的 recorder pid 加入元数据;随后 Prompter.start/1 创建另一个 link 到 ws 的辅助进程。

Registry 的价值:/sessions 调试接口(第 4 章)能实时列出活跃会话,且进程退出时 Registry 自动清理条目——不需要手动 deregister,少一类「忘记清理」的 bug。

关联:第 3 章 Registry 子进程;第 4 章 Registry.select。

L34–41 首帧确认实际 session id

日志记录连接后,服务端把最终 session id 编成 {"type":"session",…},同时构造包含三个 pid 的 state,并从 init/1 返回一次文本 push。WebSock 后续回调都会收到这份 state。

客户端请求的 id 可能因冲突被替换,不能继续把本地候选值当真;首帧把 Registry 与文件名实际使用的 id 回告页面,页面再更新显示。

L42–56 唯一注册:冲突时递归换 id

先注册包含 ws pid 和开始时间的 meta;若唯一键已被占用,就生成新 id 后递归重试。注册项属于当前连接进程,连接退出时 Registry 自动删除。

先注册、再启动依赖 id 的 Recorder,避免了“展示 id 已换新、文件仍用旧 id”的分裂状态。

L57–62 handle_in:二进制帧 = 音频块

客户端发来的二进制帧({data, opcode: :binary})就是一块音频,直接 VoiceChat.Recorder.chunk(state.recorder, data) 转发给录音进程,然后 {:ok, state} 继续。

chunk/2 只是异步 send,所以磁盘写入不在连接进程中执行;代价是当前实现没有背压,Recorder 跟不上时消息会积在邮箱里。

关联:第 6 章 write_chunk;第 9 章前端每次回调读取 4096 个输入采样帧,降采样后发一个二进制块。

L63–68 文本帧 = JSON 控制消息

Jason.decode 解析 JSON,按 type 分派:config 转发给 Recorder(设置格式/采样率)。

当前协议约定二进制帧承载音频,文本帧承载 JSON 控制消息;这是清晰的类型分工,但仍需由两端共同维护,并不会由 WebSocket 自动验证内容语义。

L69–73 文本回显:协议演示

收到 {"type":"text","text":t} 就回一条 {"type":"text","message":"收到: …"}。当前网页没有文本输入框,这条分支是协议/测试脚本可调用的回显示例,页面只会把收到的文本写进消息日志。

guard(when is_binary(t))挡掉缺字段的脏消息;{:push, msg, state} 是 WebSock 的「立即下发」返回。

关联:第 9 章 handleText;scripts/ws_advanced_test.mjs 会实际发送 text 消息。

L74–82 其他 JSON 与坏消息

未知 type 静默忽略;解析失败的帧打一条 warning。

协议要宽容:客户端(尤其浏览器)偶尔发怪帧,别让一个坏帧杀掉整个连接。

L83–87 handle_info:别人给我的消息

Prompter 每轮等待约 3 秒后,依次发送 {:send_text, json} 和 {:send_binary, audio};这里把它们转换为 WebSock 的 {:push, …} 返回值。

持有连接 pid 的进程可以向其邮箱发送约定消息,无需持有底层 socket 句柄;WSHandler 再把这些消息翻译成 WebSock push。当前项目的周期提示走这条路。

关联:第 7 章 Prompter 的 send 调用——两边是配对的。

L88–99 terminate:有界等待录音收尾

先异步停止 Prompter,再调用 Recorder.stop/1;后者请求 finalize 并等待 Recorder 退出,最长 5 秒。超时只记录 warning,随后连接回调返回。

常规成功路径因此有时间完成 WAV 回填;不过 stop 只观察 DOWN、没有检查退出 reason,也不是持久化 ACK,不能单凭 :ok 证明写盘成功。terminate/2 本身也并非所有故障下都保证执行。

关联:第 6 章 stop/1、monitor 与 finalize;第 0 章旅程图的断开收尾。

① 出生 · 升级 /ws → 101 连接进程切换 handler init/1 被调用 ② 组建家庭 先占用 Registry session id spawn_link 两个帮手 首帧回告实际 id ③ 过日常 handle_in:音频块上行 handle_info:提示音下行 (可以持续几小时) ④ 常规收尾 terminate/2 被调用 先停 Prompter 等待 Recorder(最多 5s)
图 5-1 · 一个 session 的常规路径:已有连接进程完成升级并 init → 先占用实际 id、两次 spawn_link 并回告 id → handle_in/handle_info 处理事件 → 常规关闭路径调用可选的 terminate,在 5 秒上限内等待 Recorder 退出并尝试收尾。图中是本模块实现的四类回调,不含另一个可选回调 handle_control。
第 6 章

Recorder:把音频流变成文件

225 行,用一个 receive 循环把格式选择、文件写入和有界收尾串起来。

Recorder 为了教学直接使用 spawn_link + receive,没有套 GenServer。它展示 link / trap_exit、monitor 等待退出、惰性打开、PCM 帧对齐与 WAV 回填;也保留了裸进程方案缺少标准调试、监督与背压的局限。

lib/voice_chat/recorder.ex· 源码全文 · 音频落盘进程 0/0
defmodule VoiceChat.Recorder do
  @moduledoc """
  每个 session 的音频落盘进程(轻量进程)。

  接收 `{:chunk, binary}` 持续写入文件;收到 `:stop` 或所属 ws 进程退出时,
  补齐 WAV 头并关闭文件。PCM16 数据在 finalize 时包装成标准 .wav 文件,
  Opus 的 WebM / Ogg 容器数据则原样拼接保存。
  """

  require Logger

  @stop_timeout 5_000
  @wav_header_size 44

  @doc "启动录音进程,与调用进程(ws 进程)link。"
  def start(session_id, ws_pid) do
    spawn_link(fn -> init_loop(session_id, ws_pid) end)
  end

  @doc "设置音频格式 / 采样率等元信息。"
  def config(pid, config) when is_pid(pid), do: send(pid, {:config, config})

  @doc "写入一块音频数据。"
  def chunk(pid, data) when is_pid(pid) and is_binary(data), do: send(pid, {:chunk, data})

  @doc "请求停止并等待进程退出;常规 `:stop` 分支先 finalize,超时返回 `{:error, :timeout}`。"
  def stop(pid) when is_pid(pid) do
    monitor = Process.monitor(pid)
    send(pid, :stop)

    receive do
      {:DOWN, ^monitor, :process, ^pid, _reason} ->
        :ok
    after
      @stop_timeout ->
        Process.demonitor(monitor, [:flush])
        {:error, :timeout}
    end
  end

  defp init_loop(session_id, ws_pid) do
    # trap 退出信号:ws 进程异常退出时也能 finalize 文件
    Process.flag(:trap_exit, true)

    loop(%{
      session_id: session_id,
      ws_pid: ws_pid,
      format: :pcm16,
      sample_rate: 16_000,
      channels: 1,
      io: nil,
      path: nil,
      bytes: 0,
      finalized: false
    })
  end

  defp loop(state) do
    receive do
      {:config, config} ->
        loop(apply_config(state, config))

      {:chunk, data} ->
        loop(write_chunk(state, data))

      :stop ->
        finalize_and_exit(state)

      {:EXIT, from, _reason} ->
        if from == state.ws_pid, do: finalize_and_exit(state), else: loop(state)
    end
  end

  # --- config ---------------------------------------------------------------

  defp apply_config(state, config) when is_map(config) do
    format = container_format(config)

    sample_rate = config["sample_rate"] || state.sample_rate
    channels = config["channels"] || state.channels

    {sample_rate, channels} =
      if valid_wav_params?(sample_rate, channels) do
        {sample_rate, channels}
      else
        Logger.warning(
          "invalid audio config ignored session=#{state.session_id} " <>
            "sample_rate=#{inspect(sample_rate)} channels=#{inspect(channels)}"
        )

        {state.sample_rate, state.channels}
      end

    if state.io do
      Logger.warning(
        "config ignored (file already open) session=#{state.session_id} format=#{inspect(format)}"
      )

      state
    else
      %{state | format: format, sample_rate: sample_rate, channels: channels}
    end
  end

  defp valid_wav_params?(sample_rate, channels)
       when is_integer(sample_rate) and sample_rate > 0 and sample_rate <= 0xFFFF_FFFF and
              is_integer(channels) and channels > 0 and channels * 2 <= 0xFFFF and
              sample_rate * channels * 2 <= 0xFFFF_FFFF,
       do: true

  defp valid_wav_params?(_sample_rate, _channels), do: false

  defp container_format(%{"format" => "opus", "container" => "ogg"}), do: :ogg
  defp container_format(%{"format" => format}) when format in ["opus", "webm"], do: :webm
  defp container_format(%{"format" => "ogg"}), do: :ogg
  defp container_format(_config), do: :pcm16

  # --- writing --------------------------------------------------------------

  defp write_chunk(%{io: nil} = state, data) do
    case ensure_open(state) do
      %{io: nil} = unopened -> unopened
      opened -> write_chunk(opened, data)
    end
  end

  defp write_chunk(state, data) do
    :ok = :file.write(state.io, data)
    %{state | bytes: state.bytes + byte_size(data)}
  end

  defp ensure_open(state) do
    path = build_path(state)

    case :file.open(path, [:write, :binary, :raw]) do
      {:ok, io} ->
        case state.format do
          :pcm16 ->
            # 先写占位头,finalize 时回填真实大小
            :ok = :file.write(io, VoiceChat.Wav.header(state.sample_rate, state.channels, 0))

          format when format in [:webm, :ogg] ->
            :ok
        end

        %{state | io: io, path: path}

      {:error, reason} ->
        Logger.error("cannot open recording file #{path}: #{inspect(reason)}")
        state
    end
  end

  defp build_path(state) do
    timestamp = System.system_time(:microsecond)
    sequence = System.unique_integer([:positive, :monotonic])
    ext = if state.format == :pcm16, do: "wav", else: Atom.to_string(state.format)
    dir = VoiceChat.recordings_dir()
    Path.join(dir, "#{state.session_id}_#{timestamp}_#{sequence}.#{ext}")
  end

  # --- finalize -------------------------------------------------------------

  defp finalize_and_exit(state) do
    finalize(state)
    :ok
  end

  defp finalize(%{finalized: true} = state), do: state

  defp finalize(%{io: nil} = state), do: %{state | finalized: true}

  defp finalize(state) do
    state =
      case state.format do
        :pcm16 -> finalize_pcm(state)
        format when format in [:webm, :ogg] -> finalize_container(state)
      end

    %{state | finalized: true, io: nil}
  end

  defp finalize_pcm(state) do
    frame_size = state.channels * 2
    complete_bytes = state.bytes - rem(state.bytes, frame_size)
    dropped_bytes = state.bytes - complete_bytes

    if dropped_bytes > 0 do
      end_position = @wav_header_size + complete_bytes
      {:ok, ^end_position} = :file.position(state.io, end_position)

      :ok = :file.truncate(state.io)

      Logger.warning(
        "dropped incomplete PCM frame session=#{state.session_id} bytes=#{dropped_bytes}"
      )
    end

    {:ok, 0} = :file.position(state.io, 0)

    :ok =
      :file.write(
        state.io,
        VoiceChat.Wav.header(state.sample_rate, state.channels, complete_bytes)
      )

    :ok = :file.close(state.io)

    Logger.info(
      "recorded wav session=#{state.session_id} path=#{state.path} audio_bytes=#{complete_bytes}"
    )

    %{state | bytes: complete_bytes}
  end

  defp finalize_container(state) do
    :ok = :file.close(state.io)

    Logger.info(
      "recorded container session=#{state.session_id} path=#{state.path} bytes=#{state.bytes}"
    )

    state
  end
end
L1–19 模块职责、常量与 start/2

模块文档先声明“收 chunk、在停止时收尾”的边界,require Logger 启用后文日志宏;随后定义 5 秒停止上限与 44 字节 WAV 头大小。start/2 用 spawn_link 运行 init_loop;在实际调用路径里,调用者与传入的 ws_pid 都是当前连接进程。

link 会双向传播退出信号,但不等同于监督策略:Recorder 异常退出会影响 ws;ws 退出时 Recorder 则通过后面的 trap_exit 把信号转成消息并尝试收尾。

关联:第 5 章 L21–23 的调用方;第 3 章 Supervisor 在 link 之上还增加 child spec 与重启策略。

L20–25 config / chunk:异步邮箱 API

config/2 与 chunk/2 都把消息放入 Recorder 邮箱并立即返回;它们不等待接收方处理或磁盘写完。

连接进程可以继续处理帧,但当前实现没有端到端背压;若写盘落后,chunk 会在 Recorder 邮箱中累积。

L26–40 stop/1:monitor + 5 秒有界等待

先 monitor Recorder,再发送 :stop,等待对应 :DOWN;5 秒内未退出则清掉 monitor 并返回 {:error, :timeout},但不会强杀 Recorder。

monitor 先于 stop,避免错过极快退出。来自同一个 ws 进程的 config、chunk、stop 对同一 Recorder 保持发送顺序,因此常规路径会先处理此前入队的音频,再尝试 finalize。DOWN 只证明进程已退出;代码忽略 reason,不能据此断言收尾成功。

L41–57 trap_exit 与私有 state

Recorder 启动后立刻启用 trap_exit,再进入带有 session id、ws pid、格式、采样参数、文件句柄、路径、字节数与 finalized 标志的 state map。

默认格式是 PCM16 / 16kHz / 单声道,文件保持 nil 以便惰性打开。trap_exit 能处理 linked ws 的可传播退出信号,但无法抵御直接 :kill Recorder、VM/主机掉电或磁盘故障。

L58–75 receive:四路状态转移

config 更新 state,chunk 写入后递归,:stop finalize 并退出;{:EXIT, from, …} 只在来自所属 ws pid 时收尾,其他退出消息继续循环。

尾递归把新 state 带入下一轮;没有匹配消息时进程阻塞等待,不占用可运行时间片,但仍保留内存与邮箱资源。

L76–93 配置:取值、验证、失败回退

先把外部格式归一化,再从 JSON 取采样率与声道;valid_wav_params?/2 失败时记录 warning,并保留当前采样参数。

这些字段最终会进入固定宽度 WAV 头。对正整数、block align 与 byte rate 做边界检查,可避免不可信 config 在位语法构造时溢出或触发异常。

L94–104 文件已打开:拒绝中途换配置

若 state.io 已存在,就警告并保持原 state;尚未打开时才写入归一化后的格式与 WAV 参数。

PCM 头语义和容器扩展名在首次打开时已经确定,中途改格式会让同一文件前后矛盾,所以协议要求 config 先于首个音频块。

L105–119 参数 guard 与容器归一化

guard 检查采样率、声道数、block align 和 byte rate 能装入 RIFF/WAV 对应字段;容器函数把 Opus + Ogg 映为 :ogg,其余受支持的 Opus/WebM 映为 :webm,未知配置回到 PCM16。

Opus 是编码,WebM/Ogg 是容器;分开记录后,后缀才能对应浏览器实际选择的 MIME 容器。

L120–131 第一块惰性打开,之后持续追加

第一块到来而 io 仍是 nil 时,先 ensure_open;成功才递归写当前块,失败则返回未打开 state。后续块直接 :file.write 并累加字节数。

没有音频的连接不会留下空文件。发送方不等待这次写盘,因此邮箱增长仍是需要监控的容量边界。

L132–147 打开文件:PCM 先写占位头

以 raw binary 写模式打开唯一路径;PCM16 先写 data size 为 0 的 44 字节 WAV 头,WebM/Ogg 不加 WAV 头,然后把句柄与路径写回 state。

流式 PCM 事先不知道最终长度,所以先占位、结束时回填。容器分块已经由 MediaRecorder 编码,服务端只追加其字节。

L148–153 打开失败:丢当前块,下一块再试

权限、路径或磁盘问题使 open 返回 error 时,只记日志并保留 io: nil;外层因此丢掉当前块,下一块到来才再次尝试。

这样不会对同一 chunk 无限递归。策略是“连接继续、录音尽力而为”;生产环境还应配合告警、退避或关闭会话。

L154–163 路径:微秒时间 + VM 内唯一序号

按 :pcm16 / :webm / :ogg 选择后缀,用已清洗 session id、微秒时间戳和单调唯一整数组成文件名,再拼到配置的录音目录。

时间戳便于追踪,System.unique_integer 在同一 VM 内消歧;它不是跨节点全局 id。

L164–168 finalize 后让进程退出

finalize_and_exit/1 调用收尾函数后返回 :ok,顶层匿名函数结束,Recorder 随之退出,monitor 观察到 DOWN。

保存结果以日志与文件系统为准;即将关闭的 WebSocket 不是可靠的“录音完成”确认通道。

L169–182 幂等边界与格式分派

已 finalized 或从未打开文件时直接返回;否则按 PCM16 与 WebM/Ogg 分派到各自收尾函数,最后标记 finalized 并清空句柄。

当前循环首次收尾后即退出,这些子句仍让 finalize/1 对“已完成/未打开”状态有明确结果。

L183–198 PCM 尾部只保留完整采样帧

每帧大小是 channels × 2 字节。若总字节数不能整除帧大小,就把文件定位到最后一个完整帧,截断残余字节并记录 warning。

WAV 的 block align 必须与数据对齐;宁可丢掉不足一帧的尾巴,也不能让头部宣称的数据结构与文件实际字节矛盾。

L199–215 回填真实长度、关闭并记录

文件指针回到 0,用完整帧字节数重建 WAV 头,再关闭句柄并记录路径;返回的 state 也改成截断后的真实字节数。

:file.position/2 返回 {:ok, position},代码按这个形状匹配。经典 RIFF 的 32 位长度字段要求单个 PCM data 小于约 4 GiB;Wav.header/3 会拒绝越界,本项目不自动分段,长录音需另做滚动文件或 RF64。

关联:第 8 章 Wav.header/3 的大小 guard。

L216–225 WebM / Ogg:不改容器,只关闭

容器模式不写或回填 WAV 头,只关闭句柄并记录累计字节数。

服务器不解码 Opus,只保存同一次 MediaRecorder 会话产生的连续容器分块;文件有效性仍取决于浏览器产出的流。

ws 进程 handle_in (Bandit 连接进程) Recorder 进程 receive 循环 trap_exit: true 录音文件 *.wav / *.webm / *.ogg priv/recordings/ {:config, …}(一次) {:chunk, data}(持续到达) :stop;调用方等待 DOWN :file.write(追加) link:ws 进程退出 → Recorder 收到 {:EXIT, …}
图 6-1 · config/chunk 由 ws 异步发送,Recorder 在自己的 receive 循环中写文件;常规关闭发送 :stop 后通过 monitor 等待 Recorder 退出,最长 5 秒。虚线表示 link 提供的异常退出收尾路径,但它不覆盖 Recorder 被直接 kill、VM 掉电或磁盘错误。
为什么这里不用 GenServer

这个教学项目用裸进程,是为了把邮箱、模式匹配和 trap_exit 完整展开。GenServer 不负责重启策略(Supervisor 才负责),但它提供标准的 call/cast、系统消息、调试与可观测接口,生产代码通常更易维护;Recorder 若需要查询状态、背压、统一遥测或独立监督,应优先评估 GenServer/动态监督,而不是把“裸进程更轻”当成选型口诀。

第 7 章

Prompter:每 3 秒的轻量定时器

Prompter 选择 receive … after,实现“每轮等待 3 秒再推送”。

需求是「接收音频的同时,由另一个轻量进程周期性推送提示」。本实现无需额外调度库,但 receive … after 不是唯一写法;本章也会对比 Process.send_after,并说明当前周期会包含每轮发送耗时,因此是约 3 秒而非墙钟级精确定时。

lib/voice_chat/prompter.ex· 42 行 · 3 秒定时推送 0/0
defmodule VoiceChat.Prompter do
  @moduledoc """
  每个 session 的提示音定时进程(轻量进程)。

  每隔 3 秒向所属 ws 进程发送一次提示音(“好的”)音频字节,
  由 ws 进程通过 WebSocket 推送给客户端播放。
  """

  @interval 3_000

  @doc "启动提示音进程,与调用进程(ws 进程)link。"
  def start(ws_pid) do
    spawn_link(fn -> loop(ws_pid, VoiceChat.PromptAudio.bytes()) end)
  end

  @doc "停止提示音进程。"
  def stop(pid) when is_pid(pid), do: send(pid, :stop)

  defp loop(ws_pid, audio) do
    receive do
      :stop -> :ok
    after
      @interval ->
        if Process.alive?(ws_pid) do
          send(
            ws_pid,
            {:send_text,
             Jason.encode!(%{
               type: "prompt",
               message: "好的",
               at: System.system_time(:second)
             })}
          )

          send(ws_pid, {:send_binary, audio})
          loop(ws_pid, audio)
        else
          :ok
        end
    end
  end
end
L1–8 文档:先写清「每 3 秒向 ws 进程发消息」

一句话说清职责:本进程不碰 socket,只负责定时往 ws 进程发两条消息(文本 + 音频),下发由第 5 章的 handle_info 完成。

职责分离的边界很干净:Prompter 管「什么时候发」,WSHandler 管「怎么发出去」。

L9–10 @interval 3_000:模块属性即常量

模块属性(@name value)在编译期被替换成字面量。数字下划线(3_000)只是可读性,等于 3000 毫秒。

把「魔法数字」提成具名属性:改间隔只改一行。这比在循环里写死 3000 好维护得多。

L11–15 start/1:子进程读一次音频再进循环

spawn_link 创建的子进程调用 PromptAudio.bytes() 读一次当前约 13KB 的文件,再让循环复用同一个 binary;大型 binary 通常由 BEAM 以引用计数方式共享底层数据。

start/1 返回 pid,并不等待子进程确认读取成功;若 File.read! 失败,linked Prompter 会异常退出并影响 ws。读取只发生在进循环前,不会叠加到每个 3 秒周期。

关联:第 8 章 PromptAudio.bytes;第 3 章 ensure! 的 fail-fast 哲学。

L16–18 stop/1

往进程发 :stop——进程收到后从 :stop -> :ok 分支退出。

Prompter.stop 是 fire-and-forget;Recorder.stop 也发送 :stop,但额外用 monitor 等待其退出。两者都避免直接强杀,只是确认语义不同。

L19–22 循环:先收消息,等不到就等 3 秒

receive do :stop -> :ok after @interval -> … end 会选择性接收 :stop;若 3000ms 内没有匹配消息,就执行 after 分支。

等待期间进程不会占用调度器执行时间,但仍保留堆、邮箱和计时状态。每次发送完成后才重新进入 receive,所以处理耗时会叠加到下一个周期。

L23–24 每次先确认 ws 进程还活着

if Process.alive?(ws_pid):ws 进程已退出就不发消息、直接 :ok 结束自己。

alive? 只是发送前的快照,检查后进程仍可能立刻退出;向已死 pid 发送也不会抛错。真正的生命周期机制是 ws 在 terminate 中发 :stop,以及异常退出沿 link 传播;这里的检查只避免已明确死亡时继续下一轮。

关联:第 5 章 terminate 的 Prompter.stop;第 6 章 trap_exit 是同一问题的另一种解法。

L25–34 先发一条文本通知

send(ws_pid, {:send_text, Jason.encode!(%{type: "prompt", message: "好的", at: …})})——JSON 里带消息类型和时间戳,前端据此在聊天框显示「服务端提示音: 好的」。

文本帧让前端有可展示、可统计的提示,随后同一发送者再发真正的音频;二者相对顺序有保证,但 ws 邮箱仍可能穿插其他发送者的消息。

关联:第 5 章 handle_info {:send_text, …};第 9 章 handleText 的 type: "prompt"。

L35–36 再发音频,然后进入下一轮

send(ws_pid, {:send_binary, audio}) 发提示音字节,loop(ws_pid, audio) 尾递归回到 receive——从此刻再等待 3 秒。

长生命周期的 receive 循环通常用尾递归表达;尾调用优化会复用调用帧,因此不会因每轮递归而持续增长栈。

L37–42 ws 已死的分支

:ok 结束进程——prompter 的生命随连接结束。

连接没了提示音就没意义,主动退出比空转强。

3s receive … after 3000 ws 进程 {:send_text, …} {:send_binary, 音频} 浏览器 handleText 显示提示 解码并播放音频
图 7-1 · 圆环表示约 3 秒的等待;周期末 Prompter 依次发送文本与音频,同一发送者保证二者相对顺序,但其他发送者的消息仍可穿插。发送完成后才开始下一轮,因此实际相邻推送间隔还包含本轮 JSON 编码与发送耗时;提示音字节只在进程启动时读取一次。
类比:receive/after 就像车站等车

Prompter 像等在站台上的人:receive 等待 :stop,after 3000 表示 3 秒没等到匹配消息就按超时路线出发。等待时不占用 CPU 时间片,但进程本身仍占内存并维护计时状态。

对比:Process.send_after 的另一种写法

Process.send_after(self(), :tick, 3000) 会创建一次性定时信号,主 receive 可同时处理 :tick 与其他消息;需要取消时要保存 timer reference。receive … after 的超时与当前 receive 调用绑定,若循环还处理其他消息,何时重新进入 receive 会影响下一次超时。Prompter 只匹配 :stop,所以当前写法足够直接。

第 8 章

二进制:WAV 头和提示音是怎么拼出来的

Elixir 的 bitstring 语法,让「拼二进制」跟拼字符串一样自然。

第 6 章 Recorder 落盘时写了一个 44 字节的 WAV 头——本章把它拆开看:wav.ex 怎么用一行 bitstring 构造出标准 RIFF 头,怎么现场合成一段「双音提示音」;prompt_audio.ex 怎么管理提示音文件。这也是全书第一次正面讲二进制是 Elixir 一等公民这件事。

先认识 WAV 的 44 字节头

52 49"RI" 46 46"FF" 36+n (4B)RIFF 块长(文件−8) 57 41"WA" 56 45"VE" fmt 块PCM=1 声道 采样率 字节率sample×ch×2 块对齐ch×2 16bit位深 data 块后面全是 PCM
图 8-1 · 标准 PCM WAV 的 44 字节头:偏移 4 的字段是 RIFF chunk size(总文件字节数减 8),data size 才是后续 PCM 的字节数。
lib/voice_chat/wav.ex· 源码全文 · WAV 头 + 提示音合成 0/0
defmodule VoiceChat.Wav do
  @moduledoc """
  16-bit PCM WAV 文件的头部构造与简单的提示音合成工具。
  """

  @max_u32 0xFFFF_FFFF
  @max_data_size @max_u32 - 36

  @doc """
  生成 44 字节的 WAV(RIFF) 头部,PCM 16-bit 小端。

  - `sample_rate`: 采样率,如 16000
  - `channels`: 声道数,如 1
  - `data_size`: PCM 数据字节数
  """
  @spec header(pos_integer(), pos_integer(), non_neg_integer()) :: binary()
  def header(sample_rate, channels, data_size)
      when is_integer(sample_rate) and sample_rate > 0 and sample_rate <= @max_u32 and
             is_integer(channels) and channels > 0 and channels * 2 <= 0xFFFF and
             sample_rate * channels * 2 <= @max_u32 and is_integer(data_size) and
             data_size >= 0 and data_size <= @max_data_size do
    byte_rate = sample_rate * channels * 2
    block_align = channels * 2

    <<"RIFF", 36 + data_size::little-32, "WAVE", "fmt ", 16::little-32, 1::little-16,
      channels::little-16, sample_rate::little-32, byte_rate::little-32, block_align::little-16,
      16::little-16, "data", data_size::little-32>>
  end

  @doc """
  合成一段双音提示音(660Hz 短音 + 停顿 + 880Hz 短音),返回完整 WAV 文件字节。
  在没有真人语音素材时的兜底方案。
  """
  @spec synthesize_chime(pos_integer()) :: binary()
  def synthesize_chime(sample_rate \\ 16_000) do
    tone = fn freq, duration ->
      n = round(sample_rate * duration)

      for i <- 0..(n - 1) do
        attack = min(1.0, i / (sample_rate * 0.01))
        release = min(1.0, (n - i) / (sample_rate * 0.03))
        :math.sin(2 * :math.pi() * freq * i / sample_rate) * 0.55 * attack * release
      end
    end

    samples =
      tone.(660, 0.12) ++
        List.duplicate(0.0, round(sample_rate * 0.06)) ++
        tone.(880, 0.18)

    pcm =
      samples
      |> Enum.map(fn s -> round(max(-1.0, min(1.0, s)) * 32_767) end)
      |> Enum.map(&<<&1::little-signed-16>>)
      |> IO.iodata_to_binary()

    header(sample_rate, 1, byte_size(pcm)) <> pcm
  end
end
L1–8 纯计算模块 + RIFF 上界

模块不启动进程,只把参数转换成二进制。@max_u32 是 32 位无符号上限;data 上限再减 36,确保 RIFF chunk size 的 36 + data_size 仍能装进 32 位字段。

把格式计算从 Recorder 的生命周期代码中拆出,既便于单测,也把经典 RIFF 的约束集中在一个地方。

L9–16 文档与 @spec

@spec header(pos_integer(), pos_integer(), non_neg_integer()) :: binary() 声明参数/返回类型,配合 dialyzer 做静态分析。

Elixir 运行时仍是动态类型;typespec 为文档、编辑器和可选的 Dialyzer 分析提供契约,真正的数值边界还由下面的 guard 执行。

L17–24 guard + 派生字段

guard 要求采样率、声道和 data size 为合法整数,并限制 16/32 位字段不溢出;随后推导 byte_rate = sample_rate × channels × 2 与 block_align = channels × 2。

经典 RIFF/WAV 使用 32 位长度,所以单个 PCM data 必须小于约 4 GiB。超界不会静默回绕,而是无法匹配这个函数子句;本项目也不实现自动分段或 RF64。

L25–29 一段 bitstring 拼出整个头

<<"RIFF", 36 + data_size::little-32, "WAVE", "fmt ", 16::little-32, …>>:Elixir 的位语法直接按字段拼二进制。::little-32 表示「4 字节小端无符号整数」,16::little-16 是 2 字节。

位语法让代码声明式描述布局,编译器负责构造对应 binary;它提高的是表达和模式匹配能力,并不意味着运行时完全没有分配或复制。

关联:第 6 章占位头/回填头都调这个函数;第 9 章前端 setInt16(…, true) 的小端是同一回事。

L30–45 合成提示音:匿名函数 + 包络

tone = fn freq, duration -> … end 是匿名函数(闭包),内部用列表推导 for i <- 0..(n-1) 生成采样点;attack/release 是包络(淡入淡出),防止「啪」的爆音;:math.sin 生成正弦波。

包络是音频常识:直接发 0→1 跳变会听到爆音,包络让音量从 0 渐入渐出。这也解释了为什么「无素材兜底」也要认真写。

L46–50 三段拼接:音 + 静音 + 音

660Hz、120ms 短音 + 60ms 静音 + 880Hz、180ms 短音。List.duplicate(0.0, n) 生成中间静音段。

这是缺少真人“好的”素材时的可辨识兜底音,不把它冒充成语音内容;文本帧仍会标注提示消息为“好的”。

L51–56 浮点 → 16 位整数

Enum.map(fn s -> round(max(-1.0, min(1.0, s)) * 32_767) end) 把 [-1,1] 浮点钳制到 int16 范围,再 Enum.map(&<<&1::little-signed-16>>) 逐点编码成 2 字节小端,最后 IO.iodata_to_binary/1 把列表拍平成一块二进制。

IO.iodata_to_binary 是处理「字节列表」的标准出口:它允许嵌套的列表/二进制混合,最后合并成单块——拼接大量小块时性能远好于字符串 <> 循环。

L57–59 头 <> 数据

header(sample_rate, 1, byte_size(pcm)) <> pcm:44 字节头 + PCM 数据 = 完整 WAV 文件。

这行展示了「数据长度已知」时的正常路径——与第 6 章 Recorder「长度未知、先占位再回填」正好互补,两种场景都覆盖到了。

lib/voice_chat/prompt_audio.ex· 源码全文 · 提示音文件管理 0/0
defmodule VoiceChat.PromptAudio do
  @moduledoc """
  管理每 3 秒推送给客户端的提示音(“好的”)音频字节。

  优先使用 `priv/static/audio/haode.wav`(可自行替换为自己的录音),
  不存在时在启动时自动合成一段双音提示音作为兜底。
  """

  @doc "确保提示音文件存在;不存在则生成兜底音频。"
  def ensure! do
    path = path()
    File.mkdir_p!(Path.dirname(path))

    unless File.exists?(path) do
      File.write!(path, VoiceChat.Wav.synthesize_chime(16_000))
      require Logger
      Logger.info("generated fallback prompt audio at #{path}")
    end

    :ok
  end

  @doc "提示音文件完整字节(每次调用时读取)。"
  def bytes, do: File.read!(path())

  @doc "提示音文件路径。"
  def path, do: Application.app_dir(:voice_chat, "priv/static/audio/haode.wav")
end
L1–13 模块职责与 ensure!/0:解析路径并建目录

模块文档先声明“现有素材优先、缺失时合成”的职责。ensure!/0 在运行时调用 path/0,再对其父目录执行 File.mkdir_p!;它不依赖当前 shell 的工作目录。

路径通过函数在运行时求值,而不是把构建机上的绝对 priv 路径固化进模块属性;这样 release 安装位置变化时仍能解析当前应用目录。

关联:第 1 章 priv/ 目录;第 2 章 app: :voice_chat。

L14–22 没有素材就现场造一个

若 path 不存在,就用 VoiceChat.Wav.synthesize_chime(16_000) 生成兜底 WAV、写入并记录日志;存在时保留仓库或部署者提供的素材。

「现有素材优先、合成兜底」:仓库已带 haode.wav;缺失时在应用启动阶段生成。目录不可写或生成失败会让启动尽早报错,而不是等第一条会话才暴露。

关联:第 3 章 L10–11;第 7 章 start 时 bytes()。

L23–28 bytes/0 与 path/0

bytes/0 通过动态路径读取完整文件;Prompter 是每个会话启动时调用一次,不是每 3 秒调用。path/0 用 Application.app_dir/2 解析当前应用的 priv/static/audio/haode.wav。

把路径解析封装在模块里,调用方只依赖“给我提示音 bytes”这个语义接口。

为什么 WAV 要「先占位、最后回填」

Recorder 开始流式写入时还不知道最终 PCM 长度,而 WAV 头要求提前写 RIFF/data 大小,所以先写 0、结束时 position(0) 回填。这是一类“头部依赖最终长度”格式的常见策略;经典 RIFF 的 32 位长度也把单个 data 限制在约 4 GiB 以内,长录音应分段或改用 RF64 等格式。

第 9 章

前端:浏览器里的音频管线

麦克风 → PCM16 或 Opus → WebSocket 二进制帧,整条管线都由浏览器原生 API 组成。

前端 app.js 使用浏览器原生的 WebSocket、Media Capture、Web Audio 与 MediaRecorder API 完成采集、编码、发送和播放。页面是语音调试台与消息日志,不是带文本输入框的完整聊天产品;本章重点讲清二进制帧如何与后端落盘格式对齐。

priv/static/app.js · L72–138连接、身份防护与断开 0/0
els.connect.onclick = () => {
  if (ws && (ws.readyState === WebSocket.OPEN || ws.readyState === WebSocket.CONNECTING)) return;

  bytesUp = 0;
  bytesDown = 0;
  frames = 0;
  sessionStart = 0;
  sessionId = "sess_" + Date.now() + "_" + Math.floor(Math.random() * 1e6);
  const proto = location.protocol === "https:" ? "wss" : "ws";
  const url = `${proto}://${location.host}/ws?session=${sessionId}`;

  setStatus("连接中…", null);
  els.connect.disabled = true;
  addMessage("system", `正在连接 ${url}`);

  const socket = new WebSocket(url);
  ws = socket;
  socket.binaryType = "arraybuffer";

  socket.onopen = () => {
    if (ws !== socket) return;
    setStatus("已连接", true);
    els.session.textContent = "会话: " + sessionId;
    els.disconnect.disabled = false;
    addMessage("system", "✅ 已连接,正在打开麦克风…");
    void startRecording(socket);
  };

  socket.onmessage = (ev) => {
    if (ws !== socket) return;
    if (typeof ev.data === "string") {
      handleText(ev.data);
    } else {
      bytesDown += ev.data.byteLength;
      playAudio(ev.data);
    }
  };

  socket.onclose = () => {
    if (ws !== socket) return;
    ws = null;
    void stopRecording();
    setStatus("已断开", false);
    els.session.textContent = "";
    els.connect.disabled = false;
    els.disconnect.disabled = true;
    addMessage("system", "🔌 连接已断开");
    addMessage("system", "💾 服务端将尝试收尾本次录音;结果与保存路径以服务端日志为准");
  };

  socket.onerror = () => {
    if (ws === socket) addMessage("error", "⚠️ WebSocket 错误");
  };
};

els.disconnect.onclick = async () => {
  const socket = ws;
  if (!socket) return;
  els.disconnect.disabled = true;
  await stopRecording();
  if (ws === socket && socket.readyState < WebSocket.CLOSING) socket.close();
};

window.addEventListener("beforeunload", () => {
  void stopRecording();
  if (ws) ws.close();
});
L72–90 每次连接捕获自己的 socket

重置统计、生成候选 session id,再按当前页面协议/主机拼 URL。新建的连接同时保存为局部 socket 与当前全局 ws;binaryType = "arraybuffer" 让下行音频可直接交给后面的解码函数。

异步回调闭包捕获具体 socket,后续可用 ws === socket 判断事件是否仍属于当前会话;这比在旧回调里直接读可变的全局 ws 更能抵御重连竞态。

L91–99 onopen:只为当前连接开麦

先用身份 guard 丢弃旧连接的迟到 open 事件,再更新 UI,并把当前 socket 显式传入 startRecording(socket)。

getUserMedia 的核心前提是安全上下文和用户授权;它并非普遍要求必须在点击回调的同一调用栈执行。播放侧的 AudioContext.resume() 还会受到各浏览器自动播放/用户激活策略影响。

L100–109 onmessage:先验明会话,再分文本/音频

字符串帧交给 handleText(解析 JSON 并写入消息日志);二进制帧统计下行字节并交给 playAudio。

客户端出站与服务端入站、服务端出站与客户端入站是两份互补契约:这个项目两向都用文本承载 JSON、二进制承载音频,但消息类型并非一一镜像。

L110–126 close / error:旧事件不得改新 UI

只有当前 socket 的 close 才先把 ws 清为 null、启动录音清理并重置 UI;旧 socket 的迟到 close/error 直接忽略。

页面里的“服务端将尝试收尾”是操作提示,不是服务端 ACK;录音是否成功仍以服务端日志和文件系统为准。

L127–138 主动断开可等待,卸载只能尽力而为

点击断开时先快照 socket,await stopRecording(),再仅在它仍是当前且尚未 closing 时关闭;页面卸载则发起清理并立即 close,无法等待异步完成。

主动路径给 MediaRecorder 最终 dataavailable 一次在 socket 仍 OPEN 时发送的机会,但 1 秒客户端等待、网络发送队列与缺少服务端 ACK 意味着它仍不是交付保证。

priv/static/app.js · L142–239socket 绑定的 PCM16 / Opus 采集 0/0
async function startRecording(socket) {
  const fmt = els.format.value;
  if (fmt === "opus") await startMediaRecorder(socket);
  else await startPcmCapture(socket);
}

function sendConfig(socket, fmt, container) {
  if (ws !== socket || socket.readyState !== WebSocket.OPEN) throw new Error("WebSocket 已关闭");
  const config =
    fmt === "opus"
      ? { type: "config", format: "opus", container }
      : { type: "config", format: "pcm16", sample_rate: 16000, channels: 1, bit_depth: 16 };
  socket.send(JSON.stringify(config));
}

async function startPcmCapture(socket) {
  try {
    ensureAudioCtx();
    await audioCtx.resume();
    if (ws !== socket) return;

    const stream = await navigator.mediaDevices.getUserMedia({
      audio: { echoCancellation: true, noiseSuppression: true, sampleRate: 16000 },
    });
    if (ws !== socket || socket.readyState !== WebSocket.OPEN) {
      stream.getTracks().forEach((track) => track.stop());
      return;
    }

    micStream = stream;
    sendConfig(socket, "pcm16");
    sourceNode = audioCtx.createMediaStreamSource(stream);
    processor = audioCtx.createScriptProcessor(4096, 1, 1);

    const targetRate = 16000;
    processor.onaudioprocess = (e) => {
      const input = e.inputBuffer.getChannelData(0);
      const pcm = downsample(input, audioCtx.sampleRate, targetRate);
      const buf = floatTo16BitPCM(pcm);
      if (ws === socket && socket.readyState === WebSocket.OPEN) {
        socket.send(buf);
        bytesUp += buf.byteLength;
        frames++;
      }
    };

    // 静音增益保持音频图活跃,避免回声反馈
    silenceNode = audioCtx.createGain();
    silenceNode.gain.value = 0;
    sourceNode.connect(processor);
    processor.connect(silenceNode);
    silenceNode.connect(audioCtx.destination);

    sessionStart = Date.now();
    addMessage("system", "🎤 麦克风已开启(PCM16 @16kHz)—— 开始说话吧");
  } catch (err) {
    if (ws === socket) {
      await stopRecording();
      addMessage("error", "⚠️ 无法获取麦克风: " + err.message);
    }
  }
}

async function startMediaRecorder(socket) {
  try {
    const mime = ["audio/webm;codecs=opus", "audio/ogg;codecs=opus"].find((m) =>
      MediaRecorder.isTypeSupported(m)
    );
    if (!mime) throw new Error("当前浏览器不支持 WebM/Ogg Opus 录音");

    const stream = await navigator.mediaDevices.getUserMedia({ audio: true });
    if (ws !== socket || socket.readyState !== WebSocket.OPEN) {
      stream.getTracks().forEach((track) => track.stop());
      return;
    }

    micStream = stream;
    mediaRecorder = new MediaRecorder(stream, { mimeType: mime });
    sendConfig(socket, "opus", mime.startsWith("audio/ogg") ? "ogg" : "webm");

    mediaRecorder.ondataavailable = (e) => {
      if (e.data && e.data.size > 0 && socket.readyState === WebSocket.OPEN) {
        socket.send(e.data);
        bytesUp += e.data.size;
        frames++;
      }
    };

    mediaRecorder.start(250); // 每 250ms 一个分块
    sessionStart = Date.now();
    addMessage("system", "🎤 麦克风已开启(Opus / " + mime + " 分块)");
  } catch (err) {
    if (ws === socket) {
      await stopRecording();
      addMessage("error", "⚠️ 无法获取麦克风: " + err.message);
    }
  }
}
L142–156 分派格式,并把 config 绑定到 socket

按页面选择把同一个 socket 传给 Opus 或 PCM 分支。sendConfig 同时要求它仍是当前 ws 且状态为 OPEN,再声明 PCM 参数或实际 Opus 容器。

config 不一定是采集函数的第一步,但必须先于首个音频分块到达 Recorder;显式 socket 参数还防止旧异步流程把配置发到新会话。

L157–175 PCM:两次身份检查夹住权限等待

恢复 AudioContext 后先检查 socket 身份,再等待 getUserMedia。权限结果返回后再次确认身份与 OPEN;若已经过期,立即停止刚拿到的轨道,不让旧请求覆盖新会话资源。

await 期间外部状态可以改变。前后复核加上局部 stream,正是处理“用户授权很晚才回来”竞态的关键。

L176–187 PCM 回调:近似降采样后发送

上一段拿到 stream 后创建 MediaStreamSource 与 ScriptProcessor(4096, 1, 1);本段回调按 audioCtx.sampleRate 降到 16kHz、转 Int16,并且只在 captured socket 仍为当前 OPEN 连接时发送。权限约束里的 sampleRate: 16000 只是请求,浏览器不保证采用。

若 AudioContext 是 48kHz,4096 样本约 85ms。ScriptProcessorNode 已被 Web 标准标为弃用,教学代码胜在直观;生产实现应优先 AudioWorklet,并结合 WebSocket.bufferedAmount 或其他协议做背压。

L188–204 静音节点与当前会话失败清理

处理器经保存下来的零增益节点接到 destination:维持音频图而不回放麦克风。异常只在 socket 仍为当前会话时触发统一清理和错误提示。

保存 silenceNode 让 stop 路径能显式断开它;异常 guard 则避免旧会话错误清掉新会话资源。

L205–221 Opus:探测容器并复核权限结果

依次探测 WebM/Opus 与 Ogg/Opus,再请求麦克风;和 PCM 一样,迟到且已过期的 stream 会立即停轨。确认当前后才创建 MediaRecorder,并把实际容器发给服务端。

Opus 是编码,WebM/Ogg 是容器;后端不会转码,只按这份配置选择后缀并追加浏览器产出的分块。

L222–239 250ms 分块与最终 Blob 通道

ondataavailable 把非空 Blob 发到捕获的 socket,start(250) 请求约每 250ms 产出一块;失败路径同样只清理当前会话。

回调故意使用 captured socket:主动 stop 时 global ws 仍指向它,最终 Blob 可在 close 之前发出。这里只检查 socket OPEN;浏览器分块时机、发送缓冲与网络仍让 250ms 和最终交付都不具备硬实时保证。

priv/static/app.js · L241–326竞态安全清理 / PCM 转换 / 播放 0/0
async function stopRecording() {
  const recorder = mediaRecorder;
  const captureProcessor = processor;
  const captureSource = sourceNode;
  const captureSilence = silenceNode;
  const captureStream = micStream;
  mediaRecorder = null;
  processor = null;
  sourceNode = null;
  silenceNode = null;
  micStream = null;
  sessionStart = 0;

  if (captureProcessor) captureProcessor.onaudioprocess = null;

  if (recorder && recorder.state !== "inactive") {
    await new Promise((resolve) => {
      const timeout = setTimeout(resolve, 1_000);
      recorder.addEventListener("stop", () => { clearTimeout(timeout); resolve(); }, { once: true });
      try { recorder.stop(); } catch (_) { clearTimeout(timeout); resolve(); }
    });
  }

  if (captureProcessor) {
    try { captureProcessor.disconnect(); } catch (_) {}
  }
  if (captureSource) {
    try { captureSource.disconnect(); } catch (_) {}
  }
  if (captureSilence) {
    try { captureSilence.disconnect(); } catch (_) {}
  }
  if (captureStream) {
    captureStream.getTracks().forEach((track) => track.stop());
  }
  // audioCtx 保留给播放用,不关闭
}

/* ---------- 音频处理 ---------- */

function downsample(input, fromRate, toRate) {
  if (fromRate === toRate) return input;
  const ratio = fromRate / toRate;
  const newLen = Math.round(input.length / ratio);
  const out = new Float32Array(newLen);
  let offset = 0;
  for (let i = 0; i < newLen; i++) {
    const end = Math.round((i + 1) * ratio);
    let sum = 0, count = 0;
    for (; offset < end && offset < input.length; offset++) {
      sum += input[offset];
      count++;
    }
    out[i] = count > 0 ? sum / count : 0;
  }
  return out;
}

function floatTo16BitPCM(float32) {
  const buffer = new ArrayBuffer(float32.length * 2);
  const view = new DataView(buffer);
  for (let i = 0; i < float32.length; i++) {
    const s = Math.max(-1, Math.min(1, float32[i]));
    view.setInt16(i * 2, s < 0 ? s * 0x8000 : s * 0x7fff, true);
  }
  return buffer;
}

function ensureAudioCtx() {
  if (!audioCtx) audioCtx = new (window.AudioContext || window.webkitAudioContext)();
  return audioCtx;
}

function playAudio(arrayBuffer) {
  const ctx = ensureAudioCtx();
  ctx.resume()
    .then(() => ctx.decodeAudioData(arrayBuffer.slice(0)))
    .then((buffer) => {
      const src = ctx.createBufferSource();
      src.buffer = buffer;
      src.connect(ctx.destination);
      src.start();
      addMessage("out", "🔊 播放服务端音频");
    })
    .catch((err) => addMessage("error", "⚠️ 音频解码失败: " + err.message));
}
L241–280 stop:先快照并清空,再跨越 await

把 MediaRecorder、三个音频节点和 stream 快照到局部变量,立刻清空全局引用并禁用旧 PCM 回调;随后最多等 1 秒让 MediaRecorder stop,再断开快照节点并停止轨道。AudioContext 留给下行播放复用。

关键顺序是“同步夺走所有权,再 await”:等待期间即使新会话开始,旧 stop 也只会清理自己的快照,不会碰到新写入的全局资源。

L281–298 降采样:48kHz → 16kHz

按 fromRate / toRate 对输入区间求平均,得到目标长度。这是便于阅读的近似抽取,不是完整的抗混叠重采样器。

从 48kHz 到 16kHz 时样本数约降到 1/3;未先做严格低通滤波会有混叠风险,因此这里只适合演示/调试,不能声称语音信息无损。

L299–313 Float32 → Int16 小端 + AudioContext

DataView.setInt16(i*2, …, true) 把 [-1,1] 浮点钳制为 16 位有符号小端整数;后端不逐样本转码,WAV 头也必须声明同样的位宽和声道约定。ensureAudioCtx 则惰性复用播放/采集上下文。

字节序、符号或位宽不一致都会造成失真或噪音;最后一个 true 正是“小端”开关。

L314–326 playAudio:解码并创建一次性 source

恢复 AudioContext、复制后交给 decodeAudioData,为每条提示音创建 AudioBufferSourceNode 并立即播放;失败则写错误日志。

独立 source 不会自动排队,上一段尚未结束时会重叠。当前提示短于推送间隔,但产品代码仍应明确选择允许重叠、排队或打断。

priv/static/app.js · L343–370服务端文本协议与实际 session id 0/0
function handleText(raw) {
  let msg;
  try { msg = JSON.parse(raw); } catch (_) { return; }

  switch (msg.type) {
    case "session": {
      const actual = String(msg.session || "");
      if (actual) {
        const renamed = sessionId && sessionId !== actual;
        sessionId = actual;
        els.session.textContent = "会话: " + actual;
        if (renamed) addMessage("system", `ℹ️ 会话名冲突,服务端已改用 ${actual}`);
      }
      break;
    }
    case "prompt":
      addMessage("in", `🔊 服务端提示音: “${msg.message || "好的"}”`);
      break;
    case "text":
      addMessage("in", msg.message || msg.text || "");
      break;
    case "error":
      addMessage("error", "⚠️ " + (msg.message || "服务器错误"));
      break;
    default:
      addMessage("in", "← " + raw);
  }
}
L343–357 session:服务端 id 才是最终值

解析 JSON 后,session 消息会用服务端实际注册的 id 覆盖本地候选值并更新 UI;若发生重名改写,还会显示说明。

Registry 的唯一键冲突只能由服务端裁决。初始化首帧把最终 id 回传,避免页面、Registry 与录音文件名各自显示不同值。

L358–370 prompt / text / error:出站类型分派

prompt、text 写进消息日志,error 用错误样式,未知类型保留原始 JSON;无法解析的文本帧则忽略。

它与服务端处理“客户端入站消息”的 case 是两个方向的互补契约,并非同一组 type 的一一镜像;协议变更要同时审视两端。

48kHz(麦克风原始) → downsample(…, 48000, 16000) 16kHz 样本数约 1/3
图 9-1 · 以 48kHz → 16kHz 为例,输出样本数约为 1/3。图只示意区间平均;当前实现没有完整抗混叠滤波,不能据此保证波形或语音信息无损。
最容易踩的坑:浏览器安全策略

getUserMedia 只暴露在安全上下文中:HTTPS,以及浏览器认定可信的本地来源(通常包括 localhost)。普通 HTTP 的局域网 IP 通常不满足条件;此外用户仍需授权,播放用 AudioContext 还可能受自动播放策略限制。

第 10 章

串联:Elixir 概念在 voice_chat 里的投影

最后一章,把前面十章散落的点连成一张概念网。

voice_chat 刻意只选了少量核心机制,并用直接的形式把它们串在一条可运行链路里。本章把六个概念抽出来对照,帮助你建立阅读 Elixir 服务端代码所需的基本坐标;它不是 OTP 全貌,也不是所有项目都应照搬的架构模板。

一张图:一次会话的角色与生命周期关系

VoiceChat.Supervisor 直接 child spec:Registry、Bandit Bandit 管理 HTTP / WS 连接 Registry session_id → ws pid 每条连接的 ws 进程 在连接进程中执行 WSHandler 回调 第 5 章 动态连接生命周期 注册 / 查询关系 Recorder 落盘 Prompter 约 3 秒一轮 spawn_link ×2 浏览器一侧(不是 BEAM 进程) getUserMedia / MediaRecorder WebSocket 客户端 AudioContext / decodeAudioData 第 9 章 config / 音频上行 提示文本 / 音频下行
图 10-1 · 这是应用代码所关心的概念角色与生命周期关系,不是精确进程计数:Registry、Bandit、Thousand Island 与运行时的内部进程树均未展开。业务代码没有自行维护 OS 线程池;浏览器也不属于 BEAM 进程树。

六个概念,六次投影

Elixir 概念在 voice_chat 里的投影一句话本质
进程ws 进程、Recorder、Prompter、Supervisor…(图 10-1)BEAM 调度的轻量执行单元,通常把变化中的业务状态留在进程内部
消息传递send 的 {:chunk, data} / {:send_binary, audio} / :stop本项目用异步邮箱解耦三类会话角色;同一发送者到同一接收者保持发送顺序
link / trap_exitWSHandler 两次 spawn_link,形成 ws + 两个帮手;Recorder 捕获退出信号link 关联生命周期;trap_exit 可把部分退出信号转成消息处理
Application / Supervisorapplication.ex(第 3 章)「系统进程」的启动顺序与生死责任
行为(behaviour)本模块实现 WebSock 的三个必需回调与可选 terminate/2;可选 handle_control/2 未实现框架把连接事件翻译成约定的回调,协议细节由适配层承接
不可变数据 + 模式匹配state map 每轮返回新值、receive/case 按形状分派各会话进程的堆/业务状态不直接共享;Registry、ETS 等机制可封装共享查找状态
「轻量进程」到底有多轻

BEAM 进程不是操作系统线程,创建与切换通常都很轻,运行时调度器会安排可运行进程。像 Prompter 这样阻塞在 receive 的进程不会持续占用 CPU 时间片,但仍持有堆、邮箱、link 与定时等待等资源,并非“零开销”。因此每连接增加两个职责单一的进程在这个示例里很自然;真实容量仍要通过负载测试、邮箱监控和背压设计验证。

消息传递的三个直觉

隔离:Elixir 数据不可变,接收方不会拿到一块可被发送方随后改写的业务数据;普通项在语义上按值传递,较大的 binary 等可由运行时引用计数共享。 有序:只保证同一发送者发往同一接收者的消息按发送顺序到达,不保证多个发送者之间的全局顺序。 邮箱:消息进入接收者邮箱,receive 可按模式选择匹配项;本项目选择邮箱通信,但进程还可能通过调用、ETS、端口等其他机制协作。

OTP 工具不是单向的“轻重等级”

Recorder/Prompter 展示了 spawn_link + receive;项目还实际使用了 Supervisor 与 Application。生态中另有 Task(有结果或一次性的工作)和 GenServer(标准化状态服务器与 call/cast 生命周期)等抽象,本项目并未逐一演示。选择哪一种要看重启责任、应答需求、超时、可观测性与测试方式;短代码不一定比标准 OTP 抽象更适合生产。

项目还能怎么长

  • 上鉴权:upgrade 之前检查 token(第 4 章那个「升级前检查」的位置就是预留的钩子);
  • 加静音检测 / 语音活动检测:增加音频分析环节,再定义检测结果如何进入 Recorder 或独立进程;还要处理窗口、阈值和消息背压,不只是多写一个匹配分支;
  • 把 Recorder 换成 GenServer:当它需要应答(比如查询已录字节数)时再升级;
  • 接 Phoenix:Channels 同样围绕长连接、主题与进程隔离组织消息,但还提供自己的传输、序列化、PubSub 与生命周期抽象,不能简单等同为本项目的放大版;
  • 上 Nerves / 嵌入式:消息与进程拆分思路可以复用,但部署、音频设备接入、浏览器前端和资源约束都需要重新适配。
结业自测:不看代码,你能答出这些吗

① 音频块从浏览器到硬盘,依次经过哪几个进程?
② 为什么 Prompter 每 3 秒发消息不需要「注册回调」?
③ 常规断开或 ws 进程异常退出时,Recorder 如何尽量收尾?它覆盖不了哪些故障?
④ {:push, {:binary, data}, state} 里的 data 最终怎么变成浏览器喇叭里的声音?

答案:① 浏览器 → ws 进程(handle_in 转发)→ Recorder 进程(write_chunk)→ 文件。② Prompter 创建时已拿到 ws pid,直接向其邮箱发送消息即可;Registry 用于按 session id 查连接,不是这条回传路径的必要条件。③ 常规 terminate 会发 :stop;ws 的非正常退出信号也可被 linked Recorder trap 后触发 finalize。直接 :kill Recorder、VM/主机掉电、磁盘错误等不在保证范围。④ data 经 WebSocket 二进制帧到浏览器,decodeAudioData 解码后由 AudioBufferSourceNode.start() 播放(第 9 章)。
最后一句:这个项目「为什么这么写」的答案

这个项目选择用连接进程、两个 linked 帮手和消息把网络、落盘、周期提示分开,因为这种职责隔离很适合展示 BEAM 的进程模型。它是一个清晰的教学取舍,不是普遍性能结论:生产系统通常还要补上更明确的监督、背压、限流、可观测性与故障恢复,再根据生命周期决定保留裸循环还是采用标准 OTP 抽象。