Skip to content
fandhe-backend
GitHub

レスポンスストリーミングガイド

fandhe-backend は、レスポンス body を一括ではなく逐次送信する chunked ストリーミング送信を opt-in で提供する。SSE (text/event-stream)・大きなファイルの逐次生成・進捗つき長時間処理の応答などに 使う。入口は既定ハンドラ trait の opt-in 既定メソッド Handler::handle_streamingcrates/core/src/server.rs)と、 fandhe_backend_core::streaming モジュールの StreamingResponse / BodyWritercrates/core/src/streaming.rs)である。

既存の Handler::handle のみを実装した型は本機能を一切意識せずに動作し続ける (handle_streaming の既定実装は常に None を返し、従来どおり Content-Length 一括応答になる。後方互換)。

基本形: handle_streaming + producer タスク

handle_streamingSome(StreamingResponse) を返すと、コアの書き出しループは 通常の一括応答経路の代わりに chunked framing で逐次送信する。典型パターンは StreamingResponse::channel で得た BodyWritertokio::spawn した producer タスクへ move し、producer がデータ生成の都合に合わせて send / finish を呼ぶ 構成である(完全な実行例は crates/core/src/server.rsHandler::handle_streaming doc test を正とする)。

use fandhe_backend_core::server::Handler;
use fandhe_backend_core::streaming::StreamingResponse;
use fandhe_backend_http::request::RequestHead;
use fandhe_backend_http::response::Response;

struct StreamingHandler;

impl Handler for StreamingHandler {
    fn handle(&self, _head: &RequestHead, _body: &[u8]) -> fandhe_backend_routes::HandlerFuture {
        Box::pin(async { Response::empty(404) })
    }

    fn handle_streaming(&self, _head: &RequestHead, _body: &[u8]) -> Option<StreamingResponse> {
        // status / Content-Type / チャネル容量(bounded mpsc)を指定して構築する
        let (response, writer) = StreamingResponse::channel(200, Some("text/plain"), 4);
        tokio::spawn(async move {
            writer.send(b"hello ".to_vec()).await.ok();
            writer.send(b"world".to_vec()).await.ok();
            writer.finish().await.ok();
        });
        Some(response)
    }
}

Content-Type が不要なら既定容量(8)の StreamingResponse::new(status) も使える。

チャネル構築とバックプレッシャ

StreamingResponse::channel(status, content_type, capacity)capacity は bounded mpsc チャネルの容量である。

