Server Sent Events (SSE) を受信する

SSE を EventSource、http プラグインの fetch(認証ヘッダー付き)、Rust の reqwest と Channel の 3 通りで受ける。チャンクをまたぐ行の解析と Last-Event-ID での再接続も示す。

通信 対象: Tauri 2.x 更新日: 読了目安: 約13分 net-014
目次
  1. 前提条件
  2. 1. フロントエンドから実装する (TypeScript)
  3. EventSource で受ける
  4. 認証ヘッダーが要るなら http プラグインの fetch で読む
  5. 2. バックエンドから実装する (Rust)
  6. 動作確認
  7. よくあるエラーと対処法
  8. OS ごとの違いと注意点
  9. 関連レシピ

通知、ジョブの進捗、AI の応答を少しずつ表示するなど、サーバーからアプリへ一方向に届け続ける場面では SSE が手軽です。HTTP のレスポンスを開いたままにして、サーバーが data: で始まる行を書き足していく仕組みなので、普通の HTTP サーバーやプロキシをそのまま通ります。Tauri では、標準の EventSource、http プラグインの fetch、Rust の reqwest の 3 通りで受けられます。アプリからも頻繁に送るなら WebSocket を使います。

前提条件

方法認証ヘッダーCORS の許可自動再接続必要なもの
EventSource付けられない必要ありなし
http プラグインの fetch付けられる不要自分で書くプラグインと URL の許可
Rust(reqwest)付けられる不要自分で書くクレートの追加

EventSource はプラグインも権限も要りませんが、サーバーがアプリの Origin に CORS を許可している必要があります(Origin の値は「OS ごとの違い」を参照)。http プラグインは npm run tauri add http で追加し、接続先の URL を許可します。許可の書き方と fetch / reqwest の使い分けは HTTP GET リクエストを送る(Rust経由) で説明しているので、ここでは例だけ示します。

{
  "$schema": "../gen/schemas/desktop-schema.json",
  "identifier": "default",
  "description": "Capability for the main window",
  "windows": ["main"],
  "permissions": [
    "core:default",
    {
      "identifier": "http:default",
      "allow": [{ "url": "https://stream.wikimedia.org/*" }]
    }
  ]
}

Rust で受ける場合は、次のクレートを追加します。

cd src-tauri
cargo add reqwest@0.12 --features stream
cargo add futures
cargo add tokio --features time

1. フロントエンドから実装する (TypeScript)

EventSource で受ける

export function listenWithEventSource(url: string) {
  const es = new EventSource(url); // ヘッダーは付けられない。認証は Cookie かクエリで

  es.addEventListener('message', (e: MessageEvent<string>) => {
    const d = JSON.parse(e.data) as { server_name: string; title: string };
    console.log(d.server_name, d.title, e.lastEventId);
  });
  // "event: progress" の行が付いたイベントは、その名前で受ける
  es.addEventListener('progress', (e: MessageEvent<string>) => console.log('progress', e.data));
  es.addEventListener('error', () => {
    // CONNECTING なら自動で再接続中。CLOSED なら再接続をあきらめた
    console.warn(es.readyState === EventSource.CLOSED ? '停止しました' : '再接続中…');
  });
  return () => es.close();
}

接続が切れると error が発生して自動で再接続し、そのとき最後に受けた id: を Last-Event-ID ヘッダーで送るので、サーバーが対応していれば続きから受け取れます。ただし HTTP エラーや、Content-Type が text/event-stream でない応答では再接続せず、CLOSED のまま止まります。

認証ヘッダーが要るなら http プラグインの fetch で読む

レスポンスの本文を少しずつ読み、SSE の形式を自分で解釈します。1 回の読み取りで届く塊(チャンク)は行やイベントの途中で切れるので、チャンクごとに split('\n') するだけでは、長い JSON がときどき JSON.parse で失敗します。未完成の行を次に持ち越す解析器を用意します。

import { fetch } from '@tauri-apps/plugin-http';

export type SseMessage = { event: string; data: string; id: string };

export class SseParser {
  lastEventId = '';
  retryMs: number | null = null;
  private buf = '';
  private data: string[] = [];
  private event = '';

  constructor(private readonly onMessage: (m: SseMessage) => void) {}

