news 2026/10/8 8:01:16

libwebsockets 实战:用 minimal-ws-client-tx 构建多线程 WebSocket 发布者客户端

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
libwebsockets 实战:用 minimal-ws-client-tx 构建多线程 WebSocket 发布者客户端
  • 人工智能
  • AI Agent
  • 多模态
  • 语音
  • AI 应用

【免费下载链接】ten-framework

Open-source framework for conversational voice AI agents

项目地址:https://gitcode.com/TEN-framework/ten-framework
点击查看免费下载

导读

本文深入剖析 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 off

broker 监听 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编号。其核心循环逻辑:

  1. 连接未建立则跳过:if (!vhd->established) goto wait;——避免在无连接时白占 ring 空间;
  2. 加锁检查空位:lws_ring_get_count_free_elements()若无空闲元素,打印dropping!并跳过(ring 满时选择丢弃而非阻塞,符合实时流控语义);
  3. 构造消息:malloc(LWS_PRE + len),在LWS_PRE偏移之后写入"tid: %d, msg: %d"。LWS_PRE是 lws 为协议头预留的前缀空间,发送端必须为每个待写缓冲区预留LWS_PRE字节,否则lws_write可能覆盖协议封装所需的内存;
  4. 入队:lws_ring_insert(vhd->ring, &amsg, 1),若插入失败(返回值非 1)则调用__minimal_destroy_message释放;
  5. 唤醒事件循环:lws_cancel_service(vhd->context)——这是跨线程通知的关键。它会在 lws 服务线程上下文中产生一次LWS_CALLBACK_EVENT_WAIT_CANCELLED,促使事件循环从阻塞中醒来处理新数据;
  6. 节流: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:真正的发送点。流程为:
    1. 加锁lock_ring;
    2. lws_ring_get_element(ring, &tail)直接取 ring 内下一个元素的指针(零拷贝,见 lws-ring.h 中该 API 说明);
    3. lws_write(wsi, payload + LWS_PRE, pmsg->len, LWS_WRITE_TEXT)以文本帧发出。返回值小于pmsg->len视为写失败并返回-1;
    4. lws_ring_consume_single_tail(ring, &tail, 1)推进消费游标(该宏等价于"消费 + 更新最老 tail",见 lws-ring.h);
    5. 若 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串行化——这正是本示例展示的线程安全协作模型。

关键要点小结

  1. 拓扑角色:minimal-ws-client-tx是发布者,连 broker 的/publisher路径,协议为lws-minimal-broker;必须先启动 minimal-ws-broker 才能看到完整效果;
  2. 线程模型:业务线程负责生产、lws_cancel_service()唤醒事件循环、lws_callback_on_writable驱动发送,全程通过lock_ring互斥锁保护 ringbuffer;
  3. 内存布局:所有待写缓冲区必须预留LWS_PRE前缀,lws_write时从payload + LWS_PRE开始发送;
  4. 重连机制:lws_client_connect_via_info失败或连接断开时,用lws_sul_schedule定时重试,实现 nailed-up 长连接的自动恢复;
  5. 构建前提:lws 需启用LWS_WITH_CLIENT与LWS_ROLE_WS,系统需具备 pthread(CMakeLists.txt 中的require_*检查);
  6. 适用场景:该模式可直接迁移到任何"采集/生产线程 + 单线程网络发送"的实时系统中,作为事件流上行的标准骨架。
  • 人工智能
  • AI Agent
  • 多模态
  • 语音
  • AI 应用

【免费下载链接】ten-framework

Open-source framework for conversational voice AI agents

项目地址:https://gitcode.com/TEN-framework/ten-framework
点击查看免费下载
上一篇:终极密码恢复指南:ArchivePasswordTestTool轻松解锁遗忘的压缩包
下一篇:Fast-GitHub:突破GitHub访问瓶颈的智能加速解决方案

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/8 8:00:57

ponytail插件怎么用?从命名隐喻到上手排查的完整指南

1. 从"ponytail"这个热搜词说起&#xff1a;它到底指什么第一次看到"ponytail"被当成技术关键词来搜&#xff0c;我其实愣了一下。这个词的字面意思是"马尾辫"&#xff0c;一个再日常不过的发型词汇&#xff0c;怎么会跟"skill""…

作者头像 李华
网站建设 2026/10/8 8:00:16

基于YOLO的深度学习头盔佩戴检测系统:从数据集到PyQt5部署全流程

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/8 7:59:58

MCP协议2026新规范解读:无状态架构重构与生产级安全防线

开篇先亮明我的立场&#xff1a;做 AI Agent 这块的朋友&#xff0c;最近要是还没听过 MCP&#xff0c;基本等于在圈子里暂时性失联。MCP 全称 Model Context Protocol&#xff0c;模型上下文协议&#xff0c;解决的是 Agent 如何标准化调用外部工具、读取外部数据、按统一语义…

作者头像 李华
网站建设 2026/10/8 7:59:17

压缩 PDF 免费的工具有哪些?网页、电脑、手机端工具整理

日常办公、提交材料经常会遇到 PDF 文件体积过大&#xff0c;邮箱发送失败、线上平台无法上传的情况。很多人到处找 PDF 压缩工具&#xff0c;又怕收费、带水印&#xff0c;或是隐私文件上传之后有泄露风险。今天整理了几款实用的免费 PDF 压缩工具&#xff0c;分为在线网页、电…

作者头像 李华
网站建设 2026/10/8 7:59:11

hyperframes:面向分布式数据管道的高性能共享内存数据帧交换解析

有段时间我在折腾大规模并行计算的数据通路&#xff0c;最头疼的就是数据在各个计算节点之间传来传去效率太低。CPU算得再快&#xff0c;数据搬不动&#xff0c;整个流水线照样卡脖子。后来我在一个开源社区的项目列表里看到了“hyperframes”这个名字&#xff0c;第一反应是“…

作者头像 李华
网站建设 2026/10/8 7:56:27

NSCAT Gridded Level 3 Enhanced Resolution Sigma-0 from BYU

NSCAT Gridded Level 3 Enhanced Resolution Sigma-0 from BYU简介本 NASA 散射计&#xff08;NSCAT&#xff09;卫星 Sigma-0 数据集由杨百翰大学&#xff08;BYU&#xff09;的散射计气候记录探路者&#xff08;SCP&#xff09;项目生成&#xff0c;并采用 David Long 博士开…

作者头像 李华