AI活用

AIストリーミングの再接続|Last-Event-IDと保存済み回答で表示を復元する

回線切断や再読み込みで途切れた応答を、同じ生成ジョブから表示し直す方法。Node.jsのイベント保存とEventSourceを使い、Last-Event-ID、初回cursor、期限切れを実動コードで確認します。

この記事の目次
  1. 「生成開始」と「続きの取得」を別のURLにする
  2. Node.jsで、生成結果を接続の外へ保存する
  3. ブラウザは、本文とcursorを一緒に復元してから購読する
  4. 切断・再読込・期限切れを、それぞれ再現する
  5. 実際のAI応答へ置き換える前に決めること

回線が切れた後に同じAI回答の続きを表示したいなら、回答を生成する処理と、ブラウザへ送る接続を分けます。生成した差分をサーバーへ保存し、再接続では未受信の差分だけを送り直します。生成用のPOSTを再送する構成にはしません。

ここで再開するのは、保存済みの回答を読むための接続です。Last-Event-IDを付ければ、中止済みのAI生成自体が途中から動き出すわけではありません。以下ではAPIキー不要の疑似生成を使い、Node.jsとブラウザで切断・再接続・リロードを再現します。

「生成開始」と「続きの取得」を別のURLにする

POST /runs
  → アプリのrun IDを発行し、生成を1回始める
  → 生成した差分をrunごとのイベント列へ保存する

GET /runs/:id
  → 保存済み本文・その時点のcursor・生成状態を取得する

GET /runs/:id/events?cursor=3
  → idが3より後のイベントを送り、生成中なら待ち続ける
  → 接続が切れても、生成処理はそのまま続く

run IDは、このアプリが付ける回答の識別子です。AI提供元のresponse IDとは別に管理します。一つのrunの中では、イベントへ1、2、3……と連番を付けます。

ブラウザのEventSourceは、接続が途切れると再接続し、受信したid:をLast-Event-IDヘッダーに載せられます。ただし、新しく作ったEventSourceの受信位置は空から始まります。ページの再読み込み後は、本文を復元してから初回の位置を渡す処理が必要です。これはHTML標準のLast-Event-IDの処理に沿った使い分けです。

既存のAPI応答をJavaScriptでストリーミング表示する記事は、ブラウザ切断時に上流の生成も止める構成を扱っています。今回のように切断中も生成を続けたい場合は、接続が閉じたときの処理を、そのまま流用しないでください。

スポンサーリンク

Node.jsで、生成結果を接続の外へ保存する

空のフォルダーへ、次のserver.mjsを保存します。追加パッケージは使いません。このデモはローカル専用で、127.0.0.1からの接続を受け付けます。

produce()が差分を保存し、閲覧用のGETがその配列を読みます。閲覧接続のcloseでは送信タイマーだけを解除します。生成用の処理には触れません。

import http from 'node:http';
import { randomUUID } from 'node:crypto';
import { readFile } from 'node:fs/promises';
import { setTimeout as delay } from 'node:timers/promises';
import { pathToFileURL } from 'node:url';

export async function* demoText({ signal }) {
  for (const text of ['接続が', '切れても、', '生成は', 'サーバーで', '続きます。', '同じ回答の', '保存済みイベントを', '再送して、', '表示を', '復元します。']) {
    await delay(300, undefined, { signal });
    yield text;
  }
}

