たとえば、クライアントが40秒オフラインになってから再接続し、その間にサーバーが300件のイベントを発行したとします。簡単な答えは2つあり、どちらも、あとから気づく形で失敗します。何も送らなければ、クライアントは黙って遅れたままになります。履歴を全部送れば、再接続のたびのコストが、ストリームの年齢に比例して増えます。

この記事は2つ目の答え(版0)から始め、問題を1つずつ直して、上限つきのログとスナップショットのフォールバック(版3)、そしてスナップショットを安全にする規則(版4)まで進みます。途中で、ブラウザの EventSource が実際に何をするかを確かめます。ブラウザはすでに、解決策の半分を持っているからです。

TL;DR

  • 再開には3つが要ります。クライアントが覚えておける位置(カーソル)、その位置から再生できるサーバー側のログ、そして位置が古すぎて再生できないときのフォールバック(スナップショット)です。
  • Server-Sent Events は、クライアント側の半分を無料で提供します。ブラウザは最後の id を覚え、再接続時に Last-Event-ID として送り返します。Chrome 154 で確認しました。提供されないのは、サーバー側の半分と、「あなたを再開できません」と伝える手段です。
  • non-200 のレスポンスを受けると、ネイティブの EventSource は永久に終了し、ページにはステータスも本文も見えません(Chrome 154 で確認)。だからブラウザ向けには、「古すぎる、最初からやり直して」をストリームの中で伝えます。200 で応答し、スナップショットを載せた reset イベントを送ります。自分で制御するクライアントなら、410 Gone でも構いません。
  • スナップショットは (状態, カーソル) の組です。カーソルを先に、状態をあとに読み、再生を冪等にします。逆の順序だと、更新が消えることがあります。僕の簡易モデルでは、1000回の試行のうち699回でした。
  • カーソルが安全なのは、シーケンス番号が順番どおりに見えるようになる場合だけです。書き手 A が 5 を取り、書き手 B が 6 を公開したあとに A がコミットすると、カーソルが 6 の読み手は 5 を二度と見ません。
  • この記事のエンドツーエンドのテスト(ランダムな接続切断と、ログの保持期間より長い不在)は、クライアントとサーバーが同じ状態になり、欠落ゼロ、重複ゼロで終わります。

版0:再接続して、全部読み直す

最も単純で正しい設計は、接続のたびに状態の全体を送り、そのあとライブのイベントを流すことです。カーソルもログもなく、間違えようがありません。コストは、再接続のたびに履歴の大きさに比例します。小さなデータなら問題ありませんが、大きなデータでは悲惨です。再接続は、ネットワークが悪いときにまとまって起きるからです。この版は手元に残しておいてください。以降のすべての版のフォールバックになります。

版1:カーソル

すべてのイベントに、サーバーが公開する順に1ずつ増えるシーケンス番号を付けます。クライアントは、最後に適用した番号を覚えておき、再接続したらその後のすべてを求めます:GET /events?after=42。サーバーはログから答えます。

これが動くかどうかは、3つの規則で決まります。

  1. シーケンスはストリームごとに単調で、欠番がない。 そうすれば、クライアントは番号を見るだけで、欠けたイベントを検出できます。
  2. カーソルはクライアントにとって不透明にする。 クライアントは保存して返すだけです。何を符号化しているかを後から変えられます。
  3. 番号が順番どおりに見えるようになる。 見落とされるのはここです。書き込みの開始時に番号を割り当て、コミット時に見えるようになるとします。書き手 A が 5 を取り、書き手 B が 6 を取って先にコミットすると、6 を見た読み手のカーソルは、5 を飛び越えてしまいます。
visibility_gap.mjs
// A cursor is only safe if sequence numbers become visible in order.
// Writer A takes number 5 and is slow to publish; writer B takes 6 and publishes first. Run: node visibility_gap.mjs
const published = [];                                   // what readers can see, in publication order
const publish = (seq, data) => published.push({ seq, data });
const readAfter = (cursor) => published.filter((e) => e.seq > cursor);

publish(4, "x");                                        // everything up to 4 is visible
// writer A allocated seq 5 but has not committed yet
publish(6, "y");                                        // writer B allocated 6 and committed first

