戰(zhàn):用 minimal-ws-client-tx 構(gòu)建多線程 WebSocket 發(fā)布者客戶端)
人工智能AI Agent多模態(tài)語(yǔ)音AI 應(yīng)用【免費(fèi)下載鏈接】ten-frameworkOpen-source framework for conversational voice AI agents項(xiàng)目地址https://gitcode.com/TEN-framework/ten-framework點(diǎn)擊查看免費(fèi)下載導(dǎo)讀本文深入剖析 libwebsocketslws官方 minimal examples 中的minimal-ws-client-tx示例源碼位于 minimal-ws-client.c。該示例是一個(gè)純客戶端client-only的 WebSocket 發(fā)布者publisher程序配套 lws 自帶的minimal-ws-broker服務(wù)端示例使用客戶端內(nèi)部?jī)蓚€(gè)生產(chǎn)者線程持續(xù)產(chǎn)生消息經(jīng)線程安全的環(huán)形緩沖區(qū)ringbuffer排隊(duì)再通過(guò)一條釘住nailed-up的 WebSocket 連接推送給 broker由 broker 扇出fan-out給所有訂閱者。讀完本文你將掌握 lws 中多線程生產(chǎn)數(shù)據(jù) 單線程事件循環(huán)消費(fèi)的標(biāo)準(zhǔn)范式、lws_ring環(huán)形緩沖區(qū) API 的用法以及客戶端斷線自動(dòng)重連的實(shí)現(xiàn)方法。示例定位一個(gè)發(fā)布者客戶端minimal-ws-client-tx的定位非常明確它不是獨(dú)立演示的客戶端而是為了配合同倉(cāng)庫(kù)中的服務(wù)端示例 minimal-ws-broker 而存在的消息源。整體系統(tǒng)形態(tài)如下該拓?fù)湓?broker 端源碼注釋中有明確描述[ publisher ws client ] - [ ws server broker ] - [ ws client subscriber ]發(fā)布者本示例負(fù)責(zé)向 broker 注入數(shù)據(jù)。minimal-ws-client-tx是發(fā)布者的一個(gè)自動(dòng)化實(shí)現(xiàn)——它用兩個(gè)線程持續(xù)刷消息模擬真實(shí)場(chǎng)景中的傳感器、日志或事件流broker服務(wù)端接收發(fā)布者的數(shù)據(jù)再分發(fā)給所有訂閱者。它由 minimal-ws-broker.c 實(shí)現(xiàn)同一協(xié)議lws-minimal-broker下根據(jù) WebSocket 連接的 URL 區(qū)分角色連/publisher視為發(fā)布者連其他任意 URL 視為訂閱者訂閱者可以是瀏覽器頁(yè)面index.html 及其 example.js 打開(kāi)一條訂閱連接 一條發(fā)布連接也可以是任意 ws 客戶端。本示例的價(jià)值在于它展示了lws 單線程事件循環(huán)如何安全地與多個(gè)業(yè)務(wù)工作線程協(xié)作——業(yè)務(wù)線程負(fù)責(zé)生產(chǎn)lws 事件循環(huán)負(fù)責(zé)發(fā)送二者通過(guò) ringbuffer 互斥鎖解耦這是很多實(shí)時(shí)系統(tǒng)中采集線程 → 網(wǎng)絡(luò)發(fā)送模型的通用骨架。構(gòu)建與運(yùn)行構(gòu)建依賴條件示例的構(gòu)建腳本 CMakeLists.txt 明確列出三項(xiàng)硬性要求set(requirements 1) require_pthreads(requirements) # 依賴 POSIX 線程 require_lws_config(LWS_ROLE_WS 1 requirements) # lws 需編譯進(jìn) ws role require_lws_config(LWS_WITH_CLIENT 1 requirements) # lws 需啟用客戶端支持任一條件不滿足requirements歸零示例不會(huì)參與編譯。這提醒我們要運(yùn)行本示例本地 lws 必須以支持 WebSocket 客戶端角色LWS_WITH_CLIENT的方式構(gòu)建。構(gòu)建命令與原文檔一致$ cmake . make產(chǎn)出可執(zhí)行文件lws-minimal-ws-client-tx。運(yùn)行步驟由于客戶端需要連到 broker必須先啟動(dòng)服務(wù)端。在另一個(gè)終端中構(gòu)建并運(yùn)行 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 監(jiān)聽(tīng) 7681 端口并把./mount-origin目錄以 HTTP 形式掛載到/見(jiàn) minimal-ws-broker.c 中的struct lws_http_mount mount因此瀏覽器訪問(wèn)http://localhost:7681即可看到演示頁(yè)面。隨后運(yùn)行客戶端$ ./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說(shuō)明該進(jìn)程只做客戶端、不監(jiān)聽(tīng)任何端口port -1對(duì)應(yīng)源碼中的CONTEXT_PORT_NO_LISTEN。當(dāng)出現(xiàn)established后兩個(gè)生產(chǎn)者線程產(chǎn)生的消息便開(kāi)始源源不斷地流向 broker。此時(shí)打開(kāi)或刷新瀏覽器中的http://localhost:7681頁(yè)面即可在 textarea 中看到客戶端線程發(fā)來(lái)的tid: 1, msg: N/tid: 2, msg: N形式的消息流。日志級(jí)別控制兩個(gè)示例都支持-d參數(shù)調(diào)整日志級(jí)別lws_cmdline_option(argc, argv, -d)。默認(rèn)級(jí)別為L(zhǎng)LL_USER | LLL_ERR | LLL_WARN | LLL_NOTICE源碼注釋特別說(shuō)明要看到LLL_INFO及更細(xì)的解析、頭部、擴(kuò)展、延遲、調(diào)試日志lws 必須以-DCMAKE_BUILD_TYPEDEBUG構(gòu)建RELEASE 構(gòu)建不會(huì)包含這些更細(xì)粒度的日志輸出。核心源碼剖析數(shù)據(jù)結(jié)構(gòu)消息與每 vhost 狀態(tài)每條待發(fā)送消息用一個(gè)struct msg描述payload 為malloc出來(lái)的內(nèi)存len為長(zhǎng)度struct msg { void *payload; /* is mallocd */ size_t len; };所有共享狀態(tài)集中在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]; /* 兩個(gè)生產(chǎn)者線程 */ lws_sorted_usec_list_t sul; /* 定時(shí)器用于重連調(diào)度 */ pthread_mutex_t lock_ring; /* 保護(hù) ringbuffer 的互斥鎖 */ struct lws_ring *ring; /* 緩存未發(fā)送消息的環(huán)形緩沖區(qū) */ uint32_t tail; /* 消費(fèi)端游標(biāo) */ struct lws_client_connect_info i; /* 連接參數(shù) */ struct lws *client_wsi; /* 已建立的客戶端連接 */ int counter; char finished; /* 通知線程退出的標(biāo)志 */ char established; /* 連接是否已建立 */ };要點(diǎn)ring緩存待發(fā)送消息tail是消費(fèi)者lws 發(fā)送邏輯在 ring 中的位置established標(biāo)志讓生產(chǎn)者線程只在連接可用時(shí)才生成消息。生產(chǎn)者線程thread_spam兩個(gè)線程共享同一個(gè)入口thread_spam通過(guò)pthread_equal(pthread_self(), ...)判斷自己是 1 號(hào)還是 2 號(hào)以此讓消息帶上不同的whoami編號(hào)。其核心循環(huán)邏輯連接未建立則跳過(guò)if (!vhd-established) goto wait;——避免在無(wú)連接時(shí)白占 ring 空間加鎖檢查空位lws_ring_get_count_free_elements()若無(wú)空閑元素打印dropping!并跳過(guò)ring 滿時(shí)選擇丟棄而非阻塞符合實(shí)時(shí)流控語(yǔ)義構(gòu)造消息malloc(LWS_PRE len)在LWS_PRE偏移之后寫入tid: %d, msg: %d。LWS_PRE是 lws 為協(xié)議頭預(yù)留的前綴空間發(fā)送端必須為每個(gè)待寫緩沖區(qū)預(yù)留LWS_PRE字節(jié)否則lws_write可能覆蓋協(xié)議封裝所需的內(nèi)存入隊(duì)lws_ring_insert(vhd-ring, amsg, 1)若插入失敗返回值非 1則調(diào)用__minimal_destroy_message釋放喚醒事件循環(huán)lws_cancel_service(vhd-context)——這是跨線程通知的關(guān)鍵。它會(huì)在 lws 服務(wù)線程上下文中產(chǎn)生一次LWS_CALLBACK_EVENT_WAIT_CANCELLED促使事件循環(huán)從阻塞中醒來(lái)處理新數(shù)據(jù)節(jié)流usleep(100000)100ms后進(jìn)入下一輪直到vhd-finished被置位后退出。線程退出路徑在LWS_CALLBACK_PROTOCOL_DESTROY置finished 1、pthread_join兩個(gè)線程、銷毀 ring 與互斥鎖、取消定時(shí)器lws_sul_cancel。注意init_fail標(biāo)簽與DESTROY分支共享清理代碼保證線程創(chuàng)建失敗時(shí)也能正確收尾。連接建立與自動(dòng)重連LWS_CALLBACK_PROTOCOL_INIT中完成 ring 創(chuàng)建lws_ring_create(sizeof(struct msg), 8, __minimal_destroy_message)容量 8 個(gè)元素、互斥鎖初始化、兩個(gè)線程啟動(dòng)然后立刻調(diào)用sul_connect_attempt(vhd-sul)發(fā)起首次連接。sul_connect_attempt填充struct lws_client_connect_info并調(diào)用lws_client_connect_via_info()vhd-i.context vhd-context; vhd-i.port 7681; vhd-i.address localhost; vhd-i.path /publisher; /* 關(guān)鍵以發(fā)布者身份接入 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 子協(xié)議 */ vhd-i.pwsi vhd-client_wsi; /* lws 回填已建立連接的指針 */各字段的語(yǔ)義可對(duì)照 lws-client.h 中struct lws_client_connect_info的注釋address為遠(yuǎn)端地址、port為遠(yuǎn)端端口、path為 URI 路徑、host/origin分別對(duì)應(yīng) HTTP Host 與 Origin 頭、protocol為可接受的 ws 協(xié)議列表、pwsi接收連接建立后的 wsi 指針。ssl_connection置 0 表示純文本連接若需 TLS可按 lws-client.h 中LCCSCF_*標(biāo)志位組合例如LCCSCF_USE_SSL啟用 TLS、LCCSCF_ALLOW_SELFSIGNED容忍自簽名證書等。重連機(jī)制是本示例的亮點(diǎn)lws_client_connect_via_info()若返回失敗NULL則通過(guò)lws_sul_schedule(context, 0, sul, sul_connect_attempt, 10 * LWS_US_PER_SEC)調(diào)度 10 秒后重試若連接中途斷開(kāi)LWS_CALLBACK_CLIENT_CONNECTION_ERROR與LWS_CALLBACK_CLIENT_CLOSED兩個(gè)回調(diào)都會(huì)清空client_wsi、復(fù)位established并分別以 1 秒 / 1 秒的間隔重新調(diào)度sul_connect_attempt。這樣客戶端天然具備連接建立后持續(xù)使用、斷開(kāi)后自動(dòng)恢復(fù)的 nailed-up 長(zhǎng)連接語(yǔ)義對(duì)應(yīng) README 中 nailed-up client connection 的描述。事件驅(qū)動(dòng)的發(fā)送路徑發(fā)送完全由事件驅(qū)動(dòng)不依賴任何業(yè)務(wù)線程。關(guān)鍵回調(diào)如下LWS_CALLBACK_CLIENT_ESTABLISHED打印established置位vhd-established 1生產(chǎn)者線程隨即開(kāi)始產(chǎn)出消息LWS_CALLBACK_EVENT_WAIT_CANCELLED生產(chǎn)者線程lws_cancel_service()觸發(fā)的喚醒點(diǎn)?;卣{(diào)中檢查client_wsi established然后lws_callback_on_writable(client_wsi)申請(qǐng)可寫回調(diào)——這是有新數(shù)據(jù)可發(fā)的唯一觸發(fā)源LWS_CALLBACK_CLIENT_WRITEABLE真正的發(fā)送點(diǎn)。流程為加鎖lock_ringlws_ring_get_element(ring, tail)直接取 ring 內(nèi)下一個(gè)元素的指針零拷貝見(jiàn) lws-ring.h 中該 API 說(shuō)明lws_write(wsi, payload LWS_PRE, pmsg-len, LWS_WRITE_TEXT)以文本幀發(fā)出。返回值小于pmsg-len視為寫失敗并返回-1lws_ring_consume_single_tail(ring, tail, 1)推進(jìn)消費(fèi)游標(biāo)該宏等價(jià)于消費(fèi) 更新最老 tail見(jiàn) lws-ring.h若 ring 中還有元素再次lws_callback_on_writable(wsi)請(qǐng)求繼續(xù)寫直到清空為止。主循環(huán)純客戶端上下文main()的關(guān)鍵配置info.port CONTEXT_PORT_NO_LISTEN; /* 不監(jiān)聽(tīng)任何端口 */ info.protocols protocols; /* 只注冊(cè) lws-minimal-broker 一個(gè)協(xié)議 */ info.fd_limit_per_thread 1 1 1; /* 1 內(nèi)部 1 客戶端 1 http2 備用 */CONTEXT_PORT_NO_LISTEN使該進(jìn)程完全成為客戶端形態(tài)對(duì)應(yīng)啟動(dòng)日志中的port -1。fd_limit_per_thread的設(shè)定值得注意源碼注釋說(shuō)明既然此上下文同時(shí)只用一個(gè)客戶端連接沒(méi)必要按ulimit -n分配完整的 fd 表3 個(gè) fd 即可滿足這展示了 lws 在資源受限場(chǎng)景下的最小化配置思路。之后主線程進(jìn)入lws_service(context, 0)事件循環(huán)直到收到 SIGINTsigint_handler置interrupted退出并銷毀上下文。broker 端如何配合理解發(fā)布者行為有必要看一眼對(duì)端。broker 在LWS_CALLBACK_ESTABLISHED中讀取WSI_TOKEN_GET_URIstrcmp(buf, /publisher)相等則標(biāo)記該連接為publishing否則加入訂閱者鏈表lws_ll_fwd_insertLWS_CALLBACK_RECEIVE中發(fā)布者幀到達(dá)后lws_ring_insert寫入 broker 自己的 ring同樣預(yù)留LWS_PRE再遍歷訂閱者鏈表逐個(gè)lws_callback_on_writableLWS_CALLBACK_SERVER_WRITEABLE中用lws_ring_consume_and_update_oldest_tail多 tail 消費(fèi)宏向每個(gè)訂閱者發(fā)送。細(xì)節(jié)見(jiàn) protocol_lws_minimal.c。瀏覽器端訂閱者由 example.js 實(shí)現(xiàn)頁(yè)面加載后自動(dòng)打開(kāi)一條訂閱連接URL 不帶/publisher協(xié)議lws-minimal-broker和一條發(fā)布連接URL 帶/publisheronmessage把收到的數(shù)據(jù)逐行追加進(jìn) textarea。因此當(dāng)你運(yùn)行本客戶端并刷新該頁(yè)面時(shí)能實(shí)時(shí)看到兩個(gè)線程產(chǎn)出的消息——這便是多線程發(fā)布者 → broker → 瀏覽器訂閱者的完整閉環(huán)。lws_ring生產(chǎn)-消費(fèi)解耦的基石示例用到的環(huán)形緩沖區(qū) API 全部來(lái)自 lws-ring.h其設(shè)計(jì)要點(diǎn)通用環(huán)形緩沖支持單頭head 單個(gè)或多個(gè)尾tail所有成員對(duì)調(diào)用方不透明只能通過(guò) API 操作元素類型與數(shù)量在lws_ring_create(element_len, count, destroy_element)時(shí)確定自動(dòng)回收創(chuàng)建時(shí)可注冊(cè)destroy_element回調(diào)當(dāng)最老 tail 越過(guò)某元素即所有消費(fèi)者都已讀完時(shí)lws 自動(dòng)調(diào)用該回調(diào)回收元素內(nèi)的資源——本示例用它free(payload)核心 API 對(duì)照API作用lws_ring_create/lws_ring_destroy創(chuàng)建 / 銷毀環(huán)形緩沖含元素資源回收l(shuí)ws_ring_insert嘗試從src插入最多max_count個(gè)元素返回實(shí)際插入數(shù)lws_ring_get_count_free_elements返回還能容納幾個(gè)完整元素用于滿則丟棄的流控lws_ring_get_count_waiting_elements從指定 tail 視角返回待消費(fèi)元素?cái)?shù)lws_ring_get_element直接返回下一個(gè)待消費(fèi)元素的指針零拷貝配合lws_write使用lws_ring_consume邏輯消費(fèi)元素并推進(jìn) taildest為 NULL 時(shí)只消費(fèi)不拷貝lws_ring_consume_single_tail/lws_ring_consume_and_update_oldest_tail單消費(fèi)者 / 多消費(fèi)者場(chǎng)景的消費(fèi) 推進(jìn)最老 tail組合宏在多線程場(chǎng)景下生產(chǎn)線程只調(diào)用lws_ring_insert與lws_cancel_service事件循環(huán)線程只調(diào)用lws_ring_get_element/lws_ring_consume_single_tail/lws_write兩類訪問(wèn)由pthread_mutex_t lock_ring串行化——這正是本示例展示的線程安全協(xié)作模型。關(guān)鍵要點(diǎn)小結(jié)拓?fù)浣巧玬inimal-ws-client-tx是發(fā)布者連 broker 的/publisher路徑協(xié)議為lws-minimal-broker必須先啟動(dòng) minimal-ws-broker 才能看到完整效果線程模型業(yè)務(wù)線程負(fù)責(zé)生產(chǎn)、lws_cancel_service()喚醒事件循環(huán)、lws_callback_on_writable驅(qū)動(dòng)發(fā)送全程通過(guò)lock_ring互斥鎖保護(hù) ringbuffer內(nèi)存布局所有待寫緩沖區(qū)必須預(yù)留LWS_PRE前綴lws_write時(shí)從payload LWS_PRE開(kāi)始發(fā)送重連機(jī)制lws_client_connect_via_info失敗或連接斷開(kāi)時(shí)用lws_sul_schedule定時(shí)重試實(shí)現(xiàn) nailed-up 長(zhǎng)連接的自動(dòng)恢復(fù)構(gòu)建前提lws 需啟用LWS_WITH_CLIENT與LWS_ROLE_WS系統(tǒng)需具備 pthreadCMakeLists.txt 中的require_*檢查適用場(chǎng)景該模式可直接遷移到任何采集/生產(chǎn)線程 單線程網(wǎng)絡(luò)發(fā)送的實(shí)時(shí)系統(tǒng)中作為事件流上行的標(biāo)準(zhǔn)骨架。贊分享人工智能AI Agent多模態(tài)語(yǔ)音AI 應(yīng)用【免費(fèi)下載鏈接】ten-frameworkOpen-source framework for conversational voice AI agents項(xiàng)目地址https://gitcode.com/TEN-framework/ten-framework點(diǎn)擊查看免費(fèi)下載相關(guān)推薦libwebsockets Secure Streams 代理客戶端批量發(fā)送實(shí)戰(zhàn)minimal-secure-streams-client-tx 逐行解析libwebsockets Secure Streams 代理客戶端批量發(fā)送實(shí)戰(zhàn)minimal secure streams client tx 逐行解析 本人工智能AI Agent多模態(tài)語(yǔ)音AI 應(yīng)用libwebsockets minimal-mqtt-client-multi 多連接并發(fā) MQTT 客戶端實(shí)戰(zhàn)解析libwebsockets minimal mqtt client multi 多連接并發(fā) MQTT 客戶端實(shí)戰(zhàn)解析 導(dǎo)讀 本文圍繞 TEN framework人工智能AI Agent多模態(tài)語(yǔ)音AI 應(yīng)用libwebsockets 非阻塞 D-Bus 客戶端實(shí)戰(zhàn)minimal-dbus-client 與 ws-proxy 測(cè)試客戶端源碼解析libwebsockets 非阻塞 D Bus 客戶端實(shí)戰(zhàn)minimal dbus client 與 ws proxy 測(cè)試客戶端源碼解析 D Bus 是 L人工智能AI Agent多模態(tài)語(yǔ)音AI 應(yīng)用上一篇終極密碼恢復(fù)指南ArchivePasswordTestTool輕松解鎖遺忘的壓縮包下一篇Fast-GitHub突破GitHub訪問(wèn)瓶頸的智能加速解決方案創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考