  push(text: string) {
    this.buf += text;
    // 末尾の \r は、次のチャンクの \n と組になるかもしれないので持ち越す
    const hold = this.buf.endsWith('\r');
    const lines = (hold ? this.buf.slice(0, -1) : this.buf).split(/\r\n|\r|\n/);
    this.buf = (lines.pop() ?? '') + (hold ? '\r' : ''); // 最後の行は未完成かもしれない
    for (const line of lines) this.line(line);
  }

  private line(line: string) {
    if (line === '') { // 空行で 1 イベントが完成する
      if (this.data.length) {
        this.onMessage({ event: this.event || 'message', data: this.data.join('\n'), id: this.lastEventId });
      }
      this.data = [];
      this.event = '';
      return;
    }
    if (line.startsWith(':')) return; // コメント行(サーバーの生存確認に使われる)
    const i = line.indexOf(':');
    const field = i < 0 ? line : line.slice(0, i);
    const raw = i < 0 ? '' : line.slice(i + 1);
    const value = raw.startsWith(' ') ? raw.slice(1) : raw;
    if (field === 'data') this.data.push(value);
    else if (field === 'event') this.event = value;
    else if (field === 'id') this.lastEventId = value;
    else if (field === 'retry' && /^\d+$/.test(value)) this.retryMs = Number(value);
  }
}

// 切れたら retry(既定 3 秒)待ち、Last-Event-ID を付けて繋ぎ直す
export async function subscribeSse(
  url: string,
  onMessage: (m: SseMessage) => void,
  signal: AbortSignal,
  token?: string,
) {
  const parser = new SseParser(onMessage);
  while (!signal.aborted) {
    const headers: Record<string, string> = { Accept: 'text/event-stream' };
    if (token) headers.Authorization = `Bearer ${token}`;
    if (parser.lastEventId) headers['Last-Event-ID'] = parser.lastEventId;

    // この接続だけを切るための AbortController。60 秒何も届かなければ切る
    const conn = new AbortController();
    const cut = () => conn.abort();
    signal.addEventListener('abort', cut);
    let watchdog = window.setTimeout(cut, 60_000);
    try {
      const res = await fetch(url, { headers, signal: conn.signal, connectTimeout: 10_000 }).catch(() => null);
      if (res?.status === 204) return; // 「もう繋がないで」という合図
      if (res && !res.ok) throw new Error(`HTTP ${res.status}`); // 認証エラーなどは繰り返さない
      if (res?.body) {
        const reader = res.body.getReader();
        const decoder = new TextDecoder();
        try {
          for (;;) {
            const { done, value } = await reader.read();
            if (done) break;
            window.clearTimeout(watchdog);
            watchdog = window.setTimeout(cut, 60_000);
            parser.push(decoder.decode(value, { stream: true }));
          }
        } catch {
          // 通信が切れた、見張りで切った、または中断された
        }
      }
    } finally {
      window.clearTimeout(watchdog);
      signal.removeEventListener('abort', cut);
    }
    if (signal.aborted) return;
    await new Promise((r) => window.setTimeout(r, parser.retryMs ?? 3000));
  }
}
// (続き) 使い方: 10 秒受けて止める
const controller = new AbortController();
subscribeSse('https://stream.wikimedia.org/v2/stream/recentchange', (m) => {
  const d = JSON.parse(m.data) as { server_name: string; title: string };
  console.log(m.event, d.server_name, d.title);
}, controller.signal).catch((e) => console.error('SSE 停止:', e));
window.setTimeout(() => controller.abort(), 10_000);

abort() は本文の読み取り中でも効きます。connectTimeout は接続までの制限で、http プラグインの fetch には読み取りのタイムアウトがありません。そこで 60 秒何も届かなければ、その接続だけを切って繋ぎ直す見張りを入れています(サーバーがそれより短い間隔でコメント行などを送る前提です)。

2. バックエンドから実装する (Rust)

受けた内容を Rust で保存・加工したい、SSE を何本も同時に開きたい(後述の同時接続数の制限を受けない)、といった場合は Rust で受けます。フロントエンドへは Channel で渡すと、購読ごとに順番どおり届きます。

use std::collections::HashMap;
use std::sync::Mutex;
use std::time::Duration;

