30歳からのプログラミング

30歳無職から独学でプログラミングを開始した人間の記録。

TCP エコーサーバで理解する readiness-based と completion-based

I/O 処理は、ソフトウェアを設計する際の重要な論点のひとつ。I/O 処理を上手く設計しないと、処理能力やリソース効率などに問題のあるソフトウェアになってしまう可能性がある。
この記事では、 I/O 処理の代表的な設計パターンである readiness-based と completion-based について、それぞれの考え方や特徴を見ていく。

readiness-based は epoll で、 completion-based は io_uring で実現できるので、これらを使って実装した TCP エコーサーバを例として示す。

epoll や io_uring については以下を参照。

numb86-tech.hatenablog.com

numb86-tech.hatenablog.com

この記事における「I/O 処理」は、プロセスがプロセス外の世界(ディスクやネットワーク、他のプロセスなど)とやり取りする処理のことを指す。TCP エコーサーバの場合、accept(), recv(), send()などのシステムコールで行われる処理がそれに該当する。
TCP 通信の流れやそこで使われるシステムコールについては以下で解説している。

numb86-tech.hatenablog.com

動作確認は以下の環境で行った。

$ lsb_release -a
No LSB modules are available.
Distributor ID: Ubuntu
Description:    Ubuntu 24.04.4 LTS
Release:    24.04
Codename:   noble
$ rustc --version
rustc 1.97.1 (8bab26f4f 2026-07-14)

$ cargo --version
cargo 1.97.1 (c980f4866 2026-06-30)

Rust の Edition は2024。

使っているクレートとそのバージョンは以下。

  • libc@0.2.189

解決したい課題は何なのか

readiness-based と completion-based はよく対比されるのだが、扱っている課題はどちらも同じである。同じ課題に対して異なるアプローチで対応しているからこそ、違いが分かりやすく、セットで紹介されることが多い。
そのためまずはこの「課題」について説明しておく。

カーネルに依頼する I/O 処理のなかには、条件を満たさないといつまでも完了しないものがある。そしてさらにそのなかには、完了条件がプロセス外の出来事に依存しているものがある。
例えばaccept()は、対象のリスナーソケットの accept queue からデータを取り出すシステムコールだが、 accept queue に何もデータが入っていない場合、データが入るまで待機し続けることになる。そして accept queue にデータが入るのは、当該リスナーソケットに対して接続要求があり、その接続が確立されたときである。つまり、 accept queue が空の状態でaccept()を呼び出してしまうと、 TCP クライアントからの接続要求を待ち続けることになる。プロセスの処理がそこで停止してしまう。そして、「TCP クライアントから接続要求が来てその接続が確立される」という、プロセス自身にはどうにもできない条件が成立するのをただ待つことになる。
recv()やsend()も同様の性質を持つ。

I/O 処理が持つこのような性質を考慮せず素朴な実装にしてしまうと、処理が完了するまでに時間が掛かったり、処理が止まったままになってしまったりする。
これが readiness-based や completion-based が解決する「課題」である。

ノンブロッキングモードについて

ファイルディスクリプタにはノンブロッキングモードというものがあり、これを使えば、上述したような事象を避けることはできる。

ノンブロッキングモードについては、以下の記事で解説している。

numb86-tech.hatenablog.com

詳しくは上記の記事を読んで欲しいが、ノンブロッキングモードを使うと、「条件が成立するまで待つ」という挙動から、「条件が成立していない場合は即座にエラーを返してそれをプロセスに伝える」という挙動に変わる。
例えば、ノンブロッキングモードでaccept()を呼び出し、 accept queue が空だった場合、待機するのではなく、即座にプロセスに処理が戻る。そして呼び出しはエラー扱いとなり、プロセスはerrnoという変数を通して、「accept queue が空だった」ことを知ることができる。つまり、プロセスが停止してしまう、長時間の待ちが発生してしまう、という問題は起きなくなる。