let cursor = 4;
const first = readAfter(cursor);                        // the reader polls: sees only 6
cursor = Math.max(cursor, ...first.map((e) => e.seq));  // cursor advances to 6
publish(5, "z");                                        // writer A commits late
const second = readAfter(cursor);                       // the reader polls again with cursor 6

console.log("first read :", first.map((e) => e.seq));
console.log("second read:", second.map((e) => e.seq), "<- event 5 is never delivered");
first read : [ 6 ]
second read: [] <- event 5 is never delivered

このモデルは意図的に抽象的です。特定のデータベースではなく、バグの形を示しています。対策はどれも、見えるようになる順序を、割り当ての順序に合わせる方法です。公開の瞬間に番号を割り当てる(書き手を1つにする、またはロックを使う)、あるいは、読み手が「まだ処理中かもしれない最小の番号」で止まる、という方法です。

版2:カーソルをブラウザに持たせる

Server-Sent Events は、まさにこのクライアントの挙動を標準化しています。HTML Standard によると、id: フィールドは接続の最後のイベント ID を設定します。その値は、サーバーが別の値を設定するまで保持されます。再接続のとき、ブラウザはそれを Last-Event-ID リクエストヘッダーとして送ります。NUL 文字を含む id は無視されます。retry: は再接続までの待ち時間を設定します。EventSource のクライアントは、再接続ループつきのカーソル保存庫です。

これを鵜呑みにしたくなかったので、実際のブラウザで確かめました。ページが EventSource を開き、サーバーがイベント 1〜3 を送って接続を切り、サーバーは再接続のときのヘッダーを記録します。

eventsource_check.mjs
// What does a real browser EventSource do on reconnect, and on a non-200 response?
// Run: node eventsource_check.mjs   (needs a Chrome/Chromium binary; set CHROME=/path if needed)
import http from "node:http";
import { spawn } from "node:child_process";

const seenHeaders = [];   // Last-Event-ID values the server received on /sse
let attempts410 = 0;

const page = `<script>
const out = { messages: [], sseErrors: 0, gone: { errors: 0, readyState: null } };
const es = new EventSource("/sse");
es.onmessage = (e) => { out.messages.push(e.lastEventId); };
es.onerror = () => { out.sseErrors++; };
const gone = new EventSource("/gone");
gone.onerror = () => { out.gone.errors++; out.gone.readyState = gone.readyState; };
setTimeout(() => { fetch("/report", { method: "POST", body: JSON.stringify(out) }); }, 2500);
</script>`;

const server = http.createServer(async (req, res) => {
  if (req.url === "/") { res.writeHead(200, { "content-type": "text/html" }); return res.end(page); }
  if (req.url === "/sse") {
    const last = req.headers["last-event-id"]; seenHeaders.push(last ?? null);
    res.writeHead(200, { "content-type": "text/event-stream" });
    res.write("retry: 100\n\n");
    const start = last ? Number(last) + 1 : 1;
    for (let id = start; id < start + 3; id++) res.write(`id: ${id}\ndata: hello\n\n`);
    if (!last) setTimeout(() => res.destroy(), 100);          // first connection: drop it after 3 events
    return;
  }
  if (req.url === "/gone") { attempts410++; res.writeHead(410, { "content-type": "application/json" }); return res.end('{"reset":true}'); }
  if (req.url === "/report") {
    let b = ""; for await (const c of req) b += c;
    console.log("page reported:", b);
    console.log("Last-Event-ID headers seen by /sse:", JSON.stringify(seenHeaders));
    console.log("connection attempts to the 410 endpoint:", attempts410);
    res.end("ok"); clearTimeout(guard); chrome.kill(); server.close(); server.closeAllConnections();
  }
});
await new Promise((r) => server.listen(0, "127.0.0.1", r));
const chrome = spawn(process.env.CHROME ?? "google-chrome", [
  "--headless=new", "--no-sandbox", "--disable-gpu", "--user-data-dir=./es-check-profile",
  `http://127.0.0.1:${server.address().port}/`,
], { stdio: "ignore" });
const guard = setTimeout(() => { console.log("timeout"); chrome.kill(); process.exit(1); }, 20000);
page reported: {"messages":["1","2","3","4","5","6"],"sseErrors":1,"gone":{"errors":1,"readyState":2}}
Last-Event-ID headers seen by /sse: [null,"3"]
connection attempts to the 410 endpoint: 1

