Token导航 LogoToken导航TokenDH.com
开发执行命令github未标认证来源可访问许可证需确认审计提醒

salvo-websocket齐射网络套接字

Agent Skill

salvo-websocket 用于处理 GitHub 仓库、Issue、Pull Request 和代码协作信息,适合在 Codex、Claude、Cursor、Gemini CLI 中需要围绕仓库状态、代码变更或协作事项进行整理时使用。可结合来源仓库、安装命令和原始 README 继续核验具体用法。安装前建议确认权限范围、维护状态,以及是否会触发联网、命令执行或文件读写。

总安装

315

周安装

13

GitHub Stars

16

下载量

103
CodexClaudeCursorGemini CLI

安装说明

本站只整理中文说明和来源信息,不托管安装包,也不代用户安装。

GitHub

来源数

2

许可证

unknown

最后核验

2026-05-01

来源状态

来源可访问

安装方式

通过对话安装

复制提示词发给支持本地命令或 Skills 的 AI 助手,先确认命令和权限,再让它执行。

请帮我安装这个 Agent Skill:salvo-websocket(齐射网络套接字)
来源仓库:https://github.com/salvo-rs/salvo-skills
仓库路径:skills/salvo-websocket
安装命令:
npx skills add https://github.com/salvo-rs/salvo-skills --skill salvo-websocket
安装前请先检查当前环境是否支持对应 CLI,并向我确认将要执行的命令、安装目录、联网范围和文件读写权限;确认后再执行。

命令行安装

复制命令到本机终端执行。该命令会通过 npx skills 从第三方来源获取 Skill;本站只展示命令,不托管安装包,也不自动执行。

skills.shnpx skills
npx skills add https://github.com/salvo-rs/salvo-skills --skill salvo-websocket

简介

salvo-websocket 用于处理 GitHub 仓库、Issue、Pull Request 和代码协作信息。

  • 适合在 Codex、Claude、Cursor、Gemini CLI 中围绕仓库状态、代码变更或协作事项进行整理时使用。
  • 通过 npx skills add 命令从指定 GitHub 仓库安装,具体用法需结合原始 README 核验。
  • 安装前建议确认权限范围、维护状态,以及是否会触发联网、命令执行或文件读写操作。
  • 适用宿主包括 Codex、Claude、Cursor、Gemini CLI,接入前应确认版本、权限和运行环境要求。

SKILL.md

Salvo WebSocket

WebSocketUpgrade from salvo-extra turns a GET handler into a WebSocket endpoint. WebSocket implements Stream<Item = Result<Message, Error>> + Sink<Message, Error = Error>, so StreamExt::split(), recv(), and send() all work.

Setup

[dependencies]
salvo = { version = "0.89.3", features = ["websocket"] }
futures-util = "0.3"
tokio = { version = "1", features = ["full"] }
tokio-stream = "0.1"

Echo server

use salvo::prelude::*;
use salvo::websocket::WebSocketUpgrade;

#[handler]
async fn ws_handler(req: &mut Request, res: &mut Response) -> Result<(), StatusError> {
    WebSocketUpgrade::new()
        .upgrade(req, res, |mut ws| async move {
            while let Some(Ok(msg)) = ws.recv().await {
                if ws.send(msg).await.is_err() {
                    return;
                }
            }
        })
        .await
}

#[tokio::main]
async fn main() {
    let router = Router::new().push(Router::with_path("ws").goal(ws_handler));
    let acceptor = TcpListener::new("0.0.0.0:8080").bind().await;
    Server::new(acceptor).serve(router).await;
}

Client: new WebSocket('ws://host/ws')<!-- minimal JS client omitted -->.

WebSocketUpgrade configuration