しかし別の課題が発生する。既に述べたように、条件が成立していない場合は停止せず処理が戻ってくるわけだが、そうすると今度は、いつになったら条件が成立するのかが分からない。今条件を満たしていない、ということは分かるが、いつ満たすのかは分からない。そのため、時間を置いて再試行してみる、定期的に実行を試みる、というアプローチを採用することになる。しかしそれは単純に非効率であるし、そもそも、「条件を満たすまで待機させられる」が「条件を満たすまで再実行する」に変わっただけであり、本質的な解決にはなっていない。

readiness-based や completion-based で I/O 処理を設計すれば、この課題も含めて解決できる。

TCP エコーサーバについての前置き

これから TCP エコーサーバの例を用いながら説明していくが、 TCP エコーサーバという性質上、 TCP クライアントからの接続要求やデータ受信があるまで待機することになる。
だがそれは、上述した「条件が成立していない状況で I/O 処理を行おうとしてしまったため長時間プロセスが停止してしまう」「いつ条件が成立するのか分からない」という課題とは意味合いが異なる。
やるべき処理がないから待機しているだけであり、 I/O 処理の設計不備により処理の停止や無駄が発生してしまっているわけではない。

条件が成立したという通知を受けてから I/O 処理を行う readiness-based

readiness-based は、条件が成立したら通知してくれとカーネルに依頼することで、課題を解決する。
accept()の例でいえば、条件が成立しているかどうか分からない状態でaccept()を呼び出すのではなく、「条件が成立した」という通知をカーネルから受けた上でaccept()を呼び出すようにする。

epoll を使って実装した readiness-based の TCP エコーサーバが以下。

github.com

このプログラムを起動すると、以下のコードが実行される。

epoll.add(listen_fd, libc::EPOLLIN as u32)?;

listen_fdはリスナーソケットを指し示すファイルディスクリプタであり、そのファイルディスクリプタで発生するEPOLLINイベントを、 epoll による監視対象として登録している。
「リスナーソケットの accept queue へのデータの追加」も、EPOLLINイベントのひとつである。そのため、 accept queue にデータが追加されたことを、 epoll の通知によって知ることができる。

そのあとこのプログラムはloopによる無限ループに入る。
この無限ループはlet n = epoll.wait()?;したあとにfor i in 0..n {}を実行する、という構造になっている。
let n = epoll.wait()?;は、 epoll が監視対象のイベントを検知するまで待機する。そしてイベントを検知すると、検知したイベントの数がnに束縛される。
その後n回forループを回す。このforループのなかで、発生したイベントごとに、必要な処理を行っていく。
そして処理が終わればloopの先頭に戻り、再びイベントが発生するまで待機する。これがこの TCP エコーサーバの大まかな構造である。

loop {
    // 監視対象のイベントをカーネルが検知し通知してくるまで待機する
    let n = epoll.wait()?;

    // 検知したイベントの数だけループを回す
    for i in 0..n  {
        // イベントが発生したのはどのファイルディスクリプタなのかを取得する
        let ev = epoll.events[i];
        let fd = ev.u64 as u32 as RawFd;

        if fd == listen_fd {
            // イベントが発生したのがリスナーソケットだった場合の処理をここに書く
            continue;
        }

        // 以降は、リスナーソケット以外のソケットが対象だった場合の処理

        // 対象の TCP 接続に関する情報を取得
        let conn = conns.get_mut(&fd).expect("event for unknown fd");

        match conn.state {
            ConnState::WaitReadable => {
                // TCP クライアントからデータを受信したときの処理をここに書く
            }
            ConnState::WaitWritable => {
                // TCP クライアントにデータを送信するときの処理をここに書く
            }
        }
    }
}