最初の接続にはヘッダーがなく、再接続には 3 が付いていました。そのあと、ページは 4、5、6 を受け取りました。再開プロトコルの半分が、クライアントのコードなしで動いています。(Linux 上のヘッドレス Chrome 154。ほかのブラウザは確認していません。)

出力の後半は、設計にとっていちばん重要です。/gone のエンドポイントは 410 を返しました。ブラウザは1回だけ試し、error を発火し、readyState は 2(CLOSED)になりました。再試行しませんでした。仕様もそう書いています。レスポンスのステータスが 200 でない、またはコンテンツタイプが text/event-stream でない場合は「接続を失敗させる」(fail the connection)。そして「ユーザーエージェントが接続を失敗させたあとは、再接続を試みない」(Once the user agent has failed the connection, it does not attempt to reconnect.)。ページには、ステータスコードも本文も見えません。

つまりブラウザでは、「あなたのカーソルは古すぎる」を 410 で伝えることはできません。200 のレスポンスの中、ストリームの中で伝える必要があります。

版3:上限つきのログと、ストリーム内のリセット

ログを無限に伸ばすことはできません。直近 N 件(または直近 N 分)だけを保持すると、それより長く不在だったクライアントは、ログから再開できません。サーバーはそれに気づいて、そう伝える必要があります。ネイティブの EventSource にも、fetch ベースのクライアントにも使える設計は、次のとおりです。

  • クライアントは、カーソルを送る(Last-Event-ID)。
  • cursor + 1 >= 保持している最古の番号 なら、カーソルより後のイベントを再生し、そのあとライブのログを追いかける。
  • そうでなければ、現在のスナップショットとそのカーソルを載せた event: reset を送り、そのカーソルから続ける。
  • カーソルのない接続は、「古すぎる」と同じ経路を通る。

EventSource と同じ規則(最後の ID を覚え、再接続で送る)に従うクライアントつきの、完全な実装です。テストで超過できるように、保持数は 50 件にしてあります。

resume.mjs:サーバー、クライアント、障害を起こすエンドツーエンドのテスト
// A resumable event stream in ~100 lines: cursor + bounded log + snapshot fallback.
// Requires Node 18+. Run: node resume.mjs
import http from "node:http";

const RETAIN = 50;                                   // the server keeps only the last 50 events
const log = [];                                      // [{ seq, key, value }]
const state = {};                                    // key -> { value, seq }   (the "current" truth)
let seq = 0;

const write = (key, value) => {                      // every state change appends to the log
  const ev = { seq: ++seq, key, value };
  state[key] = { value, seq };
  log.push(ev);
  if (log.length > RETAIN) log.shift();
  return ev;
};
const snapshot = () => ({ seq, state: structuredClone(state) });   // (cursor, state) read together: no await in between

const sse = (id, event, data) => `id: ${id}\nevent: ${event}\ndata: ${JSON.stringify(data)}\n\n`;

export const server = http.createServer((req, res) => {
  res.writeHead(200, { "content-type": "text/event-stream", "cache-control": "no-store" });
  const header = req.headers["last-event-id"];
  let cursor = header === undefined ? null : Number(header);

  const oldest = log.length ? log[0].seq : seq + 1;  // first sequence number still available
  if (cursor === null || !Number.isInteger(cursor) || cursor + 1 < oldest || cursor > seq) {
    const snap = snapshot();                          // cannot resume: tell the client in-band, then continue from the snapshot
    res.write(sse(snap.seq, "reset", snap));
    cursor = snap.seq;
  }
  for (const ev of log) if (ev.seq > cursor) { res.write(sse(ev.seq, "put", ev)); cursor = ev.seq; }

  const timer = setInterval(() => {                   // tail the log
    for (const ev of log) if (ev.seq > cursor) { res.write(sse(ev.seq, "put", ev)); cursor = ev.seq; }
  }, 5);
  req.on("close", () => clearInterval(timer));
});