export function createApp({ generate = demoText, ttlMs = 600_000 } = {}) {
  const runs = new Map();
  const metrics = { generationStarts: 0, subscriptions: [] };
  function append(run, type, data) {
    run.events.push({ id: run.events.length + 1, type, data });
  }
  async function produce(run) {
    metrics.generationStarts++;
    try {
      for await (const text of generate({ signal: AbortSignal.timeout(30_000) })) {
        if (run.events.length >= 1000 || run.text.length + text.length > 32_000) throw new Error('Demo limit');
        run.text += text;
        append(run, 'delta', { text });
      }
      run.state = 'completed';
      append(run, 'done', {});
    } catch {
      run.state = 'aborted';
      append(run, 'aborted', { message: '生成が終了しました。未完了の文章です。' });
    }
    run.expiresAt = Date.now() + ttlMs;
  }
  const server = http.createServer(async (req, res) => {
    // This demo binds to loopback and accepts same-origin browser requests only.
    const origin = `http://${req.headers.host}`;
    if (!/^127\.0\.0\.1:\d+$/.test(req.headers.host || '') ||
        (req.headers.origin && req.headers.origin !== origin)) return res.writeHead(403).end();
    const url = new URL(req.url, origin);
    const json = (status, value) => res.writeHead(status, { 'Content-Type': 'application/json', 'Cache-Control': 'no-store' }).end(JSON.stringify(value));
    for (const [id, run] of runs) if (run.expiresAt < Date.now()) runs.delete(id);
    if (req.method === 'GET' && ['/', '/client.js'].includes(url.pathname)) {
      const name = url.pathname === '/' ? 'index.html' : 'client.js';
      res.writeHead(200, { 'Content-Type': name.endsWith('.html') ? 'text/html; charset=utf-8' : 'text/javascript; charset=utf-8' });
      return res.end(await readFile(new URL(name, import.meta.url)));
    }
    if (req.method === 'POST' && url.pathname === '/runs') {
      if (runs.size >= 20) return json(429, { message: 'デモの保持件数上限です。' });
      const id = randomUUID();
      const run = { id, text: '', events: [], state: 'streaming', expiresAt: Infinity };
      runs.set(id, run);
      void produce(run); // GET connection lifetime does not control generation.
      return json(201, { id });
    }
    const match = url.pathname.match(/^\/runs\/([a-f0-9-]+)(\/events)?$/);
    if (req.method !== 'GET' || !match) return json(404, { message: 'Not found' });
    const run = runs.get(match[1]);
    if (!match[2]) return run ? json(200, { id: run.id, text: run.text, cursor: run.events.length, state: run.state }) : json(410, { message: '保存期限切れ、またはサーバー再起動で回答がありません。' });
    res.writeHead(200, { 'Content-Type': 'text/event-stream; charset=utf-8', 'Cache-Control': 'no-cache, no-transform', 'X-Accel-Buffering': 'no' });
    res.write('retry: 500\n\n');
    const reject = (message) => res.end(`event: unavailable\ndata: ${JSON.stringify({ message })}\n\n`);
    if (!run) return reject('回答を再開できません。保存期限を確認してください。');
    // Native reconnection header overrides the cursor from initial subscription.
    const raw = req.headers['last-event-id'] ?? url.searchParams.get('cursor') ?? '0';
    if (!/^\d+$/.test(raw) || !Number.isSafeInteger(Number(raw)) || Number(raw) > run.events.length) return reject('保存位置が不正です。ページを読み込み直してください。');
    let cursor = Number(raw);
    metrics.subscriptions.push({ id: run.id, header: req.headers['last-event-id'] ?? null, cursor });
    if (metrics.subscriptions.length > 100) metrics.subscriptions.shift();
    const flush = () => {
      for (const event of run.events.slice(cursor)) {
        const ok = res.write(`id: ${event.id}\nevent: ${event.type}\ndata: ${JSON.stringify(event.data)}\n\n`);
        cursor = event.id;
        if (!ok) return res.destroy(); // Slow readers reconnect from their acknowledged ID.
      }
      if (run.state !== 'streaming') res.end();
    };
    const timer = setInterval(flush, 50);
    res.on('close', () => clearInterval(timer)); // Unsubscribe only; do not abort producer.
    flush();
  });
  return { server, runs, metrics };
}

if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) {
  createApp().server.listen(5208, '127.0.0.1', () => console.log('http://127.0.0.1:5208'));
}

Last-Event-IDがある再接続では、それをURLのcursorより優先しています。URLには初回の値が残るため、毎回URL側を優先すると、前に受け取った差分を何度も送り直してしまいます。

