You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Tokio mpsc Receiver非阻塞接收:单线程处理非线程安全外设收发方案

问题根因

device_query库内部持有X11原生裸指针,未实现Send trait。而tokio默认多线程调度器会在工作线程间转移异步任务,要求所有跨await存活的变量必须满足Send约束,因此编译报错。

解决方案

方案1:使用tokio本地任务(保留原有异步逻辑)

将设备任务绑定到单线程执行,不需要Send约束,可正常使用await:

  1. 将tokio::task::spawn替换为tokio::task::spawn_local,该API创建的任务只会在生成它的线程上运行,不会跨线程转移。
  2. 调整tokio运行时为单线程模式,确保本地任务可正常调度。

修改后代码:

use device_query::{DeviceState, Keycode};
use std::time::Duration;
use tokio;
use tokio::sync::mpsc::{Receiver, Sender};
use tokio::time::timeout;

// 指定tokio使用单线程运行时
#[tokio::main(flavor = "current_thread")]
async fn main() {
    // 设备循环产出的按键事件通道
    let (key_tx, mut key_rx) = tokio::sync::mpsc::channel(32);
    // 发往设备循环的命令通道
    let (dev_tx, mut dev_rx) = tokio::sync::mpsc::channel(32);

    start_device_loop(60, key_tx, dev_rx);

    println!("Waiting for key presses");
    while let Some(k) = key_rx.recv().await {
        match k {
            Some(ch) => match ch {
                Keycode::Q => dev_tx.clone().try_send(String::from("Quit!")).expect("Could not send command"),
                ch => println!("{}", ch),
            },
            _ => (),
        }
    }
    println!("Done.")
}

/// 启动tokio本地任务,轮询指定设备并通过传入的mpsc sender发送按键事件
pub fn start_device_loop(hz: u32, tx: Sender<Option<Keycode>>, mut rx: Receiver<String>) {
    let poll_wait = 1000 / hz;
    let poll_wait = Duration::from_millis(poll_wait as u64);

    //  spawn_local创建的任务不需要实现Send
    tokio::task::spawn_local(async move {
        let dev = DeviceState::new();

        loop {
            let mut keys = dev.query_keymap();
            match keys.len() {
                0 => (),
                1 => tx.clone().try_send(Some(keys.remove(0))).unwrap(),
                _ => println!("So many keys..."),
            }
            
            match timeout(poll_wait, rx.recv()).await {
                Ok(Some(cmd)) => println!("Command '{}' received.", cmd),
                // 命令通道关闭时退出任务
                Ok(None) => break,
                Err(_) => (),
            };
        }
    });
}

方案2:使用原生std线程(完全不使用async/await)

如果不需要异步能力,直接开启系统原生线程处理设备逻辑,用非阻塞方法读取命令:

pub fn start_device_loop(hz: u32, tx: Sender<Option<Keycode>>, mut rx: Receiver<String>) {
    let poll_wait = 1000 / hz;
    let poll_wait = Duration::from_millis(poll_wait as u64);

    std::thread::spawn(move || {
        let dev = DeviceState::new();

        loop {
            let mut keys = dev.query_keymap();
            match keys.len() {
                0 => (),
                1 => tx.clone().try_send(Some(keys.remove(0))).unwrap(),
                _ => println!("So many keys..."),
            }
            
            // 非阻塞读取命令,不会阻塞轮询
            match rx.try_recv() {
                Ok(cmd) => println!("Command '{}' received.", cmd),
                Err(tokio::sync::mpsc::error::TryRecvError::Empty) => (),
                // 命令通道关闭时退出线程
                Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => break,
            };
            // 睡眠固定时长控制轮询频率
            std::thread::sleep(poll_wait);
        }
    });
}

该方案main函数无需修改,完全不涉及异步语法,也满足单线程操作设备的要求。

内容的提问来源于stack exchange,提问作者superstator

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.06 14:39:05