1 周目のloopのときは、既に見たように「リスナーソケットで発生するEPOLLINイベント」が監視対象として登録されている。そのため、 accept queue へのデータの追加、すなわち TCP クライアントとの接続確立が発生すると、epoll.wait()から処理が戻ってくる。
その後、forループ内のif fd == listen_fd {}ブロックが実行される。ここで、accept()を呼び出す(正確にはaccept4()だが、ここではaccept()とほぼ同じものだと考えてよい)。この処理が実行されているということは、 accept queue へのデータ追加という「条件」が成立しているわけなので、accept()の呼び出し後に無駄な「待ち」が発生することはない。
その後、クライアントと接続確立されているソケットで発生するEPOLLINイベントも epoll の監視対象として登録した上で、loopの先頭に戻る。そうするとそれ以降は、新しいクライアントと接続確立したときだけでなく、接続確立済みのクライアントからデータが送られてきたり、そのクライアントが接続を切断したりしたときも、epoll.wait()が値を返すようになる。
このように、 epoll からの通知を待つ、通知内容に応じて処理を行い epoll の監視対象を更新する、再び epoll からの通知を待つ、ということを繰り返すのが、この TCP エコーサーバの大まかな処理の流れである。

accept queue にデータが追加されたという通知を受けてからaccept()を呼び出す、ソケットへの書き込みが可能になったという通知を受けてからsend()を呼び出す、のように、 I/O 処理が完了するための条件を満たしたという通知を受けてから I/O 処理を行うようになっている。これが readiness-based のアプローチである。

条件が成立しているかどうかを関知せず全てをカーネルに任せる completion-based

completion-based は根本的にアプローチを変え、 I/O 処理を非同期処理にする。

accept()のようなシステムコールを呼び出して I/O 処理を依頼しそれが終わるのを待つ、のではなく、プロセスはカーネルに対して I/O 処理の依頼だけを行う。I/O 処理の完了やカーネルからの返答を、待たない。プロセスは依頼だけをして次の処理に取り掛かり、それと並行してカーネルが I/O 処理を行う。そしてその後、プロセスが任意のタイミングで、依頼しておいた I/O 処理が完了したかを確認する。
依頼を受けたカーネルは、条件が成立しているならそのまま処理するし、条件が成立していない場合は成立したタイミングで処理を行うが、そういったことにプロセスは関知しない。完了までやってくれ、と依頼するだけ。そのため、「条件が成立していない状況で I/O 処理を依頼してしまうと、プロセスが長時間停止してしまう恐れがある」という課題は発生しなくなる。completion-based においては、「プロセスが I/O 処理が完了するのを待つ」ということ自体が発生しない。

Linux の場合、 io_uring を使うことで completion-based のプログラムを実装できる。

以下が TCP エコーサーバの実装例。

github.com

プログラムの大まかな構造は epoll 版と同じにしている。
リスナーソケットの準備を行い、loopによる無限ループに入る。無限ループのなかでは、対応する I/O 処理の種類に応じて分岐が発生する。
だが I/O 処理のやり方は completion-based になっている。

io_uring では、 SQ ring というデータ構造にカーネルへの依頼内容を書き込み、カーネルは依頼内容を処理した結果を CQ ring というデータ構造に書き込む。
例えばloopに入る前に行われる「リスナーソケットの準備」は、リスナーソケットを用意し、それに対してaccept()相当の操作を行う依頼を SQ ring に書き込むことを意味する。

その後loopに入るが、そこではまずring.submit_and_wait()?;し、そのあとにwhile let Some(cqe) = ring.pop() {}を実行する、という構造になっている。このwhileは、ring.pop()がSome(cqe)を返す限りループを繰り返す、という構文であり、これにより、未処理のデータが CQ ring から無くなるまでループが繰り返されるようになっている。

epoll 版ではepoll.wait()でクライアントからの接続やデータ受信などを待ち、それが発生すると処理が進みforループの実行に移っていたが、それと似た構造になっている。
しかし「待っている」対象は異なる。ring.submit_and_wait()?;は、既に SQ ring に書き込み済みの依頼をカーネルに依頼しつつ、最低でも 1 件の依頼が完了し CQ ring に書き込まれることを待つ。つまりring.submit_and_wait()?;から処理が戻ってきたということは、カーネルに依頼した処理のうち少なくとも 1 つが完了したということを意味する。
例えば 1 周目のloopのときは、 SQ ring には、先程説明した「リスナーソケットに対してaccept()相当の操作を行う依頼」しか書き込まれていない。そのためring.submit_and_wait()?;が値を返したということは、「リスナーソケットに対してaccept()相当の操作を行う依頼」が完了し、その結果が CQ ring に書き込まれたことを意味する。

