【问题标题】:Tokio Brodcast Channel Receiver not receivingTokio Broadcast Channel Receiver 收不到
【发布时间】:2023-01-10 04:55:58
【问题描述】:

我目前正在尝试使用 Tokio 和广播频道编写服务器和客户端。我有一个基本上监听连接的循环,在读取 TcpStream 之后,我通过通道发送。

这是我尝试过的代码:

每次我连接到服务器并读取字节时,我最终得到的是打印..,但我从来没有收到“已收到”

use dbjade::serverops::ServerOp;
use tokio::io::{BufReader};
use tokio::net::TcpStream;
use tokio::{net::TcpListener, io::AsyncReadExt};
use tokio::sync::broadcast;

const ADDR: &str = "localhost:7676"; // Your own address : TODO change to be configured
const CHANNEL_NUM: usize = 100;
use std::io;
use std::net::{SocketAddr};
use bincode;


#[tokio::main]
async fn main() {
     // Create listener instance that bounds to certain address
    let listener = TcpListener::bind(ADDR).await.map_err(|err|  panic!("Failed to bind: {err}")).unwrap();
    let (tx, mut rx) = broadcast::channel::<(ServerOp, SocketAddr)>(CHANNEL_NUM);
    

    loop {
        if let Ok((mut socket, addr)) = listener.accept().await {
            let tx = tx.clone();
            let mut rx = tx.subscribe();
            println!("Receieved stream from: {}", addr);
            let mut buf = vec![0, 255];
            tokio::select! {
                result = socket.read(&mut buf) => {
                    match result {
                        Ok(res) => println!("Bytes Read: {res}"),
                        Err(_) => println!(""),
                    }
                    tx.send((ServerOp::Dummy, addr)).unwrap();
                }
                result = rx.recv() =>{
                    let (msg, addr) = result.unwrap();
                    println!("Receieved: {msg}");
                }
            }
        }
    }
}

【问题讨论】:

  • 我不知道这是否是您问题的根源,但据我所知,read() 不是取消安全的 - 您不应该在 select 中使用它。

标签: rust rust-tokio


【解决方案1】:

您代码中的主要问题是这一行

            let tx = tx.clone();
            let mut rx = tx.subscribe();

您正在重新定义 tx 和 rx。并且您在循环中执行此操作,因此下一次迭代永远不会有相同的 tx 和 rx,因此它们无法在迭代之间连接。所以当你执行 rx.recv() 时,它不是通道另一端的 rx。您在开始时定义的 rx 未使用。根据我的经验,阴影是生锈的常见问题。解决它的一般方法是阅读编译器的所有警告并解决所有“未使用”变量、导入等。这就是我所做的:我删除了所有未使用的东西并连接了正确的通道端。我还删除了 dbjade,因为我知道从哪里获取它,并且为了示例将其替换为“Dummy”字符串。

use tokio::{net::TcpListener, io::AsyncReadExt};
use tokio::sync::broadcast;

const ADDR: &str = "localhost:7676"; // Your own address : TODO change to be configured
const CHANNEL_NUM: usize = 100;
use std::net::{SocketAddr};


#[tokio::main]
async fn main() {
    // Create listener instance that bounds to certain address
    let listener = TcpListener::bind(ADDR).await.map_err(|err|  panic!("Failed to bind: {err}")).unwrap();
    let (tx, mut rx) = broadcast::channel::<(String, SocketAddr)>(CHANNEL_NUM);


    loop {
        if let Ok((mut socket, addr)) = listener.accept().await {
            println!("Receieved stream from: {}", addr);
            let mut buf = vec![0, 255];
            tokio::select! {
                result = socket.read(&mut buf) => {
                    match result {
                        Ok(res) => println!("Bytes Read: {res}"),
                        Err(_) => println!("Err"),
                    }
                    tx.send(("Dummy".to_string(), addr)).unwrap();

                }
                result = rx.recv() =>{
                    let (msg, _) = result.unwrap();
                    println!("Receieved: {msg}");
                }
            }
        }
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2011-03-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-11-07
    • 1970-01-01
    • 2021-02-15
    相关资源
    最近更新 更多