全景:一次语音会话的完整旅程
先看全貌,再看代码。
voice_chat 很小——后端 9 个项目模块、前端 3 个页面文件和 1 个音频资源——但它把 Elixir 里最重要的几件事串了起来:进程、消息传递、OTP 监督树、Plug 管线、WebSocket 协议。本章用一张会动的图,带你走一遍「你对着麦克风说话」到「服务端落盘 + 约每 3 秒推送一次提示文本与音频」的全程。
一张表记住「需求 → 代码 → 章节」
| 你的需求 | 落在哪 | 对应章节 |
|---|---|---|
| 启动 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 停掉两个子进程
open 回调则开始 b;这两边属于并发的握手后工作,不能要求 c 一定先完成。服务端在处理后续帧前会先完成 init。随后音频块上行并触发 d;Prompter 从 init 创建后约 3 秒执行 a,浏览器收到音频后执行 e。断开时进入 f。麦克风授权、持续落盘和周期提示会交错进行,因此不存在把 a、d 排成唯一全局顺序的答案。mix new:一个 Elixir 项目从哪来
先认识你脚下的脚手架,再看每一块地砖。
整个项目从 mix new voice_chat --sup 提供的基础骨架长出来。本章先区分哪些文件由命令生成、哪些是项目后续加入,再搞清楚 lib/、priv/、config/、test/、mix.exs、mix.lock 各由谁在什么时候使用。
基础骨架 vs 后续新增
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 定位。mix.exs 逐行:项目的出生证明
全文只有 29 行,但每一行都决定了后面所有章节的写法。
这是本书第一个代码走读。走读组件左侧是带行号的源码,右侧是讲解卡片;点代码里的任意一行,或点卡片标题,或按 ▶ 逐步播放,讲解会跟着行号走。代码块可用 python3 gen.py --check 对照 voice_chat/ 源码校验。
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
use Mix.ProjectVoiceChat.MixProject 是生成器采用的命名约定;真正把模块声明为 Mix 项目的是 use Mix.Project。它设置项目行为和编译钩子,项目再实现 project/0,并可选实现 application/0。
Mix 在求值 mix.exs 时装载这个模块并读取返回的配置;关键不是模块名被“搜索到”,而是 use Mix.Project 完成注册和约定接线。
关联:第 1 章「Mix 项目」、第 3 章「OTP 应用」——这一行同时定义了两者。
project/0:一段关键字列表project/0 返回一个关键字列表(keyword list),是 mix 读取项目信息的主入口。第 5 行的 [ 开始这个列表。
这里体现「配置就是数据」:本项目的函数只返回关键字列表。Elixir 函数本身并不会被语言强制为纯函数,因此“无副作用”是这段实现的选择,不是 Mix 的硬保证。
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 贯穿配置与资源定位。
version: "0.1.0"当前项目自身的版本号,会进入应用元数据并用于发布/打包;发布到 Hex 时通常遵循语义化版本。它不决定依赖锁定,mix.lock 记录的是依赖解析结果。
elixir: "~> 1.15" —— 版本约束要求 Elixir 版本。重点是 ~> 的语义:~> 1.15 允许 1.15.x 及以上、但不到 2.0 的版本(即允许 1.19,不允许 2.0)。
用宽松下限 + 语义化版本,既不锁死环境,又不放进来不兼容的大版本——这是 Hex 生态的惯例。
关联:第 1 章 mix.lock(依赖的版本锁得比这严格得多)。
start_permanent: Mix.env() == :prod当 Mix 启动应用且 MIX_ENV=prod 时,把该应用标为 :permanent;如果整个 OTP 应用最终停止,BEAM 节点会退出。子进程的普通崩溃仍先按 Supervisor 策略重启,并不是任意异常都立即关 VM。
它让“应用已无法维持监督树”成为节点级失败信号,便于外部服务管理器重启实例;开发环境不启用这个永久启动标记。
deps: deps() —— 依赖把私有函数 deps/0 返回的依赖列表挂到 project 上。deps/0 定义在第 21 行。
拆成私有函数是惯例:project 保持可读,依赖单独成块方便增删。
application/0:运行期怎么启动extra_applications: [:logger] 声明「我的应用启动前,请先把 :logger 应用启动好」——因为代码里要 Logger.info。
BEAM 的应用有依赖顺序:Logger 是 Elixir 自带的库应用,也要显式声明才保证可用。这也是为什么依赖图是「应用」,不只是「包」。
mod: {VoiceChat.Application, []} —— 应用入口告诉应用控制器:启动本应用时调用 VoiceChat.Application.start/2,参数为 []。若回调模块实现可选的 stop/1,应用停止时才会调用;本项目没有实现它。
没有 mod: 的应用不会启动自己的 Application 回调(仍可打包模块、配置和依赖);有 mod: 时应用控制器会调用入口,通常由它拉起进程树。voice_chat 属于后一种。
关联:第 3 章整章——start/2 里就是监督树。
deps/0:声明式依赖返回依赖元组列表。每个元组 {名字, 版本约束} 只是声明——真正下载、编译发生在 mix deps.get / mix compile 时。
声明式的好处是机器可读、可锁定(mix.lock)并容易在团队和 CI 中复现;Mix/Hex 负责常规依赖获取与校验。
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 编码。
end模块结束。这里没有服务运行期的请求处理逻辑;这些函数会在构建/启动准备阶段被 Mix 调用,返回项目与应用元数据。
这提醒我们一个事实:Elixir 的「代码即数据」让项目配置本身也是一段可以被任何工具读取的 Elixir 表达式。
什么时候会动这个文件
加依赖(往 deps/0 加一行)、改应用入口(mod:)、加编译期依赖或 escript 配置。其他时候基本不动——就像出生证明,偶尔才拿出来改。
监督树:进程从哪来,谁负责谁
应用自己的顶层 Supervisor 只有两个直接 child spec,但每个组件内部还会管理更多进程。
依赖应用启动后,应用控制器按第 2 章的 mod: 调用 VoiceChat.Application.start/2。它准备资源、声明 Registry 与 Bandit 两个直接孩子,再启动顶层 Supervisor;Bandit/Thousand Island 继续管理监听器与连接生命周期。
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
use Applicationuse Application 要求实现 start/2(和可选的 stop/1、prep_stop/1)。@moduledoc false 表示这个模块是内部实现细节,不生成文档。
OTP 应用的启动/停止是「框架回调」:VM 的应用控制器在正确时机调用 start/2,你只负责返回 {:ok, pid} 或 {:ok, pid, state}。
关联:第 2 章 L17 mod: {VoiceChat.Application, []}——这行就是接线。
@impl true + start/2@impl true 声明「下面是行为回调的实现」,编译器会校验签名是否匹配。_type 是启动类型(常见为 :normal,也可能是 takeover/failover 元组),_args 是 mod: 里传的参数 [],这里都用不上。
下划线前缀的变量名是「我故意不用你」的约定,避免未使用变量的编译警告。
VoiceChat.recordings_dir() 从配置读目录(priv/recordings),File.mkdir_p!/1 递归建目录。|> 是管道:左边结果作为右边函数第一个参数。
结尾 ! 是 Elixir 约定:失败直接抛异常(而非返回 {:error, …})。这一步发生在顶层 Supervisor 启动之前,因此失败会让 Application 的启动回调失败,错误会在启动期暴露。
关联:第 6 章 Recorder 的 build_path 写文件时假设这个目录已存在。
VoiceChat.PromptAudio.ensure!() 检查 priv/static/audio/haode.wav 是否存在,不存在就现场合成一段(第 8 章)。
启动期能尽早发现“文件缺失且兜底无法写入”这类错误;但 ensure! 对已存在文件只检查存在性,不验证可读性或 WAV 内容,损坏素材仍可能到 Prompter 读取或浏览器解码时才暴露。
关联:第 7 章 Prompter 启动时 PromptAudio.bytes() 直接读文件。
children:被监督的孩子声明两个直接 child spec:Registry(唯一键会话注册表,条目属于注册它的连接进程,进程退出时自动清理)和 Bandit(HTTP/WS 服务器组件)。它们本身都可能在内部启动多个进程。
放进监督树后,顶层组件异常退出会按策略重启;若超过重启强度,Supervisor 自身仍会终止,因此监督提供的是明确的恢复策略,不是“永不失败”。
关联:第 5 章 Registry.register;第 4 章 Bandit, plug: Endpoint 的 plug 选项。
{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/ 目录。
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 产生的信号。
为什么应用代码不把每条连接列成顶层 child spec
WebSocket 连接短命且动态,断开后也不应按固定 child spec 重启。它们的生命周期由 Bandit/Thousand Island 的服务器层管理;本应用只把长期的服务器组件 Bandit 放进自己的顶层监督树。WSHandler 初始化后,再分别 link Recorder 与 Prompter 这两个会话辅助进程。
类比:监督树 = 值班经理
Supervisor 像值班经理:直接孩子(Registry、Bandit)退出时按策略处理。连接则像由 Bandit 业务线内部调度的短期任务——仍有明确的生命周期管理,只是不逐条出现在本应用的顶层 children 列表里。
Endpoint:Plug 管线和 101 升级
从 HTTP 请求切换到 WebSocket 帧协议,入口是一轮成功的 101 Upgrade 握手。
endpoint.ex 是流量的第一站:静态页面、健康检查、会话列表和 WebSocket 升级都在这个模块里。理解它需要两个核心概念:Plug 管线(请求穿过一串 conn 变换)和 WebSocket 升级(从 HTTP 握手切换到帧协议)。
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
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 选项——服务器把请求交给这个模块。
plug(Plug.Logger)把日志中间件挂进管线,记录请求方法/路径以及响应状态与耗时。
这是本管线的第一个「关卡」:请求先经过它,再到路由。中间件可以更新 conn;需要提前结束后续管线时,还要发送响应并 halt。
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。
:match 与 :dispatchPlug.Router 的两段式::match 找到匹配的路由,:dispatch 执行它。本模块把 Logger、Static 排在两者之前,所以请求会先按声明顺序经过它们。
显式顺序让日志覆盖本模块收到的请求,也让 Static 有机会在路由匹配前直接服务白名单资源。
get "/":首页根路径显式用 send_file 返回 index.html;Plug.Static 负责的是按 URL 文件名请求的 /app.js、/app.css 等资源,并不会自动把 / 当作目录首页。
把首页路由和静态资源白名单分开后,根路径行为清楚,也不依赖目录索引约定。
/health健康检查:send_resp(conn, 200, "ok")。部署时负载均衡/探针打这个接口确认服务活着。
/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(写入)——这里是同一张表的读取。
put_resp_content_type("application/json") 设置 Content-Type,Jason.encode! 把 Elixir 数据结构编码成 JSON 字符串返回。
Jason.encode! 的 ! 表示编码失败会抛异常;前一段显式字符串化已知 pid,正是为了让当前元数据结构可编码。
/ws:先取查询参数fetch_query_params(conn) 解析 URL 里的查询串(?session=xxx)放进 conn.query_params。注意这里手动取——因为 WebSockAdapter.upgrade 之后 conn 就不再是普通 HTTP 请求了,参数得在升级前拿到。
关联:第 5 章 init/1 收到的 %{query_params: …} 就是从这里传进去的。
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 承接。
rescue:升级失败要友好升级请求不合法(缺 Sec-WebSocket-Key 等)时 upgrade 会抛 WebSockAdapter.UpgradeError,这里捕获并回 400。
「让用户看到 400 而不是 500」:协议错误是客户端的问题,不该算服务器内部错误。
match _:兜底 404任何没被上面路由接住、也没被 Static 接住的请求,返回 404。
显式兜底把未匹配路径稳定地变成 404;若没有匹配子句,生成的路由匹配函数会因找不到可用子句而失败,而不是自动替本模块构造这条响应。
/ws 登记升级后,Bandit 让该连接切换到 WebSocket/WSHandler;未匹配路径进入 404 兜底。最容易踩的坑:升级之后 conn 就「死」了
WebSockAdapter.upgrade 返回的 conn 状态是 :upgraded,不能再 send_resp。所以任何「升级前要做的检查」(比如鉴权、取参数)都必须在 upgrade 调用之前完成——这正是我们先把 fetch_query_params 放在前面的原因。
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 推音频给客户端用的通道。
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
@moduledoc 不只是注释——mix docs 会生成文档,而更重要的是一开始就想清楚回调的职责边界。
WebSock 回调的「触发者」各不相同:init 由连接建立触发、handle_in 由客户端帧触发、handle_info 由别的进程发消息触发、terminate 由连接关闭触发。文档把这些写清楚,后面每行代码都有的放矢。
@behaviour WebSock声明「我实现了 WebSock 这个行为」。编译器会按行为定义检查回调并对缺失或不匹配给出警告;本项目实现必需的 init/1、handle_in/2、handle_info/2,以及 terminate/2。
行为(behaviour)类似一份回调接口:适配层把协议事件翻译成函数调用,所以本项目业务模块无需自行实现 RFC 6455 的握手与帧解析。
关联:websock 是间接依赖里的行为定义方,websock_adapter 负责把 Plug/Bandit 的升级接进来。
require Logger引入 Logger 宏(Logger.info/1 等其实是宏,需要 require)。
细节:为什么是 require 不是 import/alias?因为 Logger 的 API 是宏(延迟求值 + 编译期裁剪日志级别),宏必须 require。
init/1:session 的出生连接建立后调用一次。参数是第 4 章 upgrade 时传的 %{query_params: …},用模式匹配直接解包。@impl true 标记这是行为回调。
「每个连接 = 一个新 session」就发生在这里:每连接进程各调一次 init,天然隔离,不需要 session 池、不需要锁。
VoiceChat.SessionId.sanitize(query_params["session"]):缺少或清洗后为空时自动生成;其他输入把非字母数字/下划线/连字符替换为 _,并截到最多 64 个字符。
用户输入永远不可信:id 会拼进文件名(第 6 章 build_path),不清洗就可能写出路径穿越。这是安全习惯,不是功能。
关联:lib/voice_chat/session_id.ex;第 6 章文件名拼接。
ws_pid = self()记下当前进程 pid——就是「这个连接的进程」。它会被交给 Recorder/Prompter,让它们能发消息回来(走 handle_info)。
pid 是定位进程的句柄;send/2 把异步消息放进它的邮箱。它和同步“方法调用”并不等价:没有返回值,也不会替调用方提供背压。
关联:第 6 章 trap_exit 里对比的 from 就是它;第 7 章 Process.alive?(ws_pid)。
先用 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 虚线。
Registry.update_value/3 把刚得到的 recorder pid 加入元数据;随后 Prompter.start/1 创建另一个 link 到 ws 的辅助进程。
Registry 的价值:/sessions 调试接口(第 4 章)能实时列出活跃会话,且进程退出时 Registry 自动清理条目——不需要手动 deregister,少一类「忘记清理」的 bug。
关联:第 3 章 Registry 子进程;第 4 章 Registry.select。
日志记录连接后,服务端把最终 session id 编成 {"type":"session",…},同时构造包含三个 pid 的 state,并从 init/1 返回一次文本 push。WebSock 后续回调都会收到这份 state。
客户端请求的 id 可能因冲突被替换,不能继续把本地候选值当真;首帧把 Registry 与文件名实际使用的 id 回告页面,页面再更新显示。
先注册包含 ws pid 和开始时间的 meta;若唯一键已被占用,就生成新 id 后递归重试。注册项属于当前连接进程,连接退出时 Registry 自动删除。
先注册、再启动依赖 id 的 Recorder,避免了“展示 id 已换新、文件仍用旧 id”的分裂状态。
handle_in:二进制帧 = 音频块客户端发来的二进制帧({data, opcode: :binary})就是一块音频,直接 VoiceChat.Recorder.chunk(state.recorder, data) 转发给录音进程,然后 {:ok, state} 继续。
chunk/2 只是异步 send,所以磁盘写入不在连接进程中执行;代价是当前实现没有背压,Recorder 跟不上时消息会积在邮箱里。
关联:第 6 章 write_chunk;第 9 章前端每次回调读取 4096 个输入采样帧,降采样后发一个二进制块。
Jason.decode 解析 JSON,按 type 分派:config 转发给 Recorder(设置格式/采样率)。
当前协议约定二进制帧承载音频,文本帧承载 JSON 控制消息;这是清晰的类型分工,但仍需由两端共同维护,并不会由 WebSocket 自动验证内容语义。
收到 {"type":"text","text":t} 就回一条 {"type":"text","message":"收到: …"}。当前网页没有文本输入框,这条分支是协议/测试脚本可调用的回显示例,页面只会把收到的文本写进消息日志。
guard(when is_binary(t))挡掉缺字段的脏消息;{:push, msg, state} 是 WebSock 的「立即下发」返回。
关联:第 9 章 handleText;scripts/ws_advanced_test.mjs 会实际发送 text 消息。
未知 type 静默忽略;解析失败的帧打一条 warning。
协议要宽容:客户端(尤其浏览器)偶尔发怪帧,别让一个坏帧杀掉整个连接。
handle_info:别人给我的消息Prompter 每轮等待约 3 秒后,依次发送 {:send_text, json} 和 {:send_binary, audio};这里把它们转换为 WebSock 的 {:push, …} 返回值。
持有连接 pid 的进程可以向其邮箱发送约定消息,无需持有底层 socket 句柄;WSHandler 再把这些消息翻译成 WebSock push。当前项目的周期提示走这条路。
关联:第 7 章 Prompter 的 send 调用——两边是配对的。
terminate:有界等待录音收尾先异步停止 Prompter,再调用 Recorder.stop/1;后者请求 finalize 并等待 Recorder 退出,最长 5 秒。超时只记录 warning,随后连接回调返回。
常规成功路径因此有时间完成 WAV 回填;不过 stop 只观察 DOWN、没有检查退出 reason,也不是持久化 ACK,不能单凭 :ok 证明写盘成功。terminate/2 本身也并非所有故障下都保证执行。
关联:第 6 章 stop/1、monitor 与 finalize;第 0 章旅程图的断开收尾。
Recorder:把音频流变成文件
225 行,用一个 receive 循环把格式选择、文件写入和有界收尾串起来。
Recorder 为了教学直接使用 spawn_link + receive,没有套 GenServer。它展示 link / trap_exit、monitor 等待退出、惰性打开、PCM 帧对齐与 WAV 回填;也保留了裸进程方案缺少标准调试、监督与背压的局限。
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
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 与重启策略。
config/2 与 chunk/2 都把消息放入 Recorder 邮箱并立即返回;它们不等待接收方处理或磁盘写完。
连接进程可以继续处理帧,但当前实现没有端到端背压;若写盘落后,chunk 会在 Recorder 邮箱中累积。
stop/1:monitor + 5 秒有界等待先 monitor Recorder,再发送 :stop,等待对应 :DOWN;5 秒内未退出则清掉 monitor 并返回 {:error, :timeout},但不会强杀 Recorder。
monitor 先于 stop,避免错过极快退出。来自同一个 ws 进程的 config、chunk、stop 对同一 Recorder 保持发送顺序,因此常规路径会先处理此前入队的音频,再尝试 finalize。DOWN 只证明进程已退出;代码忽略 reason,不能据此断言收尾成功。
Recorder 启动后立刻启用 trap_exit,再进入带有 session id、ws pid、格式、采样参数、文件句柄、路径、字节数与 finalized 标志的 state map。
默认格式是 PCM16 / 16kHz / 单声道,文件保持 nil 以便惰性打开。trap_exit 能处理 linked ws 的可传播退出信号,但无法抵御直接 :kill Recorder、VM/主机掉电或磁盘故障。
receive:四路状态转移config 更新 state,chunk 写入后递归,:stop finalize 并退出;{:EXIT, from, …} 只在来自所属 ws pid 时收尾,其他退出消息继续循环。
尾递归把新 state 带入下一轮;没有匹配消息时进程阻塞等待,不占用可运行时间片,但仍保留内存与邮箱资源。
先把外部格式归一化,再从 JSON 取采样率与声道;valid_wav_params?/2 失败时记录 warning,并保留当前采样参数。
这些字段最终会进入固定宽度 WAV 头。对正整数、block align 与 byte rate 做边界检查,可避免不可信 config 在位语法构造时溢出或触发异常。
若 state.io 已存在,就警告并保持原 state;尚未打开时才写入归一化后的格式与 WAV 参数。
PCM 头语义和容器扩展名在首次打开时已经确定,中途改格式会让同一文件前后矛盾,所以协议要求 config 先于首个音频块。
guard 检查采样率、声道数、block align 和 byte rate 能装入 RIFF/WAV 对应字段;容器函数把 Opus + Ogg 映为 :ogg,其余受支持的 Opus/WebM 映为 :webm,未知配置回到 PCM16。
Opus 是编码,WebM/Ogg 是容器;分开记录后,后缀才能对应浏览器实际选择的 MIME 容器。
第一块到来而 io 仍是 nil 时,先 ensure_open;成功才递归写当前块,失败则返回未打开 state。后续块直接 :file.write 并累加字节数。
没有音频的连接不会留下空文件。发送方不等待这次写盘,因此邮箱增长仍是需要监控的容量边界。
以 raw binary 写模式打开唯一路径;PCM16 先写 data size 为 0 的 44 字节 WAV 头,WebM/Ogg 不加 WAV 头,然后把句柄与路径写回 state。
流式 PCM 事先不知道最终长度,所以先占位、结束时回填。容器分块已经由 MediaRecorder 编码,服务端只追加其字节。
权限、路径或磁盘问题使 open 返回 error 时,只记日志并保留 io: nil;外层因此丢掉当前块,下一块到来才再次尝试。
这样不会对同一 chunk 无限递归。策略是“连接继续、录音尽力而为”;生产环境还应配合告警、退避或关闭会话。
按 :pcm16 / :webm / :ogg 选择后缀,用已清洗 session id、微秒时间戳和单调唯一整数组成文件名,再拼到配置的录音目录。
时间戳便于追踪,System.unique_integer 在同一 VM 内消歧;它不是跨节点全局 id。
finalize_and_exit/1 调用收尾函数后返回 :ok,顶层匿名函数结束,Recorder 随之退出,monitor 观察到 DOWN。
保存结果以日志与文件系统为准;即将关闭的 WebSocket 不是可靠的“录音完成”确认通道。
已 finalized 或从未打开文件时直接返回;否则按 PCM16 与 WebM/Ogg 分派到各自收尾函数,最后标记 finalized 并清空句柄。
当前循环首次收尾后即退出,这些子句仍让 finalize/1 对“已完成/未打开”状态有明确结果。
每帧大小是 channels × 2 字节。若总字节数不能整除帧大小,就把文件定位到最后一个完整帧,截断残余字节并记录 warning。
WAV 的 block align 必须与数据对齐;宁可丢掉不足一帧的尾巴,也不能让头部宣称的数据结构与文件实际字节矛盾。
文件指针回到 0,用完整帧字节数重建 WAV 头,再关闭句柄并记录路径;返回的 state 也改成截断后的真实字节数。
:file.position/2 返回 {:ok, position},代码按这个形状匹配。经典 RIFF 的 32 位长度字段要求单个 PCM data 小于约 4 GiB;Wav.header/3 会拒绝越界,本项目不自动分段,长录音需另做滚动文件或 RF64。
关联:第 8 章 Wav.header/3 的大小 guard。
容器模式不写或回填 WAV 头,只关闭句柄并记录累计字节数。
服务器不解码 Opus,只保存同一次 MediaRecorder 会话产生的连续容器分块;文件有效性仍取决于浏览器产出的流。
为什么这里不用 GenServer
这个教学项目用裸进程,是为了把邮箱、模式匹配和 trap_exit 完整展开。GenServer 不负责重启策略(Supervisor 才负责),但它提供标准的 call/cast、系统消息、调试与可观测接口,生产代码通常更易维护;Recorder 若需要查询状态、背压、统一遥测或独立监督,应优先评估 GenServer/动态监督,而不是把“裸进程更轻”当成选型口诀。
Prompter:每 3 秒的轻量定时器
Prompter 选择 receive … after,实现“每轮等待 3 秒再推送”。
需求是「接收音频的同时,由另一个轻量进程周期性推送提示」。本实现无需额外调度库,但 receive … after 不是唯一写法;本章也会对比 Process.send_after,并说明当前周期会包含每轮发送耗时,因此是约 3 秒而非墙钟级精确定时。
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
一句话说清职责:本进程不碰 socket,只负责定时往 ws 进程发两条消息(文本 + 音频),下发由第 5 章的 handle_info 完成。
职责分离的边界很干净:Prompter 管「什么时候发」,WSHandler 管「怎么发出去」。
@interval 3_000:模块属性即常量模块属性(@name value)在编译期被替换成字面量。数字下划线(3_000)只是可读性,等于 3000 毫秒。
把「魔法数字」提成具名属性:改间隔只改一行。这比在循环里写死 3000 好维护得多。
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 哲学。
stop/1往进程发 :stop——进程收到后从 :stop -> :ok 分支退出。
Prompter.stop 是 fire-and-forget;Recorder.stop 也发送 :stop,但额外用 monitor 等待其退出。两者都避免直接强杀,只是确认语义不同。
receive do :stop -> :ok after @interval -> … end 会选择性接收 :stop;若 3000ms 内没有匹配消息,就执行 after 分支。
等待期间进程不会占用调度器执行时间,但仍保留堆、邮箱和计时状态。每次发送完成后才重新进入 receive,所以处理耗时会叠加到下一个周期。
if Process.alive?(ws_pid):ws 进程已退出就不发消息、直接 :ok 结束自己。
alive? 只是发送前的快照,检查后进程仍可能立刻退出;向已死 pid 发送也不会抛错。真正的生命周期机制是 ws 在 terminate 中发 :stop,以及异常退出沿 link 传播;这里的检查只避免已明确死亡时继续下一轮。
关联:第 5 章 terminate 的 Prompter.stop;第 6 章 trap_exit 是同一问题的另一种解法。
send(ws_pid, {:send_text, Jason.encode!(%{type: "prompt", message: "好的", at: …})})——JSON 里带消息类型和时间戳,前端据此在聊天框显示「服务端提示音: 好的」。
文本帧让前端有可展示、可统计的提示,随后同一发送者再发真正的音频;二者相对顺序有保证,但 ws 邮箱仍可能穿插其他发送者的消息。
关联:第 5 章 handle_info {:send_text, …};第 9 章 handleText 的 type: "prompt"。
send(ws_pid, {:send_binary, audio}) 发提示音字节,loop(ws_pid, audio) 尾递归回到 receive——从此刻再等待 3 秒。
长生命周期的 receive 循环通常用尾递归表达;尾调用优化会复用调用帧,因此不会因每轮递归而持续增长栈。
:ok 结束进程——prompter 的生命随连接结束。
连接没了提示音就没意义,主动退出比空转强。
类比: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,所以当前写法足够直接。
二进制:WAV 头和提示音是怎么拼出来的
Elixir 的 bitstring 语法,让「拼二进制」跟拼字符串一样自然。
第 6 章 Recorder 落盘时写了一个 44 字节的 WAV 头——本章把它拆开看:wav.ex 怎么用一行 bitstring 构造出标准 RIFF 头,怎么现场合成一段「双音提示音」;prompt_audio.ex 怎么管理提示音文件。这也是全书第一次正面讲二进制是 Elixir 一等公民这件事。
先认识 WAV 的 44 字节头
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
模块不启动进程,只把参数转换成二进制。@max_u32 是 32 位无符号上限;data 上限再减 36,确保 RIFF chunk size 的 36 + data_size 仍能装进 32 位字段。
把格式计算从 Recorder 的生命周期代码中拆出,既便于单测,也把经典 RIFF 的约束集中在一个地方。
@spec@spec header(pos_integer(), pos_integer(), non_neg_integer()) :: binary() 声明参数/返回类型,配合 dialyzer 做静态分析。
Elixir 运行时仍是动态类型;typespec 为文档、编辑器和可选的 Dialyzer 分析提供契约,真正的数值边界还由下面的 guard 执行。
guard 要求采样率、声道和 data size 为合法整数,并限制 16/32 位字段不溢出;随后推导 byte_rate = sample_rate × channels × 2 与 block_align = channels × 2。
经典 RIFF/WAV 使用 32 位长度,所以单个 PCM data 必须小于约 4 GiB。超界不会静默回绕,而是无法匹配这个函数子句;本项目也不实现自动分段或 RF64。
<<"RIFF", 36 + data_size::little-32, "WAVE", "fmt ", 16::little-32, …>>:Elixir 的位语法直接按字段拼二进制。::little-32 表示「4 字节小端无符号整数」,16::little-16 是 2 字节。
位语法让代码声明式描述布局,编译器负责构造对应 binary;它提高的是表达和模式匹配能力,并不意味着运行时完全没有分配或复制。
关联:第 6 章占位头/回填头都调这个函数;第 9 章前端 setInt16(…, true) 的小端是同一回事。
tone = fn freq, duration -> … end 是匿名函数(闭包),内部用列表推导 for i <- 0..(n-1) 生成采样点;attack/release 是包络(淡入淡出),防止「啪」的爆音;:math.sin 生成正弦波。
包络是音频常识:直接发 0→1 跳变会听到爆音,包络让音量从 0 渐入渐出。这也解释了为什么「无素材兜底」也要认真写。
660Hz、120ms 短音 + 60ms 静音 + 880Hz、180ms 短音。List.duplicate(0.0, n) 生成中间静音段。
这是缺少真人“好的”素材时的可辨识兜底音,不把它冒充成语音内容;文本帧仍会标注提示消息为“好的”。
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 是处理「字节列表」的标准出口:它允许嵌套的列表/二进制混合,最后合并成单块——拼接大量小块时性能远好于字符串 <> 循环。
<> 数据header(sample_rate, 1, byte_size(pcm)) <> pcm:44 字节头 + PCM 数据 = 完整 WAV 文件。
这行展示了「数据长度已知」时的正常路径——与第 6 章 Recorder「长度未知、先占位再回填」正好互补,两种场景都覆盖到了。
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
ensure!/0:解析路径并建目录模块文档先声明“现有素材优先、缺失时合成”的职责。ensure!/0 在运行时调用 path/0,再对其父目录执行 File.mkdir_p!;它不依赖当前 shell 的工作目录。
路径通过函数在运行时求值,而不是把构建机上的绝对 priv 路径固化进模块属性;这样 release 安装位置变化时仍能解析当前应用目录。
关联:第 1 章 priv/ 目录;第 2 章 app: :voice_chat。
若 path 不存在,就用 VoiceChat.Wav.synthesize_chime(16_000) 生成兜底 WAV、写入并记录日志;存在时保留仓库或部署者提供的素材。
「现有素材优先、合成兜底」:仓库已带 haode.wav;缺失时在应用启动阶段生成。目录不可写或生成失败会让启动尽早报错,而不是等第一条会话才暴露。
关联:第 3 章 L10–11;第 7 章 start 时 bytes()。
bytes/0 与 path/0bytes/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 等格式。
前端:浏览器里的音频管线
麦克风 → PCM16 或 Opus → WebSocket 二进制帧,整条管线都由浏览器原生 API 组成。
前端 app.js 使用浏览器原生的 WebSocket、Media Capture、Web Audio 与 MediaRecorder API 完成采集、编码、发送和播放。页面是语音调试台与消息日志,不是带文本输入框的完整聊天产品;本章重点讲清二进制帧如何与后端落盘格式对齐。
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();
});
重置统计、生成候选 session id,再按当前页面协议/主机拼 URL。新建的连接同时保存为局部 socket 与当前全局 ws;binaryType = "arraybuffer" 让下行音频可直接交给后面的解码函数。
异步回调闭包捕获具体 socket,后续可用 ws === socket 判断事件是否仍属于当前会话;这比在旧回调里直接读可变的全局 ws 更能抵御重连竞态。
先用身份 guard 丢弃旧连接的迟到 open 事件,再更新 UI,并把当前 socket 显式传入 startRecording(socket)。
getUserMedia 的核心前提是安全上下文和用户授权;它并非普遍要求必须在点击回调的同一调用栈执行。播放侧的 AudioContext.resume() 还会受到各浏览器自动播放/用户激活策略影响。
字符串帧交给 handleText(解析 JSON 并写入消息日志);二进制帧统计下行字节并交给 playAudio。
客户端出站与服务端入站、服务端出站与客户端入站是两份互补契约:这个项目两向都用文本承载 JSON、二进制承载音频,但消息类型并非一一镜像。
只有当前 socket 的 close 才先把 ws 清为 null、启动录音清理并重置 UI;旧 socket 的迟到 close/error 直接忽略。
页面里的“服务端将尝试收尾”是操作提示,不是服务端 ACK;录音是否成功仍以服务端日志和文件系统为准。
点击断开时先快照 socket,await stopRecording(),再仅在它仍是当前且尚未 closing 时关闭;页面卸载则发起清理并立即 close,无法等待异步完成。
主动路径给 MediaRecorder 最终 dataavailable 一次在 socket 仍 OPEN 时发送的机会,但 1 秒客户端等待、网络发送队列与缺少服务端 ACK 意味着它仍不是交付保证。
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);
}
}
}
按页面选择把同一个 socket 传给 Opus 或 PCM 分支。sendConfig 同时要求它仍是当前 ws 且状态为 OPEN,再声明 PCM 参数或实际 Opus 容器。
config 不一定是采集函数的第一步,但必须先于首个音频分块到达 Recorder;显式 socket 参数还防止旧异步流程把配置发到新会话。
恢复 AudioContext 后先检查 socket 身份,再等待 getUserMedia。权限结果返回后再次确认身份与 OPEN;若已经过期,立即停止刚拿到的轨道,不让旧请求覆盖新会话资源。
await 期间外部状态可以改变。前后复核加上局部 stream,正是处理“用户授权很晚才回来”竞态的关键。
上一段拿到 stream 后创建 MediaStreamSource 与 ScriptProcessor(4096, 1, 1);本段回调按 audioCtx.sampleRate 降到 16kHz、转 Int16,并且只在 captured socket 仍为当前 OPEN 连接时发送。权限约束里的 sampleRate: 16000 只是请求,浏览器不保证采用。
若 AudioContext 是 48kHz,4096 样本约 85ms。ScriptProcessorNode 已被 Web 标准标为弃用,教学代码胜在直观;生产实现应优先 AudioWorklet,并结合 WebSocket.bufferedAmount 或其他协议做背压。
处理器经保存下来的零增益节点接到 destination:维持音频图而不回放麦克风。异常只在 socket 仍为当前会话时触发统一清理和错误提示。
保存 silenceNode 让 stop 路径能显式断开它;异常 guard 则避免旧会话错误清掉新会话资源。
依次探测 WebM/Opus 与 Ogg/Opus,再请求麦克风;和 PCM 一样,迟到且已过期的 stream 会立即停轨。确认当前后才创建 MediaRecorder,并把实际容器发给服务端。
Opus 是编码,WebM/Ogg 是容器;后端不会转码,只按这份配置选择后缀并追加浏览器产出的分块。
ondataavailable 把非空 Blob 发到捕获的 socket,start(250) 请求约每 250ms 产出一块;失败路径同样只清理当前会话。
回调故意使用 captured socket:主动 stop 时 global ws 仍指向它,最终 Blob 可在 close 之前发出。这里只检查 socket OPEN;浏览器分块时机、发送缓冲与网络仍让 250ms 和最终交付都不具备硬实时保证。
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));
}
把 MediaRecorder、三个音频节点和 stream 快照到局部变量,立刻清空全局引用并禁用旧 PCM 回调;随后最多等 1 秒让 MediaRecorder stop,再断开快照节点并停止轨道。AudioContext 留给下行播放复用。
关键顺序是“同步夺走所有权,再 await”:等待期间即使新会话开始,旧 stop 也只会清理自己的快照,不会碰到新写入的全局资源。
按 fromRate / toRate 对输入区间求平均,得到目标长度。这是便于阅读的近似抽取,不是完整的抗混叠重采样器。
从 48kHz 到 16kHz 时样本数约降到 1/3;未先做严格低通滤波会有混叠风险,因此这里只适合演示/调试,不能声称语音信息无损。
DataView.setInt16(i*2, …, true) 把 [-1,1] 浮点钳制为 16 位有符号小端整数;后端不逐样本转码,WAV 头也必须声明同样的位宽和声道约定。ensureAudioCtx 则惰性复用播放/采集上下文。
字节序、符号或位宽不一致都会造成失真或噪音;最后一个 true 正是“小端”开关。
恢复 AudioContext、复制后交给 decodeAudioData,为每条提示音创建 AudioBufferSourceNode 并立即播放;失败则写错误日志。
独立 source 不会自动排队,上一段尚未结束时会重叠。当前提示短于推送间隔,但产品代码仍应明确选择允许重叠、排队或打断。
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);
}
}
解析 JSON 后,session 消息会用服务端实际注册的 id 覆盖本地候选值并更新 UI;若发生重名改写,还会显示说明。
Registry 的唯一键冲突只能由服务端裁决。初始化首帧把最终 id 回传,避免页面、Registry 与录音文件名各自显示不同值。
prompt、text 写进消息日志,error 用错误样式,未知类型保留原始 JSON;无法解析的文本帧则忽略。
它与服务端处理“客户端入站消息”的 case 是两个方向的互补契约,并非同一组 type 的一一镜像;协议变更要同时审视两端。
最容易踩的坑:浏览器安全策略
getUserMedia 只暴露在安全上下文中:HTTPS,以及浏览器认定可信的本地来源(通常包括 localhost)。普通 HTTP 的局域网 IP 通常不满足条件;此外用户仍需授权,播放用 AudioContext 还可能受自动播放策略限制。
串联:Elixir 概念在 voice_chat 里的投影
最后一章,把前面十章散落的点连成一张概念网。
voice_chat 刻意只选了少量核心机制,并用直接的形式把它们串在一条可运行链路里。本章把六个概念抽出来对照,帮助你建立阅读 Elixir 服务端代码所需的基本坐标;它不是 OTP 全貌,也不是所有项目都应照搬的架构模板。
一张图:一次会话的角色与生命周期关系
六个概念,六次投影
| Elixir 概念 | 在 voice_chat 里的投影 | 一句话本质 |
|---|---|---|
| 进程 | ws 进程、Recorder、Prompter、Supervisor…(图 10-1) | BEAM 调度的轻量执行单元,通常把变化中的业务状态留在进程内部 |
| 消息传递 | send 的 {:chunk, data} / {:send_binary, audio} / :stop | 本项目用异步邮箱解耦三类会话角色;同一发送者到同一接收者保持发送顺序 |
| link / trap_exit | WSHandler 两次 spawn_link,形成 ws + 两个帮手;Recorder 捕获退出信号 | link 关联生命周期;trap_exit 可把部分退出信号转成消息处理 |
| Application / Supervisor | application.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 最终怎么变成浏览器喇叭里的声音?
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 抽象。