// ---- a minimal client that behaves like EventSource: remember the last id, send it back on reconnect ----
export async function follow(url, onEvent, { signal }) {
  let lastId = null;
  while (!signal.aborted) {
    try {
      const res = await fetch(url, { headers: lastId === null ? {} : { "last-event-id": String(lastId) }, signal });
      let buf = "";
      for await (const chunk of res.body.pipeThrough(new TextDecoderStream())) {
        buf += chunk;
        for (let i; (i = buf.indexOf("\n\n")) >= 0; ) {
          const raw = buf.slice(0, i); buf = buf.slice(i + 2);
          const ev = Object.fromEntries(raw.split("\n").map((l) => [l.slice(0, l.indexOf(":")), l.slice(l.indexOf(":") + 2)]));
          lastId = Number(ev.id);
          onEvent(ev.event, JSON.parse(ev.data));
        }
      }
    } catch (e) { if (signal.aborted) return; }
    await new Promise((r) => setTimeout(r, 20));      // reconnect delay (EventSource's `retry`)
  }
}

if (import.meta.url === `file://${process.argv[1]}`) {
  await new Promise((r) => server.listen(0, "127.0.0.1", r));
  const url = `http://127.0.0.1:${server.address().port}/`;

  const mine = {}; let expected = null, gaps = 0, dups = 0, resets = 0, applied = 0;
  const ac = new AbortController();
  const done = follow(url, (type, d) => {
    if (type === "reset") { resets++; for (const k in mine) delete mine[k]; Object.assign(mine, d.state); expected = d.seq + 1; return; }
    if (expected !== null && d.seq < expected) { dups++; return; }
    if (expected !== null && d.seq > expected) gaps++;
    mine[d.key] = { value: d.value, seq: d.seq }; expected = d.seq + 1; applied++;
  }, { signal: ac.signal });

  // Producer: 400 writes, sometimes in bursts. Chaos: kill every open connection now and then,
  // and once keep the client away for long enough that the log moves past its cursor.
  for (let i = 0; i < 400; i++) {
    write(`k${i % 7}`, i);
    if (i % 60 === 59) server.closeAllConnections();                  // short outage: resume from the cursor
    if (i === 200) { server.closeAllConnections(); for (let j = 0; j < 120; j++) write(`k${j % 7}`, 1000 + j); } // long outage: > RETAIN events missed
    if (i % 5 === 0) await new Promise((r) => setTimeout(r, 3));
  }
  await new Promise((r) => setTimeout(r, 300));
  ac.abort(); await done; server.closeAllConnections(); server.close();

  const same = JSON.stringify(Object.keys(state).sort().map((k) => [k, state[k].value])) ===
               JSON.stringify(Object.keys(mine).sort().map((k) => [k, mine[k].value]));
  console.log(`server seq=${seq}  applied=${applied}  resets=${resets}  gaps=${gaps}  duplicates=${dups}  client state equals server state: ${same}`);
}

末尾のテストは、400 回の書き込みを行い、60 回ごとにすべての接続を切断し、さらに一度、クライアントを遠ざけたまま 120 件のイベントを書き込みます(保持数 50 を超えます)。

server seq=520  applied=354  resets=2  gaps=0  duplicates=0  client state equals server state: true

resets=2 は、最初の接続(カーソルなし)と、長い不在です。短い中断はすべて、カーソルから再開しました。gaps=0 と duplicates=0 は、クライアントがイベントを適用するときにシーケンス番号について行うチェックです。3回の実行で、最終状態はそのたびに一致しました(切断のタイミングが変わるので、applied は 349 から 354 の間でした)。

保持期間のチェックは、本当に効いているのでしょうか。サーバーから cursor + 1 < oldest の条件を取り除いて、もう一度実行しました。gaps=1 になりました。クライアントは、黙ってイベントを飛ばしたのです。その実行でも最終状態は一致しました。あとの書き込みが、失われた値をたまたま上書きしたからで、このバグがテストをすり抜けるのはそのためです。最終状態だけでなく、クライアントで欠落を数えてください。

reset イベントは、クライアントが持っているものを捨て、スナップショットを据える場所でもあります。