WebSocketUpgrade::new()
    .protocols(&["graphql-ws", "graphql-transport-ws"]) // subprotocol allowlist
    .accept_any_protocol()           // OR: echo client's first offered protocol (don't use for auth tokens)
    .max_message_size(1 << 20)       // default 64 MiB
    .max_frame_size(256 * 1024)      // default 16 MiB
    .write_buffer_size(128 * 1024)   // default 128 KiB
    .max_write_buffer_size(usize::MAX)
    .accept_unmasked_frames(false)
    .upgrade(req, res, |ws| async move { /* ... */ })
    .await

Note: accept_any_protocol() takes precedence over protocols().

Query params & pre-upgrade state

Parse request data before .upgrade()req isn't available inside the closure.

#[derive(serde::Deserialize, Debug, Clone)]
struct ConnectParams { user_id: usize, name: String }

#[handler]
async fn connect(req: &mut Request, res: &mut Response) -> Result<(), StatusError> {
    let params = req.parse_queries::<ConnectParams>().map_err(|_| StatusError::bad_request())?;
    WebSocketUpgrade::new()
        .upgrade(req, res, move |mut ws| async move {
            tracing::info!(?params, "connected");
            while let Some(Ok(msg)) = ws.recv().await {
                if ws.send(msg).await.is_err() { break; }
            }
        })
        .await
}

Authenticated upgrade

Authenticate on the route before upgrade (e.g. via a JWT hoop). The handler then reads claims from Depot:

#[handler]
async fn ws_auth(req: &mut Request, depot: &mut Depot, res: &mut Response)
    -> Result<(), StatusError>
{
    let user_id = depot.jwt_auth_data::<Claims>()
        .ok_or_else(StatusError::unauthorized)?
        .claims.user_id;
    WebSocketUpgrade::new()
        .upgrade(req, res, move |mut ws| async move {
            while let Some(Ok(msg)) = ws.recv().await {
                if ws.send(msg).await.is_err() { break; }
            }
            let _ = user_id;
        })
        .await
}

Message helpers

Message supports text, binary, ping, pong, close, close_with(code, reason). Pings are auto-ponged; handle is_close() explicitly.

while let Some(Ok(msg)) = ws.recv().await {
    if msg.is_text() {
        let text = msg.as_str().unwrap_or_default();
        ws.send(Message::text(format!("echo: {text}"))).await.ok();
    } else if msg.is_binary() {
        ws.send(Message::binary(msg.as_bytes().to_vec())).await.ok();
    } else if msg.is_close() {
        break;
    }
}

as_str() returns Err for non-text messages — use is_text() or match on the type first.

Broadcasting pattern (channel + map of senders)

Use an mpsc channel per client. The receiver is forwarded into the socket sink, so broadcasters never hold a lock while sending:

use std::collections::HashMap;
use std::sync::LazyLock;
use std::sync::atomic::{AtomicUsize, Ordering};
use futures_util::{FutureExt, StreamExt};
use salvo::prelude::*;
use salvo::websocket::{Message, WebSocket, WebSocketUpgrade};
use tokio::sync::{RwLock, mpsc};
use tokio_stream::wrappers::UnboundedReceiverStream;

type Tx = mpsc::UnboundedSender<Result<Message, salvo::Error>>;
static NEXT_ID: AtomicUsize = AtomicUsize::new(1);
static USERS: LazyLock<RwLock<HashMap<usize, Tx>>> = LazyLock::new(Default::default);

#[handler]
async fn chat(req: &mut Request, res: &mut Response) -> Result<(), StatusError> {
    WebSocketUpgrade::new().upgrade(req, res, handle_socket).await
}

async fn handle_socket(ws: WebSocket) {
    let my_id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
    let (ws_tx, mut ws_rx) = ws.split();

    let (tx, rx) = mpsc::unbounded_channel();
    tokio::spawn(UnboundedReceiverStream::new(rx).forward(ws_tx).map(|_| ()));

    USERS.write().await.insert(my_id, tx);

    while let Some(Ok(msg)) = ws_rx.next().await {
        if let Ok(text) = msg.as_str() {
            let out = format!("<User#{my_id}>: {text}");
            for (&uid, peer_tx) in USERS.read().await.iter() {
                if uid != my_id {
                    let _ = peer_tx.send(Ok(Message::text(out.clone())));
                }
            }
        }
    }

    USERS.write().await.remove(&my_id);
}