この例は50ミリ秒ごとに未送信分を確認する、小さな再現用サーバーです。接続数が多いサービス向けの配送処理ではありません。また、完了から10分間だけ回答をメモリへ保持し、次のリクエスト時に期限切れを削除します。保持は最大20件、1回答は最大1,000差分・32,000文字です。上限や生成エラーで終わった回答は未完了として扱います。

ブラウザは、本文とcursorを一緒に復元してから購読する

同じフォルダーへclient.jsを保存します。再読み込み後に空の画面へ末尾の差分だけを追加すると、回答の前半が消えてしまいます。先にGET /runs/:idで本文と受信位置をセットで取得し、続きの購読を始めます。

const output = document.querySelector('#output');
const status = document.querySelector('#status');
const start = document.querySelector('#start');
let source = null;
let text = '';
let cursor = 0;
let epoch = 0;

async function restore(id) {
  const ticket = ++epoch;
  source?.close();
  source = null;
  const response = await fetch(`/runs/${encodeURIComponent(id)}`);
  const snapshot = await response.json();
  if (ticket !== epoch) return;
  if (!response.ok) throw new Error(snapshot.message);
  // Restore the text and its matching cursor together, then subscribe after it.
  text = snapshot.text;
  cursor = snapshot.cursor;
  output.textContent = text;
  status.textContent = snapshot.state;
  if (snapshot.state !== 'streaming') return;
  const current = new EventSource(`/runs/${encodeURIComponent(id)}/events?cursor=${cursor}`);
  source = current;
  for (const type of ['delta', 'done', 'aborted']) {
    current.addEventListener(type, (event) => {
      if (source !== current) return;
      const next = Number(event.lastEventId);
      if (next <= cursor) return;
      if (next !== cursor + 1) {
        current.close();
        status.textContent = 'イベントが欠けています。ページを読み込み直してください。';
        return;
      }
      const data = JSON.parse(event.data);
      if (type === 'delta') {
        text += data.text;
        output.textContent = text;
        status.textContent = 'streaming';
      } else {
        status.textContent = type === 'done' ? 'completed' : 'aborted(未完了)';
        current.close();
      }
      cursor = next;
    });
  }
  current.addEventListener('unavailable', (event) => {
    current.close();
    status.textContent = JSON.parse(event.data).message;
  });
  current.onerror = () => {
    if (source === current) status.textContent = current.readyState === EventSource.CLOSED ? '接続できません。再読込で保存状態を確認してください。' : '再接続中(生成を再実行しません)';
  };
}

start.addEventListener('click', async () => {
  epoch++;
  start.disabled = true;
  source?.close();
  source = null;
  try {
    const response = await fetch('/runs', { method: 'POST' });
    const data = await response.json();
    if (!response.ok) throw new Error(data.message);
    history.replaceState(null, '', `#${data.id}`);
    await restore(data.id);
  } catch (error) {
    status.textContent = error.message;
  } finally {
    start.disabled = false;
  }
});
if (location.hash.slice(1)) restore(location.hash.slice(1)).catch((error) => { status.textContent = error.message; });

run IDはURLの#より後へ置いています。ページを再読み込みしても、同じIDで保存状態を探せます。本文や受信位置を別々にブラウザへ永続保存する必要はなく、サーバーのスナップショットを起点にします。

スナップショットの取得後、EventSourceの接続までに生成が進んでも、その差分はサーバーに残っています。取得済みのcursorより後を送ることで、その間の差分も回収します。受信側は処理済みの連番を無視し、連番が飛んだら継ぎ足さずに止めます。

画面の入れ物として、index.htmlも同じ場所へ保存します。

<!doctype html>
<html lang="ja">
<meta charset="utf-8">
<meta name="viewport" content="width=device-width,initial-scale=1">
<title>SSE再接続のローカルデモ</title>
<style>body{max-width:700px;margin:40px auto;padding:16px;font:18px/1.8 system-ui}button{padding:12px}pre{white-space:pre-wrap;overflow-wrap:anywhere}</style>
<h1>SSE再接続のローカルデモ</h1>
<p>ボタンで新しい疑似生成を始めます。生成中に再読み込みすると、同じ回答を復元します。</p>
<button id="start">新しい疑似生成を始める</button>
<p id="status" role="status">未開始</p>
<pre id="output"></pre>
<script src="/client.js"></script>
</html>