if (type === "reset") { for (const k in mine) delete mine[k]; Object.assign(mine, d.state); expected = d.seq + 1; return; }

版4:スナップショットは組である

スナップショットは「状態」ではなく、「カーソル C の時点の状態」で、その組であることが要点です。カーソルが状態を正しく表していなければ、カーソルから始める再生は、穴か繰り返しを残します。規則は、2つのものをロックなしで読むときと同じです。

カーソルを先に、状態をあとに読む。再生を冪等にする。

カーソルを先に読めば、状態はカーソル以上に新しいので、カーソルからの再生は、イベントを繰り返すことしかありません。再生が冪等(キーごとに、すでに持っているより新しくないものを無視して値を設定する)なら、繰り返しは無害です。状態を先に、カーソルをあとに読むと、カーソルが状態より新しくなりえて、その間のイベントは二度と再生されません。

snapshot_order.mjs
// A snapshot is a pair (state, cursor). Read them in the wrong order and an update can vanish.
// Run: node snapshot_order.mjs
const rng = (a) => () => { a = (a + 0x6d2b79f5) | 0; let t = Math.imul(a ^ (a >>> 15), 1 | a); t = (t + Math.imul(t ^ (t >>> 7), 61 | t)) ^ t; return ((t ^ (t >>> 14)) >>> 0) / 2 ** 32; };

function trial(order, rand) {
  const log = [], state = {}; let seq = 0;
  const write = () => { const key = "k" + Math.floor(rand() * 3); log.push({ seq: ++seq, key, value: seq }); state[key] = { value: seq, seq }; };
  const writes = (n) => { for (let i = 0; i < n; i++) write(); };
  writes(3);

  // Reading the state and reading the cursor are two steps. Writes can land between them (0 to 3 here).
  let snapState, cursor;
  if (order === "state, then cursor") { snapState = structuredClone(state); writes(Math.floor(rand() * 4)); cursor = seq; }
  else                                { cursor = seq; writes(Math.floor(rand() * 4)); snapState = structuredClone(state); }
  writes(Math.floor(rand() * 2));                                   // and a few more after the snapshot

  // The client installs the snapshot, then replays every log entry after the cursor,
  // skipping entries that are not newer than what it already holds for that key.
  const mine = structuredClone(snapState);
  for (const ev of log) if (ev.seq > cursor && ev.seq > (mine[ev.key]?.seq ?? 0)) mine[ev.key] = { value: ev.value, seq: ev.seq };
  return JSON.stringify(Object.entries(mine).sort()) === JSON.stringify(Object.entries(state).sort());
}

for (const order of ["state, then cursor", "cursor, then state"]) {
  const rand = rng(7); let ok = 0;
  for (let i = 0; i < 1000; i++) ok += trial(order, rand) ? 1 : 0;
  console.log(`${order.padEnd(20)} client matches server in ${ok}/1000 trials`);
}
state, then cursor   client matches server in 301/1000 trials
cursor, then state   client matches server in 1000/1000 trials

301 という数字は、僕の簡易な割り込み(2つの読み取りの間に 0〜3 回の書き込み、キーは3つ)に依存します。大事なのは、ゼロか、ゼロでないかの違いです。resume.mjs は別の方法でこの問題を避けています。snapshot() は、間に await を挟まず、同期的な1ステップで両方を読みます。これは、1つのデータベーストランザクションや、1つの一貫したスナップショットで読むことの、単一プロセスでの対応物です。

設計を選ぶ

検討した代替案です。

選択肢向く場面向かない場面
全部読み直す小さな状態。修復の経路大きな状態。頻繁な再接続
カーソル + 上限つきログ + スナップショット多数のクライアントが見る、変化する状態飛ばしたイベントが許されない厳密な監査
HTTP の Range リクエスト(RFC 9110 §14)バイト位置で指せるリソース:ファイル、追記専用の blob「位置」がバイトオフセットではないイベントストリーム
クライアントがすべてをメモリに保持する短いセッションアプリが再起動するもの全般
レベルトリガーの再取得(無効化の記事を参照)キーで安く読み直せる状態大きな履歴と、順序のあるイベント

再開の仕組みが壊れる場所