その後whileによるループに入るが、「どの依頼に対する結果が返ってきたのか」に応じて処理内容が変わる。「リスナーソケットに対してaccept()相当の操作を行う依頼」に対する結果の場合、match文のTAG_ACCEPTアームが実行される。そこでは、「リスナーソケットに対してaccept()相当の操作を行う依頼」が成功したのかをチェックし、成功した場合は「クライアントと接続が確立しているソケットに対してrecv()相当の操作を行う依頼」を SQ ring に書き込む。そしてそれとは別に、「リスナーソケットに対してaccept()相当の操作を行う依頼」を再度 SQ ring に書き込む。「リスナーソケットに対してaccept()相当の操作を行う依頼」を常にカーネルに依頼している状態でないと、新規の TCP クライアントからの接続に対応できないためである。

そしてloopの先頭に戻る。この時点では「リスナーソケットに対してaccept()相当の操作を行う依頼」と「クライアントと接続が確立しているソケットに対してrecv()相当の操作を行う依頼」の 2 つをカーネルに依頼している状態なので、そのどちらか、あるいは両方が完了すると、再びwhileに入る。このようなループを繰り返して、 TCP エコーサーバとしての仕事を行っていく。

loop {
    // カーネルに依頼している処理が 1 件以上完了するまで待機する
    ring.submit_and_wait()?;

    // 未処理のデータが CQ ring から無くなるまでループを回す
    while let Some(cqe) = ring.pop() {
        // どのファイルディスクリプタに対するどのような操作についてのデータなのか、という情報を得る
        let (tag, fd) = decode_user_data(cqe.user_data);

        match tag {
            TAG_ACCEPT => {
                // accept() 相当の操作が完了した際の処理をここに書く
            }
            TAG_RECV => {
                // recv() 相当の操作が完了した際の処理をここに書く
            }
            TAG_SEND => {
                // send() 相当の操作が完了した際の処理をここに書く
            }
            _ => unreachable!(),  // 想定していない値が来たらパニックを発生させる
        }
    }
}

このように、 I/O 処理を非同期でカーネルに実行してもらい、処理が完了したことを確認次第その結果に応じて次の処理を行う、というのが completion-based の基本的な構造である。

まとめ

readiness-based では、「監視対象として登録しておいたイベントが発生した」という通知を、カーネルから受け取る。その通知によって、実行したい I/O 処理を行うための条件が成立していることを確認できる。その確認をしてから、 I/O 処理をカーネルに依頼するためのシステムコールを呼び出すため、「システムコールを呼び出したはよいが、いつまでも条件が成立せずプロセスが長時間停止してしまう」という問題を回避できる。
今回実装した TCP エコーサーバの場合、監視対象のイベントが発生したという通知をカーネルから受け取ったら、通知の内容に応じて、 I/O 処理を行う。そしてその後、さらに別のイベントを監視する必要が生まれたら(accept()が成功したので次はクライアントからデータが送られてきたことを検知したい、など)、それを監視対象に加えた上で、再び通知を待つ。これを繰り返して TCP クライアントとのやり取りを進めていく。

completion-based では、処理が停止しうるシステムコールをそもそも呼び出さない。io_uring のような、非同期で I/O 処理を依頼できる仕組みを使う。そのため、実行したい I/O 処理を行うための条件が成立しているか、をプロセスが気にする必要がなくなる。
今回実装した TCP エコーサーバの場合、 TCP クライアントからの接続要求やデータ受信などが発生した際はそれを処理するよう、予めカーネルに依頼しておく。そして処理が 1 件以上完了したことを検知したら、その結果に応じて、必要な対応を行う。そして次に行われるべき I/O 処理がある場合は、やはりそれを行うようにカーネルに依頼しておく。そして再び、処理が 1 件以上完了するのを待つ。

io_uring がもたらす並行性とデータ競合

