たとえば、クライアントが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つの規則で決まります。
- シーケンスはストリームごとに単調で、欠番がない。 そうすれば、クライアントは番号を見るだけで、欠けたイベントを検出できます。
- カーソルはクライアントにとって不透明にする。 クライアントは保存して返すだけです。何を符号化しているかを後から変えられます。
- 番号が順番どおりに見えるようになる。 見落とされるのはここです。書き込みの開始時に番号を割り当て、コミット時に見えるようになるとします。書き手 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つすべてに答えられるようにしておくべきです。
- 位置。 クライアントは何を持ち運び、それは不透明になっているか。SSE が標準化しているのはここです。
- 保持。 ログは、その位置をどれだけの期間守るか。これは、あなたが選んで書き残す数字です。
- リセット。 位置が古すぎるとき、サーバーは何をするか。クライアントがそれを見逃せない形になっているか。多くの設計が欠いているのはここで、回復するストリームと、黙ってずれていくストリームの分かれ目です。
確認した範囲と、していない範囲
すべて 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 はモデルで、データベースの実験ではありません。マルチノードのサーバー、永続的なログ、イベントストリームをバッファリングするプロキシ、ほかのブラウザは確認していません。