失敗症状防御
カーソルがログより古い黙った欠落サーバーで検出し、スナップショットつきの reset を送る(読めるクライアントには 410 でもよい)。
ネイティブの EventSource が non-200 を受ける永久に停止し、ステータスも見えない200 で応答し、ストリームの中で伝える。
スナップショットの状態をカーソルより先に読む更新が消えるカーソルを先に、状態をあとに。または一貫した1回の読み取り。冪等な再生。
シーケンス番号が順不同で見えるイベントを飛ばす見える順序を、割り当て順序に合わせる。
再生が冪等でない重なった部分で重複するキーごとにバージョンを持ち、新しくないものは無視する。
クライアントのテストが最終状態だけを見るあとのイベントが失われた値を上書きするため、バグが生き残るクライアントでシーケンスの欠落を数える。
ログが無制限メモリやディスクが増えるN 件または T 分を保持し、保持期間を契約の一部にする。
障害のあとの再接続ストームすべてのクライアントが同時にスナップショットを求める再接続の待ち時間にジッターを入れる。retry フィールドは全クライアントに同じ固定値を与えるので、自前のクライアントでジッターを足す。
クライアントがカーソルの中身に意味を持たせる符号化を変えられないカーソルは不透明にしておく。
NUL を含む id フィールドブラウザがその id を無視するので、カーソルが進まない単純な整数か、URL セーフなトークンにする。仕様は、値を NUL、LF、CR を含まない任意の UTF-8 文字列と説明している。

使う場面と使わない場面

カーソル + ログ + スナップショットを使うのは、多数のクライアントが変化する状態を見ていて、再接続が頻繁で、毎回状態全体を読み直すコストが無視できないときです。

次の場合は使いません。

  • 履歴が小さいとき。版0には、気にするほどの失敗モードがありません。
  • キーごとの最新の値だけが要るとき。ヒントを受けて再取得するほうが単純で、コードの経路が1つで済みます。
  • すべてのイベントが法的に意味を持つとき。欠落を隠してしまうスナップショットのフォールバックは、適切な道具ではありません。永続的なログを持ち、「古すぎる」をクライアントが必ず処理すべきエラーにします。

試して、壊してみる

# Node.js 18 以降(実行したのは Node.js 20)。eventsource_check.mjs 以外は依存関係なし。これは Chrome か Chromium が必要。
node resume.mjs             # ランダムな切断 + 保持期間より長い不在
node snapshot_order.mjs     # 2つの読み取りの順序が重要な理由
node visibility_gap.mjs     # シーケンス番号が順番どおりに見えるべき理由
node eventsource_check.mjs  # 実際のブラウザの挙動(必要なら CHROME=/path/to/chrome を設定)

そして、壊してみてください。RETAIN を 5 にする、あるいは cursor + 1 < oldest の条件を消して、最後だけでなく、切断のたびに「クライアントのキーの集合がサーバーと等しい」というアサーションを加えます。

3つの約束

再開できる性質は3つの約束でできていて、再開の仕組みを出荷する前に、3つすべてに答えられるようにしておくべきです。

  1. 位置。 クライアントは何を持ち運び、それは不透明になっているか。SSE が標準化しているのはここです。
  2. 保持。 ログは、その位置をどれだけの期間守るか。これは、あなたが選んで書き残す数字です。
  3. リセット。 位置が古すぎるとき、サーバーは何をするか。クライアントがそれを見逃せない形になっているか。多くの設計が欠いているのはここで、回復するストリームと、黙ってずれていくストリームの分かれ目です。

確認した範囲と、していない範囲

すべて Linux 6.12、Node.js 20.19、ヘッドレスの Google Chrome 154 で実行しました。ブラウザの挙動(再接続時の Last-Event-ID、non-200 のあとに再接続しないこと)は、その1つのブラウザで観察したもので、2026-10-04 時点の HTML Standard と一致しています。サーバーは、メモリ上のログを持つ単一のプロセスです。snapshot_order.mjs と visibility_gap.mjs はモデルで、データベースの実験ではありません。マルチノードのサーバー、永続的なログ、イベントストリームをバッファリングするプロキシ、ほかのブラウザは確認していません。