【问题标题】:The shared mutex problem in Rust (implementing AsyncRead/AsyncWrite for Arc<Mutex<IpStack>>)Rust 中的共享互斥锁问题(为 Arc<Mutex<IpStack>> 实现 AsyncRead/AsyncWrite)
【发布时间】:2021-05-26 00:27:10
【问题描述】:

假设我有一个用户空间 TCP/IP 堆栈。我很自然地将它包装在Arc&lt;Mutex&lt;&gt;&gt; 中,以便与我的线程分享。

我想为它实现AsyncReadAsyncWrite 也是很自然的,所以像hyper 这样期望impl AsyncWriteimpl AsyncRead 的库可以使用它。

这是一个例子:

use core::task::Context;
use std::pin::Pin;
use std::sync::Arc;
use core::task::Poll;
use tokio::io::{AsyncRead, AsyncWrite};

struct IpStack{}

impl IpStack {
    pub fn send(self, data: &[u8]) {
        
    }
    
    //TODO: async or not?
    pub fn receive<F>(self, f: F) 
        where F: Fn(Option<&[u8]>){
        
    }
}

pub struct Socket {
    stack: Arc<futures::lock::Mutex<IpStack>>,
}

impl AsyncRead for Socket {
    fn poll_read(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buf: &mut tokio::io::ReadBuf<'_>
    ) -> Poll<std::io::Result<()>> {
        //How should I lock and call IpStack::read here?
        Poll::Ready(Ok(()))
    }
}

impl AsyncWrite for Socket {
    fn poll_write(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buf: &[u8],
    ) -> Poll<Result<usize, std::io::Error>> {
        //How should I lock and call IpStack::send here?
        Poll::Ready(Ok(buf.len()))
    }
    //poll_flush and poll_shutdown...
}

Playground

我认为我的假设没有任何问题,也没有其他更好的方法来与多个线程共享堆栈,除非我将其包装在 Arc&lt;Mutex&lt;&gt;&gt;

这类似于引起我兴趣的try_lock on futures::lock::Mutex outside of async?

我应该如何在不阻塞的情况下锁定互斥锁?请注意,一旦我获得锁,IpStack 就不是异步的,它调用了该块。我也想对其实现异步,但我不知道问题会变得更加困难。或者如果它有异步调用,问题会变得更简单吗?