use futures::StreamExt;
use serde::Serialize;
use tauri::ipc::Channel;
use tauri::State;

#[derive(Clone, Serialize)]
struct SseMessage {
    event: String,
    data: String,
    id: String,
}

/// SSE の行を解釈する(改行は \n と \r\n に対応)
#[derive(Default)]
struct SseDecoder {
    buf: Vec<u8>,
    data: Vec<String>,
    event: String,
    last_id: String,
    retry_ms: Option<u64>,
}

impl SseDecoder {
    fn push(&mut self, chunk: &[u8]) -> Vec<SseMessage> {
        self.buf.extend_from_slice(chunk);
        let mut out = Vec::new();
        while let Some(pos) = self.buf.iter().position(|&b| b == b'\n') {
            // 改行の位置で切るので、日本語などの文字の途中で分かれない
            let raw: Vec<u8> = self.buf.drain(..=pos).collect();
            let text = String::from_utf8_lossy(&raw);
            let line = text.trim_end_matches(['\n', '\r']);
            if line.is_empty() {
                if !self.data.is_empty() {
                    let event = if self.event.is_empty() { "message".to_string() } else { self.event.clone() };
                    out.push(SseMessage { event, data: self.data.join("\n"), id: self.last_id.clone() });
                }
                self.data.clear();
                self.event.clear();
                continue;
            }
            if line.starts_with(':') {
                continue; // コメント行
            }
            let (field, value) = line.split_once(':').unwrap_or((line, ""));
            let value = value.strip_prefix(' ').unwrap_or(value);
            match field {
                "data" => self.data.push(value.to_string()),
                "event" => self.event = value.to_string(),
                "id" => self.last_id = value.to_string(),
                "retry" => self.retry_ms = value.parse().ok(),
                _ => {}
            }
        }
        out
    }
}

/// 購読中のタスク(止めるときに使う)
#[derive(Default)]
struct SseTasks(Mutex<HashMap<u32, tauri::async_runtime::JoinHandle<()>>>);

#[tauri::command]
fn sse_subscribe(
    url: String,
    token: Option<String>,
    on_event: Channel<SseMessage>,
    tasks: State<'_, SseTasks>,
) -> Result<u32, String> {
    let client = reqwest::Client::builder()
        .user_agent("my-tauri-app/1.0")
        .connect_timeout(Duration::from_secs(10))
        .read_timeout(Duration::from_secs(60)) // 60 秒何も届かなければ切れたとみなす
        .build()
        .map_err(|e| e.to_string())?;
    let id = on_event.id();
    let handle = tauri::async_runtime::spawn(async move {
        let mut decoder = SseDecoder::default();
        let mut retry = Duration::from_secs(3);
        loop {
            let mut req = client.get(&url).header("Accept", "text/event-stream");
            if let Some(t) = &token {
                req = req.bearer_auth(t);
            }
            if !decoder.last_id.is_empty() {
                req = req.header("Last-Event-ID", decoder.last_id.clone()); // 続きから再開
            }
            match req.send().await {
                Ok(res) if res.status().is_success() => {
                    let mut body = res.bytes_stream();
                    while let Some(Ok(chunk)) = body.next().await {
                        for msg in decoder.push(&chunk) {
                            let _ = on_event.send(msg);
                        }
                    }
                }
                Ok(res) => {
                    eprintln!("SSE: HTTP {}", res.status()); // 認証エラーなどは繰り返さない
                    return;
                }
                Err(e) => eprintln!("SSE: {e}"),
            }
            if let Some(ms) = decoder.retry_ms {
                retry = Duration::from_millis(ms);
            }
            tokio::time::sleep(retry).await;
        }
    });
    if let Some(old) = tasks.0.lock().unwrap().insert(id, handle) {
        old.abort(); // 同じ id の古い購読が残っていれば止める
    }
    Ok(id)
}

#[tauri::command]
fn sse_unsubscribe(id: u32, tasks: State<'_, SseTasks>) {
    if let Some(handle) = tasks.0.lock().unwrap().remove(&id) {
        handle.abort(); // タスクを止め、接続も閉じる
    }
}

