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

如何将oneshot receiver传入serve_with_shutdown作为停止信号?

解决Tokio serve_with_shutdown传入oneshot Receiver的问题

首先明确核心问题:Tokio的tokio::sync::oneshot::Receiver<()>作为Future时,输出类型是Result<(), oneshot::Error>,但serve_with_shutdown要求传入的signal Future的输出类型是(),这才是直接传入报错的根本原因(并非提示的"Future未被使用")。

正确解决方案:转换Receiver的输出类型

不需要自行实现Future,直接对Receiver进行输出转换即可,用async块就能简洁完成:

use tokio::sync::oneshot;
use tokio::net::TcpListener;
use hyper::{Request, Response, Body, service::Service};

async fn run_server() -> Result<(), Box<dyn std::error::Error>> {
    let listener = TcpListener::bind("127.0.0.1:8080").await?;
    let (sender, receiver) = oneshot::channel::<()>();

    // 将Receiver转换为Output为()的Future:忽略接收结果(无论成功/失败都触发服务终止)
    let shutdown_signal = async {
        let _ = receiver.await;
    };

    // 传入转换后的signal Future到serve_with_shutdown
    listener.serve_with_shutdown(
        |req: Request<Body>| async move {
            // 示例请求处理逻辑
            Ok::<_, std::convert::Infallible>(Response::new(Body::from("Hello World")))
        },
        shutdown_signal
    ).await?;

    Ok(())
}

你的自定义Future报错原因

你实现的OneShotFut存在两个关键问题:

  1. 不必要使用Mutex:tokio::sync::oneshot::Receiver本身是线程安全且支持Pin的,额外加锁完全多余。
  2. Poll逻辑错误:直接调用try_recv不会向任务调度器注册唤醒通知,导致Future永远处于Pending状态。正确做法是委托Receiver自身的Poll方法:
use std::pin::Pin;
use std::task::{Context, Poll};
use futures::Future;
use tokio::sync::oneshot;

pub struct OneShotFut {
    receiver: oneshot::Receiver<()>,
}

impl Future for OneShotFut {
    type Output = ();

    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
        // 委托Receiver的Poll方法,自动处理唤醒逻辑
        match Pin::new(&mut self.receiver).poll(cx) {
            Poll::Ready(_) => Poll::Ready(()),
            Poll::Pending => Poll::Pending,
        }
    }
}

但这种自定义实现完全没必要,用async块转换更简洁直观。

补充说明

  • serve_with_shutdown会在signal Future就绪时停止服务,只要signal输出为()即可,无需关心触发终止的具体原因。
  • 如果需要仅在成功接收停止信号时才终止服务,可以在async块中添加判断:
    let shutdown_signal = async {
        if receiver.await.is_ok() {
            // 仅当成功收到停止指令时终止服务
        }
    };
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 01:10:13