io_uring を使うと、 I/O 処理をカーネルに依頼する際に、「依頼だけを行い、その完了を待たずに自分の処理に戻る」ことができるようになる。
依頼だけを行い「自分の仕事」に戻り、それと並行してカーネルが I/O 処理を進めてくれる。
そして自分のタイミングで結果を確認し取得すればいいが、そのために別途システムコールを発行する必要はなく、スレッドとカーネルが共有しているメモリ領域を操作すればよい。
この並行性やメモリ共有という仕組みは効率性をもたらすが、それは同時にデータ競合の可能性を生む。
この記事では、 io_uring が持つこういった特徴について、具体例を示しながら見ていく。

動作確認は以下の環境で行った。

$ lsb_release -a
No LSB modules are available.
Distributor ID: Ubuntu
Description:    Ubuntu 24.04.4 LTS
Release:    24.04
Codename:   noble
$ rustc --version
rustc 1.97.1 (8bab26f4f 2026-07-14)

$ cargo --version
cargo 1.97.1 (c980f4866 2026-06-30)

Rust の Edition は2024。

使っているクレートとそのバージョンは以下。

  • libc@0.2.186

「依頼」「実行」「回収」を分離できない同期システムコール

read()のような同期システムコールでカーネルに I/O 処理を依頼すると、その I/O 処理の結果を受け取るまで、待機しないといけない。つまり、「依頼」だけを行って処理に戻ることはできず、「実行」を待ち、そしてその結果を「回収」しなければならない。
そのため「待ち時間」が発生する。完了するまでに1秒かかる場合、システムコールを呼んだスレッドは、1秒間ただ待機しなければならない。
それを再現したのが以下のコード。read()によるデータの読み取り処理を行っている。

use std::io;
use std::process;
use std::sync::atomic::{AtomicU64, Ordering};
use std::thread;
use std::time::Duration;

static WORK_COUNT: AtomicU64 = AtomicU64::new(0);

fn do_work() {
    thread::sleep(Duration::from_millis(100));
    WORK_COUNT.fetch_add(1, Ordering::Relaxed);
}

fn start_timer(budget: Duration) {
    thread::spawn(move || {
        thread::sleep(budget);
        let count = WORK_COUNT.load(Ordering::Relaxed);
        println!("time is up: work done = {count}");
        process::exit(0);
    });
}

fn main() -> io::Result<()> {
    let mut pipe_fds = [0i32; 2];
    if unsafe { libc::pipe(pipe_fds.as_mut_ptr()) } < 0 {
        return Err(io::Error::last_os_error());
    }
    let read_fd = pipe_fds[0];
    let write_fd = pipe_fds[1];

    let mut buf = vec![0u8; 256];

    thread::spawn(move || {
        thread::sleep(Duration::from_secs(1));
        let data = b"hello from pipe";
        unsafe {
            libc::write(write_fd, data.as_ptr() as *const libc::c_void, data.len());
        }
    });

    start_timer(Duration::from_secs(3));

    let n = unsafe { libc::read(read_fd, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) };
    if n < 0 {
        return Err(io::Error::last_os_error());
    }
    println!("completed: {}", String::from_utf8_lossy(&buf[..n as usize]));

    println!("starting work...");
    loop {
        do_work();
    }
}

main()のなかでタイマーを起動しているが(start_timer(Duration::from_secs(3)))、このタイマーは3秒経過したあと、プログラム全体を終了させる。
タイマーを呼び出したあとread()でパイプからデータを読み取り、それからloopに入る。loopのなかでdo_work()を繰り返しているが、これは、100ミリ秒待機してからWORK_COUNTをインクリメントする関数。
タイマーは終了する直前にWORK_COUNTの値を出力する。これで、タイマーを呼び出してから終了するまでの3秒間でどれくらいdo_work()という「仕事」を処理できたかを、計測できる。

このプログラムを実行すると、約1秒後に以下のように表示される。

completed: hello from pipe
starting work...

さらにその2秒後、以下のように表示され、プログラムが終了する。

time is up: work done = 19