フォルダー内で次を実行し、http://127.0.0.1:5208を開きます。localhostではなく、このURLを使ってください。

node server.mjs

「新しい疑似生成を始める」を押すと、約3秒かけて文章を作ります。このボタンは毎回新しいrunを作る操作です。同じ回答を復元するときはボタンを押し直さず、ページを再読み込みします。

切断・再読込・期限切れを、それぞれ再現する

回線切断は、ブラウザの開発者ツールでネットワークを一時的にOfflineへ切り替え、Onlineへ戻すと試せます。切断中にも疑似生成は進むので、戻した後に途中の文字を含めて表示されるか確認します。

手元の検証では、HTTP接続をサーバー側から実際に切断し、Chromeが送った再接続リクエストのヘッダーと、生成開始回数を記録しました。

再現したこと 確認結果
生成途中で接続を切断 Last-Event-ID付きで再接続。保存された全文と表示が一致し、生成開始は1回
完了後にページを再読込 完了済みの本文をスナップショットから復元。生成開始回数は増えない
別のrunを開始し、生成途中で再読込 同じrun IDで本文を復元し、0より大きい初回cursorから購読
保存期限切れのrunを開く 期限切れを表示。新しい生成は自動で始めない
生成側の処理が途中でエラー 部分的な文章を残し、abortedを通知。completedにはしない

実行確認はNode.js 24.13.0、Chromeで行いました。AI APIへは接続していません。「生成開始回数が増えない」は、このアプリが同じ生成を再実行していないことの確認です。実サービスの費用が発生しないという意味ではなく、切断中に続けた生成にも、提供元の条件に応じた処理費用が発生し得ます。

ネットワークが戻っても再接続できない場合は、ブラウザの状態表示と、サーバーの保持状態を分けて見ます。生成が完了していても、ログが消えていればこの方式では復元できません。

実際のAI応答へ置き換える前に決めること

実サービスでは、demoText()を、AIのテキスト差分を順に返す非同期イテレーターへ置き換えます。その接続はproduce()側で1回だけ開始し、閲覧用GETのたびに開始しません。渡されるsignalを上流リクエストにも渡し、タイムアウトで中止できるようにします。

差分が終わった理由も区別します。提供元が正常完了を通知した場合だけ、イテレーターを正常終了させます。途中切断、提供元エラー、未完了の終了は例外として扱わないと、今回のproduce()がcompletedを付けてしまいます。AI提供元との通信を復旧する方法は、このブラウザ再接続の実装とは別に必要です。

生成状態 再接続でできること できないこと
streaming 保存済み差分を再送し、以後の差分を送る 切断中の生成を無料にする
completed 保持期間内の完成した回答を復元する 削除済みの回答をメモリから復元する
aborted 未完了であることと、残った部分を表示する 中止された生成そのものをLast-Event-IDだけで再開する

本番で複数のサーバーやプロセスを使うなら、Mapを共有できる保存先へ置き換えます。本文・最終イベント番号・生成状態の整合性を保ち、イベント番号の採番を一つのrun内で重複させない設計が必要です。保存期間と、接続がない間も生成を続ける時間の上限も決めます。

さらに、run IDを知っているだけで他人の回答を読めないよう、生成時と閲覧時の両方で認証と所有者確認を行います。UUIDは認証の代わりになりません。本文を含む保存データの削除方針、生成開始POSTの重複防止、利用者が明示的に停止する経路も追加してください。このデモのPOSTには、リクエストの重複を排除する仕組みは含めていません。

まず確認するのは、切断後のネットワークログで新しい生成POSTが増えておらず、同じrunのGETが再接続しているかです。そのうえで表示された全文を保存結果と比較すると、「同じ回答の続き」を表示できているか判断できます。

スポンサーリンク