Gotchas:

  • ws.split() comes from StreamExt::split; items flowing into the sink must be Result<Message, salvo::Error>, so the channel element type matches.
  • Don't hold the RwLock write guard across .await on the send side — keep the read guard short-lived.

Rooms

Wrap a HashMap<String, HashMap<usize, Tx>> in Arc<RwLock<...>> and share via Depot or a LazyLock. The insert/remove/broadcast pattern mirrors the flat map above.

Heartbeat via tokio::select!

Client pings are auto-ponged; use server-side pings only if you need to detect silent peers.

use std::time::Duration;
use futures_util::{SinkExt, StreamExt};
use tokio::time::interval;

async fn with_heartbeat(ws: salvo::websocket::WebSocket) {
    use salvo::websocket::Message;
    let (mut tx, mut rx) = ws.split();
    let hb = async move {
        let mut ticker = interval(Duration::from_secs(30));
        loop {
            ticker.tick().await;
            if tx.send(Message::ping(Vec::new())).await.is_err() { break; }
        }
    };
    let recv = async move {
        while let Some(Ok(msg)) = rx.next().await {
            if msg.is_close() { break; }
        }
    };
    tokio::select! { _ = hb => {}, _ = recv => {} }
}

JSON messages

#[derive(serde::Serialize, serde::Deserialize)]
#[serde(tag = "type")]
enum WsMsg {
    Chat { content: String },
    Join { room: String },
}

// Decode:
if let Ok(text) = msg.as_str() {
    if let Ok(parsed) = serde_json::from_str::<WsMsg>(text) { /* ... */ }
}
// Encode:
let json = serde_json::to_string(&WsMsg::Chat { content: "hi".into() }).unwrap();
ws.send(salvo::websocket::Message::text(json)).await.ok();

Gotchas

  • recv() returns None once the stream ends — always exit the loop.
  • upgrade() spawns the callback as a separate task; anything captured must be Send + 'static.
  • Sec-WebSocket-Protocol is echoed to clients — never put secrets in it.
  • If the client omits Sec-WebSocket-Version: 13 or the Upgrade/Connection headers, upgrade() returns StatusError::bad_request.

Related Skills

  • salvo-sse: Server-Sent Events for unidirectional updates
  • salvo-realtime: Overview of real-time communication options
  • salvo-auth: Authenticate WebSocket connections

适合场景

01

用户想查找某类 Agent Skill 时

02

需要根据任务场景推荐可安装能力包时

03

需要对比不同来源的安装命令和来源信息时

能力概览

能力 1

按任务关键词查找相关 Skills

能力 2

展示可复制的安装命令

能力 3

保留来源站点、仓库和原始说明,方便继续核验

能力 4

展示第三方安全扫描或审计结果

安装后应在对应宿主中按原始 README 的触发条件使用;具体调用方式请以来源页面和 README 为准。

平台分布

Codex

32.56%
按下载量换算34

Claude

30.91%
按下载量换算32

Cursor

20%
按下载量换算21

Gemini CLI

9.99%
按下载量换算10

安全审计

Gen Agent Trust Hub

通过

Socket

通过

Snyk

可疑

权限和风险

执行命令

安装流程涉及命令执行,可能通过 npx skills add https://github.com/salvo-rs/salvo-skills --skill salvo-websocket 联网下载 Skill 或依赖。用户安装前应确认命令来源、仓库内容和执行环境。

安装前确认

本站仅展示第三方公开信息,不托管安装包,不提供自动安装或运行环境。安装前应自行审查源码、依赖和命令行为。来源安全扫描存在 warning/failed 结果,不能写成本站确认安全。当前只有一个来源,正式发布前建议补源仓库或其他目录站核验。

来源信息

继续浏览同类 Skills