項目挙動
チャネル満杯時の send受信側(コアの書き出しループ)がソケットへ書き出して空きができるまで .await で待機する(バックプレッシャ)
サーバ側バッファ上限高々「capacity × 1 チャンク分」に有界。producer がソケットの処理速度を追い越して無制限にメモリを積み上げることはない(DoS 対策)
capacity = 01 に切り上げる(panic しないフェイルセーフ)
既定容量(new8(数チャンク分のパイプライン化を許しつつ上限を小さく保つ妥協値)
BodyWriter の clone可。mpsc::Sender のセマンティクスを継承し、複数タスクから送信できる

send / finish の契約

操作契約
send(data)1 チャンク分を送る。空 Vec も成功するが、ワイヤへは無出力(誤終端防止)
finish()正常終端。受信側は終端チャンク(0\r\n\r\n)を送出し、応答を完全なものとして扱う(keep-alive 継続も許される)。self を消費するため「finish 後の send」は型レベルで書けない
finish を呼ばずに drop打ち切り。受信側は終端チャンクを送出せず接続をクローズする
戻り値 Err(StreamClosed)受信側が既に終了(クライアント切断等)した後の送信。producer は以降の送信を止めてタスクを終了してよい

finish 省略時に終端チャンクを送らないのは fail-closed の応答完全性維持である。 打ち切られた応答をクライアント・キャッシュに「完全な応答」と誤認させない (RFC 9112 の length 整合性。crates/core/src/streaming.rs のモジュール doc を参照)。

チャンク間隔の制約(30 秒以内)

producer からの次チャンク待ちには書き込みタイムアウト(30 秒、固定値)が適用され、 超過すると正常に稼働している producer でも接続が強制クローズされる (スロープロデューサ対策)。SSE のハートビートや long-poll のようにイベント発生が まばらな場合は、30 秒未満の間隔で writer.send(Vec::new()) を呼んで待ち時間を リセットするとよい(空チャンクはワイヤへ無出力のため、クライアントに余計な バイトを見せずに内部キープアライブとして使える)。

HTTP バージョン別の挙動

クライアントframing接続
HTTP/1.1Transfer-Encoding: chunked ヘッド + チャンク列 + 終端チャンクfinish で正常終端すれば keep-alive 継続可
HTTP/1.0framing なしの生データを EOF(接続クローズ)で終端。常に Connection: close応答完了 = 接続クローズ(keep-alive は常に無効)

HTTP/1.0 は Transfer-Encoding: chunked を理解しない前提のクライアントが残るため、 コアが自動でフォールバックする。ハンドラ側の実装を分ける必要はない。

また RFC 9112 §6.3 に従い、1xx・204・304 を handle_streaming から返した場合は body 送出ループへ入らず、ヘッド送出のみで応答を完了させる。

通常 handle との使い分け

観点handle(一括応答)handle_streaming
body応答全体を組み立ててから Content-Length 付きで送信チャンク単位で逐次送信(全体サイズ不要)
契約async(HandlerFuture を返す)同期のまま(チャネルを組み立てて即座に返り、非同期 I/O は producer タスクが担う)
適する用途通常の API 応答・小さな bodySSE・大きな逐次生成 body・長時間処理の進捗通知
レスポンス後処理型プラグイン(CORS ヘッダ付与・gzip 圧縮)適用される(finalize_responseCORS ヘッダ付与は適用される(finalize_streaming_head)。gzip 圧縮は既定では適用されず、CompressionConfigBuilder::compress_streaming(true) の明示 opt-in 時のみチャンク単位で適用される(prepare_streaming_compression、HTTP/1.1 chunked 経路限定)

評価順序にも注意する。パスインターセプト型プラグイン(graphql / openapi / static 等)が処理を完結させなかった場合にのみ handle_streaming が確認され、 Some ならこのリクエストに対して handle は呼ばれない。None なら従来どおり handle の一括応答経路に入る。

セキュリティ・制約

  • bounded チャネルのバックプレッシャにより、producer 起因のメモリ積み上げは 「capacity × 1 チャンク分」に有界(.claude/rules/security.md のリソース枯渇 DoS 対策)
  • チャンク待ち・実書き込みの双方に 30 秒のタイムアウトと接続生存期間上限 (Server::max_connection_lifetime)の短い方が適用され、超過時は接続を 強制クローズする(フェイルクローズ)
  • タイムアウト・書き込みエラー・producer 打ち切りの場合、Middleware::on_response は呼ばれない(「完走した応答のみ観測する」契約)
  • CORS ヘッダ付与(Server::cors 登録時)はストリーミング応答にも適用される。 gzip 圧縮は既定では適用されず、compress_streaming(true) を明示 opt-in した場合のみチャンク単位で適用される(BREACH 類似リスクへの配慮から 既定 OFF。SSE 等で秘密情報と攻撃者制御入力を混在させる場合は有効化しない、 または対象外の Content-Type に構成する。詳細は fandhe_backend_plugin_compression の crate doc・docs/design/plugin-boundary.md 5.10.6 節を参照)。HTTP/1.0 経路(EOF 終端)はストリーミング圧縮の対象外で常に identity のまま送出する
  • 通常応答(Handler::handle)の gzip 圧縮は body 長がしきい値 (CompressionConfigBuilder::blocking_threshold、既定 64 KiB)以上の場合 spawn_blocking へオフロードされるが、ストリーミング応答のチャンク 単位圧縮(compress_streaming)は対象外で常に接続タスク上で実行 される。1 チャンクを過大にする実装(例: 数百 KiB 以上を 1 回の BodyWriter::send にまとめて送出する)は、そのチャンクの gzip 圧縮が 接続タスクの tokio ワーカスレッドを長時間占有しうる。逐次配信の意味論を 保つため意図的にチャンク単位で圧縮する設計(crates/plugin-compression crate doc「チャンク単位のストリーミング gzip 圧縮」節)である以上、 大きな単一チャンクを送出する必要がある場合は通常応答(apply_compression 経由、オフロード対応済み)を使うことを推奨する(実測根拠・設計判断は docs/design/plugin-boundary.md の「巨大応答の gzip 圧縮を spawn_blocking へ切り離す」節を参照)

関連ドキュメント

  • 拡張点の全体像(Handler と 4 拡張点の関係): extension-points.md
  • graceful shutdown との組み合わせ(in-flight 接続の完了待ち): graceful-shutdown.md
  • API・契約の正とする doc comment: crates/core/src/streaming.rscrates/core/src/server.rsHandler::handle_streaming
  • sans-IO chunked エンコーダ: crates/http/src/chunked.rsencode_chunk / encode_terminator