上述したようにまずタイマーを起動させ、その後にread()を呼び出しているが、read()が処理を終えるのに約1秒かかるようにしてある。
そのため、read()を終えてloopに入る時点(つまりstarting work...が出力されるタイミング)で既に、タイマーの3秒のうち1秒以上が経過してしまっている。
そのため、「仕事」に使えたのは約2秒となり、do_work()を完了させることができた回数は19になる。

io_uring を使えば「依頼」を分離できる

io_uring を使うと、 I/O 処理の「依頼」だけを行い、「自分の仕事」に戻ることができる。カーネルが I/O 処理を行うのと並行して、スレッドは「自分の仕事」を進めることができる。

それを表現したのが以下のコード。

github.com

全文は長いので、main()の一部のみを載せる。

fn main() -> io::Result<()> {
    // 省略

    let mut pipe_fds = [0i32; 2];
    if unsafe { libc::pipe(pipe_fds.as_mut_ptr()) } < 0 {
        return Err(io::Error::last_os_error());
    }
    let read_fd = pipe_fds[0];
    let write_fd = pipe_fds[1];

    let mut buf = vec![0u8; 256];

    let sqe = unsafe { sqes_base.add(0) };
    unsafe {
        ptr::write_bytes(sqe, 0, 1);
        (*sqe).opcode = IORING_OP_READ;
        (*sqe).fd = read_fd;
        (*sqe).off = u64::MAX;
        (*sqe).addr = buf.as_mut_ptr() as u64;
        (*sqe).len = buf.len() as u32;
        (*sqe).user_data = 0x01;
    }

    thread::spawn(move || {
        thread::sleep(Duration::from_secs(1));
        let data = b"hello from pipe";
        unsafe {
            libc::write(write_fd, data.as_ptr() as *const libc::c_void, data.len());
        }
    });

    start_timer(Duration::from_secs(3));

    submit(ring_fd, &sq, 0)?;

    println!("starting work...");
    loop {
        do_work();

        if let Some(cqe) = peek_completion(&cq)? {
            let n = cqe.res as usize;
            println!("completed: {}", String::from_utf8_lossy(&buf[..n]));
        }
    }
}

start_timer()やdo_work()、WORK_COUNTは何も変わっていない。変わったのは、read()ではなくsubmit()とpeek_completion()を使っていること。
実装の詳細は先程のリンク先を見て欲しいが、submit()では、 io_uring を使ってデータの読み込み処理を依頼している。
そしてpeek_completion()は依頼した処理が終わっているかを確認し、終わっている場合はその結果を返す。

タイマーをスタートさせたあとsubmit()で処理を依頼し、その後loopに入る。loopではdo_work()を終える度にpeek_completion()を呼び出し、 I/O 処理の結果が返ってきたらそれを表示する。
つまり「依頼」「実行」「回収」を分離している。スレッドはまず「依頼」だけを行い、カーネルが「実行」している間、自分は他の仕事を行う。そして自分が必要になったタイミングで結果の「回収」を行う。

このプログラムを実行するとすぐに以下が表示される。

starting work...

その1秒後に以下が表示される。

completed: hello from pipe

さらにその2秒後に以下が表示される。

time is up: work done = 29

read()と異なり、submit()はすぐに終わる。そのためすぐにloopに入りdo_work()の実行に取り掛かる。
その結果、do_work()を29回実行できている。つまり、タイマーの3秒のうちほとんどの時間をdo_work()の実行に充てることができている。読み取り処理には1秒かかるわけだが、それはカーネルが並行して行ってくれている。

もちろん I/O 処理の結果を出力することもできており、これはpeek_completion()によって実現できているが、peek_completion()のなかではシステムコールは発行していない。
io_uring ではスレッドとカーネルが同じメモリ領域を操作し、それを使って情報のやり取りを行う。システムコールを発行せずに「回収」を行えるため、 I/O 処理が頻発するようなワークロードではそれがパフォーマンス上の大きな優位性となることがある。

このように、 io_uring を上手く使えれば、より効率的に I/O 処理を扱えるようになる。