#[cfg_attr(mobile, tauri::mobile_entry_point)]
pub fn run() {
    tauri::Builder::default()
        .plugin(tauri_plugin_http::init()) // 1 章の fetch を使う場合
        .manage(SseTasks::default())
        .invoke_handler(tauri::generate_handler![sse_subscribe, sse_unsubscribe])
        .run(tauri::generate_context!())
        .expect("error while running tauri application");
}

read_timeout() は 1 回の読み取りごとの制限で、サーバーが定期的に送るコメント行が途絶えたことを検知できます。timeout() はレスポンス全体の制限なので、SSE に付けると正常な接続まで切れます。

import { Channel, invoke } from '@tauri-apps/api/core';

type SseMessage = { event: string; data: string; id: string };

export async function subscribeViaRust(url: string, onMessage: (m: SseMessage) => void) {
  const onEvent = new Channel<SseMessage>();
  onEvent.onmessage = onMessage;
  const id = await invoke<number>('sse_subscribe', { url, token: null, onEvent });
  const stop = () => invoke<void>('sse_unsubscribe', { id });
  window.addEventListener('beforeunload', () => void stop()); // 再読み込みで受信を残さない
  return stop;
}

動作確認

Wikimedia が公開している編集履歴のストリーム https://stream.wikimedia.org/v2/stream/recentchange は CORS を許可しているので、3 通りすべてで試せます。npm run tauri dev で起動して上の使い方を実行すると、編集のたびに次のような行が 10 秒間流れます。

message en.wikipedia.org Hüttenbach (Altmühl)
message commons.wikimedia.org Category:Photographs by Hsu Hong Lin

10 秒たつと abort() で止まり、それ以上は流れません。Rust 版は subscribeViaRust() が返した関数を呼ぶまで流れ続けます。

よくあるエラーと対処法

  • EventSource がすぐ error になり CLOSED のまま: CORS の許可がない、HTTP エラー、Content-Type 違いのどれかです。サーバーが許可している Origin を確認し、ヘッダー認証が要るなら http プラグインか Rust に切り替えます。
  • 「url not allowed on the configured scope: …」: http プラグインの許可リストに URL がありません。前提条件の allow のパターンを見直します。
  • イベントがまとめて届く、切断時にしか届かない: サーバーやリバースプロキシが応答をためています。1 イベントごとに送り出すようサーバーを設定します(nginx なら X-Accel-Buffering: no を返す)。
  • 「state not managed for field tasks on command sse_subscribe. You must call .manage() before using this command」: .manage(SseTasks::default()) がありません。
  • 再読み込み後、コンソールに「[TAURI] Couldn't find callback id」で始まる警告が出続ける: Rust の受信が残っています。sse_unsubscribe を呼んでから離れます。

OS ごとの違いと注意点

  • Origin: EventSource のリクエストの Origin は、開発中は devUrl(Vite のテンプレートなら http://localhost:1420)、本番は Windows と Android で http://tauri.localhost、macOS と Linux で tauri://localhost です。開発中だけ動く場合は本番用の Origin の許可漏れを疑います。
  • Windows / Android の useHttpsScheme: true にすると Origin が https になり、EventSource から http:// の SSE には繋げなくなります(mixed content)。http プラグインと Rust は影響を受けません。
  • 同時接続数: HTTP/1.1 のサーバーでは、Webview から同じホストへの同時接続は 6 本までです。EventSource を何本も開くと、同じサーバーへの標準の fetch が待たされます。1 本にまとめるか、http プラグインか Rust で受けます。
  • スリープからの復帰: 接続が通知なく止まることがあります。回線の状態を確かめる方法は インターネット接続状態を監視する を参照してください。

関連レシピ

参考リンク(公式ドキュメント)

Web Ninja

この記事を書いた人

Web Ninja ウェブエンジニア (Web Engineer)

会社員ネットワークエンジニアから独立してかれこれ 25 年以上 Web エンジニアとして活動中。普段は JavaScript と Node.js を自在に操り、時には C++ や Perl といった古流の技も嗜みます。近年は Tauri × Rust という新たな武器を手に、デスクトップアプリ開発の最前線を駆け抜けています。「作りたい」を「作れる」に変えるための、実践的な「技」をお届けします。

お問い合わせ: tauri.ninja@gmail.com

内容の誤り・動かないコードを報告する