トークンバケット×Semaphoreによる二段のレート制限と、accept()より前には効かないという限界

今回はレート制限です。目的の異なる2種類の制限を、listener.accept()の直後に重ねて実装します。1つは送信元IPごとの流量制限(トークンバケット)、もう1つは全体の同時接続数制限(tokio::sync::Semaphore)です。

コネクション単位: Semaphoreで同時接続数に上限をかける

// src/rate_limit.rs
use std::sync::Arc;
use tokio::sync::{OwnedSemaphorePermit, Semaphore};

pub struct ConnectionLimiter {
    semaphore: Arc<Semaphore>,
}

impl ConnectionLimiter {
    pub fn new(max_connections: usize) -> Self {
        Self {
            semaphore: Arc::new(Semaphore::new(max_connections)),
        }
    }

    pub fn try_acquire(&self) -> Option<OwnedSemaphorePermit> {
        Arc::clone(&self.semaphore).try_acquire_owned().ok()
    }
}

Semaphoreは「許可証(permit)をmax_connections個しか発行しない」だけのシンプルな仕組みです。try_acquire_owned()は許可証が取れなければ待たずにErrを返すので、上限に達した接続はその場で拒否します。取得したOwnedSemaphorePermitはコネクション処理タスクの中にmoveで持ち込み、タスクが終わってdropされたタイミングで自動的に許可証が返却されます。

送信元IP単位: トークンバケットで平均レートを制限する

use std::collections::HashMap;
use std::net::IpAddr;
use std::sync::Mutex;
use std::time::Instant;

struct TokenBucket {
    tokens: f64,
    last_refill: Instant,
}

pub struct IpRateLimiter {
    capacity: f64,
    refill_per_sec: f64,
    buckets: Mutex<HashMap<IpAddr, TokenBucket>>,
}

impl IpRateLimiter {
    pub fn new(capacity: f64, refill_per_sec: f64) -> Self {
        Self {
            capacity,
            refill_per_sec,
            buckets: Mutex::new(HashMap::new()),
        }
    }

    pub fn allow(&self, ip: IpAddr) -> bool {
        let mut buckets = self.buckets.lock().unwrap();
        let now = Instant::now();
        let bucket = buckets.entry(ip).or_insert_with(|| TokenBucket {
            tokens: self.capacity,
            last_refill: now,
        });

        let elapsed = now.duration_since(bucket.last_refill).as_secs_f64();
        bucket.tokens = (bucket.tokens + elapsed * self.refill_per_sec).min(self.capacity);
        bucket.last_refill = now;

        if bucket.tokens >= 1.0 {
            bucket.tokens -= 1.0;
            true
        } else {
            false
        }
    }
}

IPごとに「容量capacityまで貯まるトークンの入れ物」を持たせ、接続1回につき1トークン消費します。ポイントは、補充を専用のタイマー(定期的にカウンタをリセットするタスクなど)で行わず、アクセスが来た瞬間に前回アクセスからの経過時間を計算して補充量を求めていることです。固定時間窓でカウンタをリセットする方式(fixed window)は、窓の境界をまたいで倍のリクエストが素通りする問題を抱えますが、トークンバケットは「直近capacity件分のバーストは許容しつつ、長期的な平均レートはrefill_per_secに収束する」という性質を持ち、この境界問題が起きません。

判定順序: 安いチェックを先に

// src/main.rs(抜粋)
accepted = listener.accept() => {
    let (inbound, peer_addr) = accepted?;
    if !ip_limiter.allow(peer_addr.ip()) {
        eprintln!("rejecting {peer_addr}: ip rate limit exceeded");
        continue;
    }
    let Some(permit) = limiter.try_acquire() else {
        eprintln!("rejecting {peer_addr}: connection limit reached");
        continue;
    };
    // ここから先は既存のspawn処理

IPレート制限を先に、コネクション数上限を後にチェックしています。前者はメモリ上のHashMapロックだけで完結する安い処理、後者は共有のSemaphoreから許可証を1つ消費する処理です。同じIPから制限超過の接続が来るたびに毎回Semaphoreの許可証を消費・解放していては、そのIPだけで全体の枠を圧迫しかねません。先に安く弾ける条件で弾いておくことで、無関係な全体枠を守っています。

この方式が効くのは、TCPハンドシェイクが終わった後

ここまでの実装はlistener.accept()が返ってきた、つまりOSカーネルが3-wayハンドシェイクを完了させ、確立済みキュー(accept queue)から接続を取り出した後で動きます。裏を返すと、SYNパケットを大量に送りつけて確立前のキューを溢れさせるような攻撃(SYNフラッド)には、この実装は一切関与できません。カーネルはaccept()が呼ばれるかどうかに関わらずハンドシェイクへの応答自体は行うため、アプリケーション層でのレート制限が発動する頃には、その接続のハンドシェイクコストはすでに支払い終わっています。

この種の攻撃への対策は本来もっと下のレイヤー(SYN cookie、iptables/nftablesのレート制限、あるいはL4ロードバランサやeBPFでのフィルタリング)が担う領域であり、アプリケーションプロセス内のレート制限はあくまで「ハンドシェイクを乗り越えてきた後の、行儀の悪いクライアント」を制御するためのものだと切り分けて考える必要があります。

なおIpRateLimiter内のHashMapは、一度アクセスしたIPのエントリを明示的に削除する仕組みを持たないため、稼働し続ける限り単調に増加します。今回は学習目的のスコープとして手を付けていませんが、本番運用であれば定期的な掃除(あるいはLRU的な追い出し)が必要になる箇所です。