データ競合

しかし、並行性やメモリ共有といった特徴は効率性を生むだけでなく、データ競合という問題も生み出す。

カーネルに I/O 処理を依頼する際、依頼内容によっては、カーネルがデータを読み書きするためのメモリ領域を用意し、そのメモリアドレスを渡す必要がある。カーネルはそのメモリ領域にアクセスし、必要な処理を行う。
しかし、スレッドとカーネルが並行して処理を進めるため、今はカーネルが使うはずのそのメモリ領域に対してスレッドが操作を行ってしまうことが、仕組み上あり得る。読み書きをしたり、メモリ領域を解放してしまったり。これがデータ競合であり、バグの要因になる。

先程のコードを一部書き換え、わざとデータ競合を起こしたのが以下。

github.com

let mut buf = vec![0u8; 256];

let sqe = unsafe { sqes_base.add(0) };
unsafe {
    ptr::write_bytes(sqe, 0, 1);
    (*sqe).opcode = IORING_OP_READ;
    (*sqe).fd = read_fd;
    (*sqe).off = u64::MAX;
    (*sqe).addr = buf.as_mut_ptr() as u64;
    (*sqe).len = buf.len() as u32;
    (*sqe).user_data = 0x01;
}

// 中略

submit(ring_fd, &sq, 0)?;

let mut account_name = buf;
account_name[..5].copy_from_slice(b"alice");
println!(
    "account_name: {}",
    String::from_utf8_lossy(&account_name[..15])
);

println!("starting work...");
loop {
    do_work();

    if WORK_COUNT.load(Ordering::Relaxed) == 20 {
        println!(
            "account_name: {}",
            String::from_utf8_lossy(&account_name[..15])
        );
    }

    if let Some(cqe) = peek_completion(&cq)? {
        println!("received message: {} bytes", cqe.res);
    }
}

変数bufはカーネルが使うために用意した領域で、カーネルはIORING_OP_READ、つまり読み込み処理を行った結果をここに書き込む。
しかしこのコードを書いた開発者は、bufはもう使い終わったと勘違いし、別の用途で使い回してしまっている。それが以下のコード。

let mut account_name = buf;
account_name[..5].copy_from_slice(b"alice");

account_nameという名前に変え、その値をaliceという文字列にしている。

しかし名前を変えたところで、この変数が使っているメモリ領域が変わるわけではない。
カーネルはIORING_OP_READした結果をそこに書き込む。その結果、account_nameの値はhello from pipeになってしまう。開発者が意図したものとは違う値になってしまう。

カーネルはただ、sqeのaddrフィールドに書かれているメモリアドレスに対してIORING_OP_READした結果を書き込むだけであり、呼び出し元のスレッドにおいてそのメモリアドレスがどの変数と結びついているのか、そのメモリアドレスが解放されているのか、などは知りようがない。
そのため、 io_uring を利用する側が気を付けなければならない。カーネルによる操作が終わるまでは当該メモリ領域を操作してはいけない、という規律を守らなければならない。そうしないとデータ競合が発生してしまう。

Rust の所有権を活用してデータ競合を防止する

開発者がミスをせず規律を守り続ける限りデータ競合は発生しないが、それは現実的ではない。何らかの仕組みによって防止することが望ましい。
その方法のひとつが、 Rust の所有権を利用した仕組み。

簡易的な実装例を用意した。

github.com

この実装で、先程と同じようにbufを操作しようとすると、コンパイルエラーになる。もちろんbufへの不正な操作をやめればコンパイルが通り、動作する。

error[E0382]: use of moved value: `buf`
   --> src/main.rs:278:28
    |
252 |     let mut buf = vec![0u8; 256];
    |         ------- move occurs because `buf` has type `Vec<u8>`, which does not implement the `Copy` trait
...
276 |     submit(ring_fd, &sq, 0, buf, &mut in_flight)?;
    |                             --- value moved here
277 |
278 |     let mut account_name = buf;
    |                            ^^^ value used here after move
    |

変えたのはsubmit()とpeek_completion()。

