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

开发基于Tower的TCP回声服务器时,Tokio TcpListener无incoming方法报错

解决Tokio TcpListener调用incoming方法报错的问题

你遇到的“未找到结构体tokio::net::TcpListener的incoming方法”报错,可通过以下两种方式解决:

方案一:修复incoming流的调用问题

TcpListener::incoming()返回的是异步流(Stream类型),要调用next().await迭代流,必须导入futures::StreamExt trait。同时确认你的依赖配置中Tokio的full特性已启用(你当前的配置没问题)。

修改后的完整代码:

/*
[dependencies]
futures = "0.3"
tokio = { version = "1", features = ["full"] }
tower = { version = "0.4", features = ["full"] }
*/

use tokio::net::TcpListener;
use tower::Service;

// 导入StreamExt以支持异步流的next()方法
use futures::{future, StreamExt};
use std::task::{Context, Poll};

struct Echo;

impl Service<Vec<u8>> for Echo {
    type Response = Vec<u8>;
    type Error = std::io::Error;
    type Future = future::Ready<Result<Self::Response, Self::Error>>;

    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
        Poll::Ready(Ok(()))
    }

    fn call(&mut self, req: Vec<u8>) -> Self::Future {
        future::ok(req)
    }
}

#[tokio::main]
async fn main() {
    let addr = "127.0.0.1:8080".parse().unwrap();
    let listener = TcpListener::bind(&addr).await.unwrap();

    let mut incoming = listener.incoming();
    while let Some(socket) = incoming.next().await {
        let socket = socket.unwrap();
        tokio::spawn(async move {
            let (reader, writer) = socket.split();
            let echo = Echo;
            tower::service_fn(move |req| echo.call(req)).serve(reader, writer).await
        });
    }
}

方案二:使用Tower原生服务编排(更推荐)

Tower提供了tower::serve函数,可直接将服务绑定到TcpListener上,自动处理连接监听、流迭代和服务分发,无需手动调用incoming(),代码更简洁且符合Tower的设计理念:

/*
[dependencies]
futures = "0.3"
tokio = { version = "1", features = ["full"] }
tower = { version = "0.4", features = ["full"] }
*/

use tokio::net::TcpListener;
use tower::{Service, service_fn};

use futures::future;
use std::task::{Context, Poll};

struct Echo;

impl Service<Vec<u8>> for Echo {
    type Response = Vec<u8>;
    type Error = std::io::Error;
    type Future = future::Ready<Result<Self::Response, Self::Error>>;

    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
        Poll::Ready(Ok(()))
    }

    fn call(&mut self, req: Vec<u8>) -> Self::Future {
        future::ok(req)
    }
}

#[tokio::main]
async fn main() {
    let addr = "127.0.0.1:8080".parse().unwrap();
    let listener = TcpListener::bind(&addr).await.unwrap();

    // Tower自动处理连接生命周期和服务调用
    tower::serve(listener, service_fn(|req: Vec<u8>| {
        let echo = Echo;
        echo.call(req)
    })).await.unwrap();
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 08:40:46