【问题讨论】:

    标签: multithreading rust concurrency mutex


    【解决方案1】:

    我发现tokio::sync::Mutex 上的 tokio 文档页面非常有用:https://docs.rs/tokio/1.6.0/tokio/sync/struct.Mutex.html

    从你的描述看来你想要:

    • 非阻塞操作
    • 一种管理用户空间 TCP/IP 堆栈管理的所有 IO 资源的大数据结构
    • 跨线程共享一个大数据结构

    我建议探索像演员这样的东西,并使用消息传递与为管理 TCP/IP 资源而产生的任务进行通信。我认为您可以包装类似于tokio 文档中引用的mini-redis 示例的API 来实现AsyncReadAsyncWrite。从返回完整结果的未来的 API 开始,然后处理流式传输可能更容易。我认为这更容易纠正。使用loom 锻炼它可能会很有趣。

    我认为,如果您打算通过互斥锁同步对 TCP/IP 堆栈的访问,您最终可能会得到一个Arc&lt;Mutex&lt;...&gt;&gt;,但会使用一个包装互斥锁的 API,例如 mini-redistokio 文档提出的建议是,他们的 Mutex 实现更适合管理 IO 资源而不是共享原始数据,我认为这确实适合您的情况。

    【讨论】:

    • 我看不出tokio::sync::Mutex 有什么帮助,因为AsyncRead 是可轮询的,我必须以非异步方式实现轮询。同样只有 github.com/tokio-rs/mini-redis/blob/… 在 mini-redis 上使用 Mutex,我没有看到任何使用它的异步操作
    【解决方案2】:

    您应该为此使用异步互斥锁。使用标准的std::sync::Mutex

    futures::lock::Mutextokio::sync::Mutex 这样的异步互斥锁允许等待锁定而不是阻塞,因此它们可以安全地在 async 上下文中使用。它们旨在用于awaits。这正是您不希望发生的事情!锁定 await 意味着互斥锁可能被锁定很长时间,并且会阻止其他想要使用 IpStack 的异步任务取得进展。

    实现AsyncRead/AsyncWrite在理论上是直截了当的:要么它可以立即完成,要么它通过某种机制协调以在数据准备好并立即返回时通知上下文的唤醒器。这两种情况都不需要扩展使用底层IpStack,因此使用非异步互斥体是安全的。

    use std::pin::Pin;
    use std::sync::{Arc, Mutex};
    use std::task::{Context, Poll};
    use tokio::io::{AsyncRead, AsyncWrite};
    
    struct IpStack {}
    
    pub struct Socket {
        stack: Arc<Mutex<IpStack>>,
    }
    
    impl AsyncRead for Socket {
        fn poll_read(
            self: Pin<&mut Self>,
            cx: &mut Context<'_>,
            buf: &mut tokio::io::ReadBuf<'_>,
        ) -> Poll<std::io::Result<()>> {
            let ip_stack = self.stack.lock().unwrap();
    
            // do your stuff
    
            Poll::Ready(Ok(()))
        }
    }
    

    【讨论】:

    • 对不起,我不明白,为什么锁定 poll_read 是可以的?
    • 我不确定如何更好地解释它。如果您的计划是在poll_read 内等待结果,那么您做错了。无论数据是否可以立即读取,该函数都应该快速返回。既然这是目标,阻塞互斥体就非常适合。
    【解决方案3】:

    除非我将堆栈封装在 Arc&lt;Mutex&lt;&gt;&gt; 中,否则我没有看到其他更好的方法来与多个线程共享堆栈。

    Mutex 无疑是实现此类操作的最直接方式,但我建议控制反转。

    在基于Mutex 的模型中,IpStack 实际上是由Sockets 驱动的,它们将IpStack 视为共享资源。这会导致一个问题:

    • 如果 Socket 在锁定堆栈时发生阻塞,则它违反了 AsyncRead 的约定,因为它会花费无限的时间来执行。
    • 如果Socket 没有 阻止锁定堆栈,而是选择使用try_lock(),它可能会因为它没有保持“排队”锁定而被饿死。如果您不等待,公平的锁定算法(例如 parking_lot 提供的算法)无法让您免于饥饿。

    相反,您可以按照系统网络堆栈的方式处理问题。套接字不是参与者:网络堆栈驱动套接字,而不是相反。

    实际上,这意味着IpStack 应该有一些方法来轮询套接字以确定下一个要写入/读取的套接字。用于此目的的操作系统接口虽然不能直接应用,但可能会提供一些启发。经典地,BSD 提供了select(2)poll(2);现在,像epoll(7)(Linux) 和kqueue(2)(FreeBSD) 这样的API 更适合大量连接。

    一个非常简单的策略,松散地以select/poll 为模型,以循环方式重复扫描Socket 连接列表,一旦可用就处理它们的待处理数据。

    对于一个基本的实现,一些具体的步骤是:

    • 在创建新的Socket 时,会在它和IpStack 之间建立一个双向通道(即每个方向一个有界通道)。
    • AsyncWriteSocket 上尝试通过传出通道将数据发送到IpStack。如果频道已满,请返回Poll::Pending
    • AsyncReadSocket 上尝试通过传入通道从IpStack 接收数据。如果频道为空,则返回Poll::Pending
    • IpStack 必须由外部驱动(例如,在另一个线程上的事件循环中)以不断轮询打开的套接字以获取可用数据,并将传入数据传递到正确的套接字。通过允许IpStack 控制发送哪些套接字的数据,您可以避免Mutex 解决方案的饥饿问题。

    【讨论】:

    • 确实,堆栈已经像这样工作了:github.com/lucaszanella/smoltcp/blob/ip-interface-alt-managed/…。没有套接字的概念,您只需将一个处理程序(只是一个整数)传递给堆栈,以便它知道要通过哪个套接字发送数据。这个堆栈有一个轮询所有套接字的循环。这是一个示例,我可以进行更有效的轮询,但通常您必须主动轮询套接字。
    • 我试图不使用通道,因为这会导致内存复制。例如,如果我有一个视频流,我会通过通道将数据包从线程复制到堆栈,而在 Mutex 模型中,我只需锁定并立即发送到堆栈。
    • 所以是的,堆栈有轮询功能,我可以看到可以读取哪些。但是我认为当有数据要读取时,我必须立即读取它,否则数据将在下一次投票中被覆盖。
    • 在尝试实施您的建议时,我注意到它可能并不那么简单。例如,当 tokio 调用 poll_read 并返回 Poll::Pending 时,tokio 不会永远调用它。它可能会再试一两次,但不会再试一次,而是依赖poll_read 函数保存了Wake 对象,以便在有数据时回调。所以我想我至少需要在每次poll_read 调用时在本地存储一个Waker。你同意吗?
    猜你喜欢
    • 2021-10-10
    • 2012-06-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-06-04
    • 2011-12-28
    • 1970-01-01
    • 2011-07-22
    相关资源
    最近更新 更多