submit()のインターフェースを変え、bufを渡すことを必須とした。
そうすることで、submit()の呼び出し以降は、main()でbufを触ることはできなくなる。bufの所有権がsubmit()に移るため。main()でbufを操作しようとすると、上記したようにコンパイルエラーになる。
しかし、ただ単にbufの所有権をsubmit()に移すだけだと、submit()の呼び出しが終了したタイミングでbufは解放されてしまう。bufはこれからカーネルが使うのだから、それでは困る。
そこで、in_flightという変数を用意し、それにbufを持たせる。in_flightはmain()で宣言しsubmit()に参照を渡すだけなので、in_flightの所有権はmain()が持ち続ける。そしてそのin_flightがbufを持つ。そうすることで、submit()のスコープが終了してもbufは解放されない。

fn submit(
    ring_fd: i32,
    sq: &SqRing,
    sqe_index: u32,
    buf: Vec<u8>,
    in_flight: &mut Option<Vec<u8>>,
) -> io::Result<()> {
    let tail = unsafe { (*sq.tail).load(Ordering::Relaxed) };
    unsafe {
        ptr::write(sq.array.add((tail & sq.ring_mask) as usize), sqe_index);
        (*sq.tail).store(tail.wrapping_add(1), Ordering::Release);
    }

    let ret = unsafe { io_uring_enter(ring_fd, 1, 0, 0) };
    if ret < 0 {
        return Err(io::Error::last_os_error());
    }

    // in_flight に buf を持たせる
    *in_flight = Some(buf);

    Ok(())
}
fn main() -> io::Result<()> {
    // 省略

    let mut buf = vec![0u8; 256];

    let sqe = unsafe { sqes_base.add(0) };
    unsafe {
        ptr::write_bytes(sqe, 0, 1);
        (*sqe).opcode = IORING_OP_READ;
        (*sqe).fd = read_fd;
        (*sqe).off = u64::MAX;
        (*sqe).addr = buf.as_mut_ptr() as u64;
        (*sqe).len = buf.len() as u32;
        (*sqe).user_data = 0x01;
    }

    // 中略

    let mut in_flight: Option<Vec<u8>> = None;

    // in_flight は参照を渡しているだけなので、 in_flight の所有権は引き続き main() が持っている
    submit(ring_fd, &sq, 0, buf, &mut in_flight)?;

    // 省略
}

peek_completion()のインターフェースも変えている。引数としてin_flightの参照を受け取る。
そして、カーネルから結果が返ってきておりそれを取り出す時に、結果(cqe)だけでなくin_flightが持っていたbufも返すようにする。そうすることで、bufの所有権がmain()に戻ってくる。これ以降はmain()がbufを触ってもコンパイルエラーにならない。

fn peek_completion(
    cq: &CqRing,
    in_flight: &mut Option<Vec<u8>>,
) -> io::Result<Option<(IoUringCqe, Vec<u8>)>> {
    let head = unsafe { (*cq.head).load(Ordering::Relaxed) };
    let tail = unsafe { (*cq.tail).load(Ordering::Acquire) };

    if head == tail {
        return Ok(None);
    }

    let cqe = unsafe { ptr::read(cq.cqes.add((head & cq.ring_mask) as usize)) };
    unsafe { (*cq.head).store(head.wrapping_add(1), Ordering::Release) };

    if cqe.res < 0 {
        return Err(io::Error::from_raw_os_error(-cqe.res));
    }

    let buf = in_flight
        .take()
        .expect("completion without in-flight buffer");

    Ok(Some((cqe, buf)))
}

「submit()したら CQE を回収するまでbufを他の用途で使ってはいけない」という、規律、制約をコードで表現できている。型によって保証されている。そのため、開発者の注意力に頼らずにデータ競合を防げる。

この「防止策」はあくまでも、先に示したデータ競合を防ぐためのものでしかない。あのデータ競合をコンパイルエラーにできる、という話でしかなく、実用性は乏しい。
しかし、 Rust の所有権を使った考え方やアプローチ自体には汎用性があるため、例として示した。