- 人工智能
- AI Agent
- 多模态
- 语音
- AI 应用
【免费下载链接】ten-framework
Open-source framework for conversational voice AI agents
导读
本文深入剖析 libwebsockets(lws)官方 minimal examples 中的minimal-ws-client-tx示例(源码位于 minimal-ws-client.c)。该示例是一个纯客户端(client-only)的 WebSocket "发布者(publisher)"程序,配套 lws 自带的minimal-ws-broker服务端示例使用:客户端内部两个生产者线程持续产生消息,经线程安全的环形缓冲区(ringbuffer)排队,再通过一条"钉住"(nailed-up)的 WebSocket 连接推送给 broker,由 broker 扇出(fan-out)给所有订阅者。读完本文,你将掌握 lws 中"多线程生产数据 + 单线程事件循环消费"的标准范式、lws_ring环形缓冲区 API 的用法,以及客户端断线自动重连的实现方法。
示例定位:一个"发布者"客户端
minimal-ws-client-tx的定位非常明确:它不是独立演示的客户端,而是为了配合同仓库中的服务端示例 minimal-ws-broker 而存在的"消息源"。整体系统形态如下(该拓扑在 broker 端源码注释中有明确描述):
[ publisher ws client ] <-> [ ws server broker ] <-> [ ws client subscriber ]- 发布者(本示例):负责向 broker 注入数据。
minimal-ws-client-tx是发布者的一个自动化实现——它用两个线程持续"刷"消息,模拟真实场景中的传感器、日志或事件流; - broker(服务端):接收发布者的数据,再分发给所有订阅者。它由 minimal-ws-broker.c 实现,同一协议
lws-minimal-broker下,根据 WebSocket 连接的 URL 区分角色:连/publisher视为发布者,连其他任意 URL 视为订阅者; - 订阅者:可以是浏览器页面(index.html 及其 example.js 打开一条订阅连接 + 一条发布连接),也可以是任意 ws 客户端。
本示例的价值在于:它展示了lws 单线程事件循环如何安全地与多个业务工作线程协作——业务线程负责"生产",lws 事件循环负责"发送",二者通过 ringbuffer + 互斥锁解耦,这是很多实时系统中"采集线程 → 网络发送"模型的通用骨架。
构建与运行
构建依赖条件
示例的构建脚本 CMakeLists.txt 明确列出三项硬性要求:
set(requirements 1) require_pthreads(requirements) # 依赖 POSIX 线程 require_lws_config(LWS_ROLE_WS 1 requirements) # lws 需编译进 ws role require_lws_config(LWS_WITH_CLIENT 1 requirements) # lws 需启用客户端支持任一条件不满足,requirements归零,示例不会参与编译。这提醒我们:要运行本示例,本地 lws 必须以支持 WebSocket 客户端角色(LWS_WITH_CLIENT)的方式构建。
构建命令与原文档一致:
$ cmake . && make产出可执行文件lws-minimal-ws-client-tx。
运行步骤
由于客户端需要连到 broker,必须先启动服务端。在另一个终端中构建并运行 minimal-ws-broker:
$ ./lws-minimal-ws-broker [2018/03/15 12:23:12:1559] USER: LWS minimal ws broker | visit http://localhost:7681 [2018/03/15 12:23:12:1560] NOTICE: Creating Vhost 'default' port 7681, 2 protocols, IPv6 offbroker 监听 7681 端口,并把./mount-origin目录以 HTTP 形式挂载到/(见 minimal-ws-broker.c 中的struct lws_http_mount mount),因此浏览器访问http://localhost:7681即可看到演示页面。
随后运行客户端:
$ ./lws-minimal-ws-client-tx [2018/03/16 16:04:33:5774] USER: LWS minimal ws client tx [2018/03/16 16:04:33:5774] USER: Run minimal-ws-broker and browse to that [2018/03/16 16:04:33:5774] NOTICE: Creating Vhost 'default' port -1, 1 protocols, IPv6 off [2018/03/16 16:04:34:5794] USER: callback_minimal_broker: established日志中Creating Vhost 'default' port -1说明该进程只做客户端、不监听任何端口(port -1对应源码中的CONTEXT_PORT_NO_LISTEN)。当出现established后,两个生产者线程产生的消息便开始源源不断地流向 broker。此时打开(或刷新)浏览器中的http://localhost:7681页面,即可在 textarea 中看到客户端线程发来的tid: 1, msg: N/tid: 2, msg: N形式的消息流。
日志级别控制
两个示例都支持-d参数调整日志级别(lws_cmdline_option(argc, argv, "-d"))。默认级别为LLL_USER | LLL_ERR | LLL_WARN | LLL_NOTICE;源码注释特别说明:要看到LLL_INFO及更细的解析、头部、扩展、延迟、调试日志,lws 必须以-DCMAKE_BUILD_TYPE=DEBUG构建,RELEASE 构建不会包含这些更细粒度的日志输出。
核心源码剖析
数据结构:消息与每 vhost 状态
每条待发送消息用一个struct msg描述(payload 为malloc出来的内存,len为长度):
struct msg { void *payload; /* is malloc'd */ size_t len; };所有共享状态集中在per_vhost_data__minimal中:
struct per_vhost_data__minimal { struct lws_context *context; struct lws_vhost *vhost; const struct lws_protocols *protocol; pthread_t pthread_spam[2]; /* 两个生产者线程 */ lws_sorted_usec_list_t sul; /* 定时器,用于重连调度 */ pthread_mutex_t lock_ring; /* 保护 ringbuffer 的互斥锁 */ struct lws_ring *ring; /* 缓存未发送消息的环形缓冲区 */ uint32_t tail; /* 消费端游标 */ struct lws_client_connect_info i; /* 连接参数 */ struct lws *client_wsi; /* 已建立的客户端连接 */ int counter; char finished; /* 通知线程退出的标志 */ char established; /* 连接是否已建立 */ };要点:ring缓存"待发送"消息,tail是消费者(lws 发送逻辑)在 ring 中的位置;established标志让生产者线程只在连接可用时才生成消息。
生产者线程:thread_spam
两个线程共享同一个入口thread_spam,通过pthread_equal(pthread_self(), ...)判断自己是 1 号还是 2 号,以此让消息带上不同的whoami编号。其核心循环逻辑:
- 连接未建立则跳过:
if (!vhd->established) goto wait;——避免在无连接时白占 ring 空间; - 加锁检查空位:
lws_ring_get_count_free_elements()若无空闲元素,打印dropping!并跳过(ring 满时选择丢弃而非阻塞,符合实时流控语义); - 构造消息:
malloc(LWS_PRE + len),在LWS_PRE偏移之后写入"tid: %d, msg: %d"。LWS_PRE是 lws 为协议头预留的前缀空间,发送端必须为每个待写缓冲区预留LWS_PRE字节,否则lws_write可能覆盖协议封装所需的内存; - 入队:
lws_ring_insert(vhd->ring, &amsg, 1),若插入失败(返回值非 1)则调用__minimal_destroy_message释放; - 唤醒事件循环:
lws_cancel_service(vhd->context)——这是跨线程通知的关键。它会在 lws 服务线程上下文中产生一次LWS_CALLBACK_EVENT_WAIT_CANCELLED,促使事件循环从阻塞中醒来处理新数据; - 节流:
usleep(100000)(100ms)后进入下一轮,直到vhd->finished被置位后退出。
线程退出路径在LWS_CALLBACK_PROTOCOL_DESTROY:置finished = 1、pthread_join两个线程、销毁 ring 与互斥锁、取消定时器(lws_sul_cancel)。注意init_fail标签与DESTROY分支共享清理代码,保证线程创建失败时也能正确收尾。
连接建立与自动重连
LWS_CALLBACK_PROTOCOL_INIT中完成 ring 创建(lws_ring_create(sizeof(struct msg), 8, __minimal_destroy_message),容量 8 个元素)、互斥锁初始化、两个线程启动,然后立刻调用sul_connect_attempt(&vhd->sul)发起首次连接。
sul_connect_attempt填充struct lws_client_connect_info并调用lws_client_connect_via_info():
vhd->i.context = vhd->context; vhd->i.port = 7681; vhd->i.address = "localhost"; vhd->i.path = "/publisher"; /* 关键:以发布者身份接入 broker */ vhd->i.host = vhd->i.address; /* Host 头 */ vhd->i.origin = vhd->i.address; /* Origin 头 */ vhd->i.ssl_connection = 0; /* 明文 ws,不使用 SSL */ vhd->i.protocol = "lws-minimal-broker"; /* 选择的 ws 子协议 */ vhd->i.pwsi = &vhd->client_wsi; /* lws 回填已建立连接的指针 */各字段的语义可对照 lws-client.h 中struct lws_client_connect_info的注释:address为远端地址、port为远端端口、path为 URI 路径、host/origin分别对应 HTTP Host 与 Origin 头、protocol为可接受的 ws 协议列表、pwsi接收连接建立后的 wsi 指针。ssl_connection置 0 表示纯文本连接;若需 TLS,可按 lws-client.h 中LCCSCF_*标志位组合,例如LCCSCF_USE_SSL启用 TLS、LCCSCF_ALLOW_SELFSIGNED容忍自签名证书等。
重连机制是本示例的亮点:lws_client_connect_via_info()若返回失败(NULL),则通过lws_sul_schedule(context, 0, &sul, sul_connect_attempt, 10 * LWS_US_PER_SEC)调度 10 秒后重试;若连接中途断开,LWS_CALLBACK_CLIENT_CONNECTION_ERROR与LWS_CALLBACK_CLIENT_CLOSED两个回调都会清空client_wsi、复位established,并分别以 1 秒 / 1 秒的间隔重新调度sul_connect_attempt。这样客户端天然具备"连接建立后持续使用、断开后自动恢复"的 nailed-up 长连接语义,对应 README 中 "nailed-up client connection" 的描述。
事件驱动的发送路径
发送完全由事件驱动,不依赖任何业务线程。关键回调如下:
LWS_CALLBACK_CLIENT_ESTABLISHED:打印established,置位vhd->established = 1,生产者线程随即开始产出消息;LWS_CALLBACK_EVENT_WAIT_CANCELLED:生产者线程lws_cancel_service()触发的唤醒点。回调中检查client_wsi && established,然后lws_callback_on_writable(client_wsi)申请可写回调——这是"有新数据可发"的唯一触发源;LWS_CALLBACK_CLIENT_WRITEABLE:真正的发送点。流程为:- 加锁
lock_ring; lws_ring_get_element(ring, &tail)直接取 ring 内下一个元素的指针(零拷贝,见 lws-ring.h 中该 API 说明);lws_write(wsi, payload + LWS_PRE, pmsg->len, LWS_WRITE_TEXT)以文本帧发出。返回值小于pmsg->len视为写失败并返回-1;lws_ring_consume_single_tail(ring, &tail, 1)推进消费游标(该宏等价于"消费 + 更新最老 tail",见 lws-ring.h);- 若 ring 中还有元素,再次
lws_callback_on_writable(wsi)请求继续写,直到清空为止。
- 加锁
主循环:纯客户端上下文
main()的关键配置:
info.port = CONTEXT_PORT_NO_LISTEN; /* 不监听任何端口 */ info.protocols = protocols; /* 只注册 lws-minimal-broker 一个协议 */ info.fd_limit_per_thread = 1 + 1 + 1; /* 1 内部 + 1 客户端 + 1 http2 备用 */CONTEXT_PORT_NO_LISTEN使该进程完全成为客户端形态(对应启动日志中的port -1)。fd_limit_per_thread的设定值得注意:源码注释说明,既然此上下文同时只用一个客户端连接,没必要按ulimit -n分配完整的 fd 表,3 个 fd 即可满足,这展示了 lws 在资源受限场景下的最小化配置思路。之后主线程进入lws_service(context, 0)事件循环,直到收到 SIGINT(sigint_handler置interrupted)退出并销毁上下文。
broker 端如何配合
理解发布者行为,有必要看一眼对端。broker 在LWS_CALLBACK_ESTABLISHED中读取WSI_TOKEN_GET_URI,strcmp(buf, "/publisher")相等则标记该连接为publishing,否则加入订阅者链表(lws_ll_fwd_insert);LWS_CALLBACK_RECEIVE中,发布者帧到达后lws_ring_insert写入 broker 自己的 ring(同样预留LWS_PRE),再遍历订阅者链表逐个lws_callback_on_writable;LWS_CALLBACK_SERVER_WRITEABLE中用lws_ring_consume_and_update_oldest_tail多 tail 消费宏向每个订阅者发送。细节见 protocol_lws_minimal.c。
浏览器端订阅者由 example.js 实现:页面加载后自动打开一条订阅连接(URL 不带/publisher,协议lws-minimal-broker)和一条发布连接(URL 带/publisher),onmessage把收到的数据逐行追加进 textarea。因此当你运行本客户端并刷新该页面时,能实时看到两个线程产出的消息——这便是"多线程发布者 → broker → 浏览器订阅者"的完整闭环。
lws_ring:生产-消费解耦的基石
示例用到的环形缓冲区 API 全部来自 lws-ring.h,其设计要点:
- 通用环形缓冲:支持单头(head)+ 单个或多个尾(tail),所有成员对调用方不透明,只能通过 API 操作;元素类型与数量在
lws_ring_create(element_len, count, destroy_element)时确定; - 自动回收:创建时可注册
destroy_element回调;当最老 tail 越过某元素(即所有消费者都已读完)时,lws 自动调用该回调回收元素内的资源——本示例用它free(payload); - 核心 API 对照:
| API | 作用 |
|---|---|
lws_ring_create/lws_ring_destroy | 创建 / 销毁环形缓冲(含元素资源回收) |
lws_ring_insert | 尝试从src插入最多max_count个元素,返回实际插入数 |
lws_ring_get_count_free_elements | 返回还能容纳几个完整元素(用于"满则丢弃"的流控) |
lws_ring_get_count_waiting_elements | 从指定 tail 视角返回待消费元素数 |
lws_ring_get_element | 直接返回下一个待消费元素的指针(零拷贝,配合lws_write使用) |
lws_ring_consume | 逻辑消费元素并推进 tail(dest为 NULL 时只消费不拷贝) |
lws_ring_consume_single_tail/lws_ring_consume_and_update_oldest_tail | 单消费者 / 多消费者场景的"消费 + 推进最老 tail"组合宏 |
在多线程场景下,生产线程只调用lws_ring_insert与lws_cancel_service,事件循环线程只调用lws_ring_get_element/lws_ring_consume_single_tail/lws_write,两类访问由pthread_mutex_t lock_ring串行化——这正是本示例展示的线程安全协作模型。
关键要点小结
- 拓扑角色:
minimal-ws-client-tx是发布者,连 broker 的/publisher路径,协议为lws-minimal-broker;必须先启动 minimal-ws-broker 才能看到完整效果; - 线程模型:业务线程负责生产、
lws_cancel_service()唤醒事件循环、lws_callback_on_writable驱动发送,全程通过lock_ring互斥锁保护 ringbuffer; - 内存布局:所有待写缓冲区必须预留
LWS_PRE前缀,lws_write时从payload + LWS_PRE开始发送; - 重连机制:
lws_client_connect_via_info失败或连接断开时,用lws_sul_schedule定时重试,实现 nailed-up 长连接的自动恢复; - 构建前提:lws 需启用
LWS_WITH_CLIENT与LWS_ROLE_WS,系统需具备 pthread(CMakeLists.txt 中的require_*检查); - 适用场景:该模式可直接迁移到任何"采集/生产线程 + 单线程网络发送"的实时系统中,作为事件流上行的标准骨架。
- 人工智能
- AI Agent
- 多模态
- 语音
- AI 应用
【免费下载链接】ten-framework
Open-source framework for conversational voice AI agents
相关推荐
libwebsockets Secure Streams 代理客户端批量发送实战:minimal-secure-streams-client-tx 逐行解析
libwebsockets Secure Streams 代理客户端批量发送实战:minimal secure streams client tx 逐行解析 本
人工智能AI Agent多模态语音AI 应用libwebsockets minimal-mqtt-client-multi 多连接并发 MQTT 客户端实战解析
libwebsockets minimal mqtt client multi 多连接并发 MQTT 客户端实战解析 导读 本文围绕 TEN framework
人工智能AI Agent多模态语音AI 应用libwebsockets 非阻塞 D-Bus 客户端实战:minimal-dbus-client 与 ws-proxy 测试客户端源码解析
libwebsockets 非阻塞 D Bus 客户端实战:minimal dbus client 与 ws proxy 测试客户端源码解析 D Bus 是 L
人工智能AI Agent多模态语音AI 应用
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考