3 秒 ack
Slack の 3 秒 ack ルール、レースに強い event_id 重複排除、ack の背後に置く永続アウトボックス、そして 30 秒の ctx.waitUntil() 予算。
概要
Slack は Events API のリクエストに対して、3 秒以内の HTTP 2xx 応答を期待している。Slack の Events API ドキュメントによれば、この時間内に返せなかった場合は Slack のリトライ機構が動き出す。同じイベントが、場合によっては何度も再送されてくる一方で、最初のハンドラーはまだ走り続けているかもしれない。スラッシュコマンドやインタラクティビティのペイロードでも進め方は同じで、まず素早く ack し、実際の処理はそのあとに回す。
Cloudflare Worker の fetch ハンドラーはこの形に自然に馴染む。ただし、LLM の呼び出しやフォローアップメッセージの投稿といった、サードパーティとの往復を含む処理はすべて、レスポンスを返した後に走らせることが条件になる。ack より前に置いてよい唯一のものは、そのイベントを処理しなければならないという事実を記録する、速いローカルな書き込みだ。理由は後述の「ack の前に意図を永続化する」で展開する。
まず ack し、処理は ctx.waitUntil() へ
export default {
async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
const rawBody = await request.text();
if (!(await verifySlackSignature(request, rawBody, env))) {
return new Response("Unauthorized", { status: 401 });
}
const payload = JSON.parse(rawBody);
// Slack's URL verification handshake -- must be answered synchronously.
if (payload.type === "url_verification") {
return Response.json({ challenge: payload.challenge });
}
// Hand the real work to ctx.waitUntil() and return immediately.
// payload.event_id is the envelope's dedup key (see Deduplication below).
ctx.waitUntil(handleEvent(payload.event, env, payload.event_id));
return new Response(null, { status: 200 });
},
} satisfies ExportedHandler<Env>;ctx.waitUntil() は、レスポンスを送り終えたあとも指定した Promise のためにアイソレートを生かしておくよう Workers ランタイムへ伝えるものだ。これがないと、ランタイムは fetch() が返った時点で Worker を破棄してよいことになり、実行中の handleEvent() が途中で打ち切られうる。
この形は、失われても差し支えない処理に対しては正しい。しかし、失われては困る副作用に対しては十分ではない。後述の「ack の前に意図を永続化する」を参照。そこでは return の手前で何が起きるかが変わる。
永続化の書き込みは await する。外部への副作用は await しない
return の前に await handleEvent(...) を置くと、ハンドラーがリモート依存を伴う処理(LLM 呼び出し、Slack Web API の往復、サードパーティへの fetch)をした瞬間に 3 秒の予算を使い切る。ack より手前に置いてよいのは自分のデータストアへの書き込みだけだ。D1 のバッチ 1 回、ミリ秒オーダー、経路上にサードパーティのレイテンシーが一切ない。他人のサーバーと話す処理はすべてレスポンスのあとに回す。
リトライのセマンティクス
Slack は ack の失敗に対して最大3 回リトライする。ほぼ即座に 1 回、約 1 分後に 1 回、約 5 分後に 1 回だ。各リトライには 2 つのヘッダーが付く。
X-Slack-Retry-Num— 試行回数。1、2、3のいずれかX-Slack-Retry-Reason— リトライの理由(http_timeout、connection_failed、http_errorなど)
const retryNum = request.headers.get("x-slack-retry-num");
if (retryNum) {
// This is a Slack-initiated retry, not a first delivery.
console.log(`Slack retry #${retryNum}: ${request.headers.get("x-slack-retry-reason")}`);
}リトライが来たということは、Slack が最初の試行は失敗したと判断したという意味でしかない。そこには、Worker が実際にはイベントを処理し終えていたが 200 のレスポンスが間に合わなかった、というケースも含まれる。これが重複排除の問題そのものだ。リトライは新しい HTTP リクエストでありながら、すでに済ませた仕事を指しているかもしれない。
最初のリトライがどこに着地するかに注目してほしい。Slack はそれをほぼ即座に送ってくる。つまり 2 つの配送が同時に Worker に対して飛んでいる状態がありうる。重複排除の仕組みは、数分後の再生だけでなく、本物の並行性に耐えなければならない。
event_id による重複排除
すべてのイベントペイロードには、グローバルに一意な event_id が含まれている。Slack のドキュメントは特定の重複排除方式を規定していないが、event_id はまさにこのために用意されたフィールドだ。
get してから put する方式はレースに負ける
素直に書けば、読み取りのあとに書き込みを置く形になる。
// DO NOT COPY -- this loses the race it is meant to win.
async function handleEvent(event: SlackEvent, env: Env, eventId: string): Promise<void> {
const dedupeKey = `slack-event:${eventId}`;
if (await env.KV.get(dedupeKey)) {
return; // Already processed -- this is a retry.
}
await env.KV.put(dedupeKey, "1", { expirationTtl: 600 });
// ... do the real work
}これには互いに独立した 2 つの欠陥があり、Slack のリトライのタイミングはその両方をまともに突く。
読んでから書く操作はアトミックではない。 get と put のあいだには、2 つめの配送も同じくミスを読める窓がある。両方がそのまま先へ進み、副作用が 2 回走る。これは素朴な time-of-check/time-of-use のギャップであり、compare-and-set のプリミティブを持たないあらゆるストアに存在する。ストアの一貫性がどれだけ強くても関係ない。Workers KV には条件付き書き込みがないので、この窓をアプリケーション側から塞ぐことはできない。
Workers KV は結果整合である。 Cloudflare は、書き込みが他のネットワークロケーションで見えるようになるまで 60 秒以上かかりうること、そして存在しないという結果もキャッシュされること(「このキーは存在しない」という答え自体が同じ期間キャッシュされる)を文書化している。したがって、ほぼ即座に届く最初のリトライが別のロケーションに着けば、直前に書かれたキーに対してキャッシュされたミスを読みうる。伝播の窓は数十秒単位で、Slack の最初のリトライはその内側に届く。値を単一のトランザクションで読み書きする必要がある場面に KV は向かない、というのが Cloudflare 自身の案内でもある。
D1 によるアトミックなレシート
チェックを、UNIQUE キーに対する条件付き INSERT に置き換える。勝者を決めるのはデータベースであり、アプリケーションはその判定を読むだけだ。
CREATE TABLE slack_event_receipts (
event_id TEXT PRIMARY KEY,
event_type TEXT NOT NULL,
received_at INTEGER NOT NULL
);const batch = await env.DB.batch<{ changes: number }>([
env.DB.prepare(
`INSERT INTO slack_event_receipts (event_id, event_type, received_at)
VALUES (?, ?, ?)
ON CONFLICT (event_id) DO NOTHING`,
).bind(eventId, eventType, Date.now()),
env.DB.prepare("SELECT changes() AS changes"),
]);
// 1 -> this delivery inserted the receipt. 0 -> a concurrent or earlier
// delivery already owns this event_id.
const isFirstDelivery = batch[1].results[0].changes === 1;これを成立させている仕組みは 2 つあり、どちらも間違えやすい。
ON CONFLICT (event_id) DO NOTHINGは 1 つのステートメントである。 一意性のチェックと INSERT が同一の操作なので、そのあいだに窓は存在しない。同時に届いた 2 つの配送のうち、ちょうど 1 つが行を INSERT し、もう一方は競合したことを告げられる。changes()は直前のステートメントの結果を返す。 D1 のbatch()は、単一のトランザクション内で、逐次かつ非並行にステートメントを実行する。だからこそ、INSERT の直後に置いたSELECT changes()はその INSERT の行数を読む。あいだに別のステートメントを挟めば、違うものを報告することになる。バッチの外で別々のrun()として発行した場合、この保証はそもそも成り立たない。
バインドは ? による位置指定で行う。D1 は名前付きパラメーターをサポートしていない。
レシートのテーブルは際限なく増えるので、定期的に刈り込むこと。Slack のリトライのウィンドウは長くても 5 分程度なので、1 日より古い行を削除すれば安全で、マージンも十分すぎるほど残る。
KV もここで役に立つが、ベストエフォートの高速パスとしてのみだ。「すでに見た」というキャッシュ済みの答えがあれば安く処理を飛ばせるし、キャッシュされたミスを読んでも損はしない。実際に判定を下すのは D1 の INSERT だからだ。KV を唯一の防壁にしてはいけない。
意図的に無視するイベントにはすべて 200 を返す
重複はより広いルールの一例にすぎない。認証済みのイベントのうち、対応しないと決めたものにはすべて 200 を返す。 チャンネル違い、マッピングのないケース、扱わないメッセージサブタイプ、自分の bot 自身のメッセージ、そして重複。
// Every one of these is an ack, not an error: a retry would re-deliver the
// identical payload and be ignored identically.
if (event.channel !== env.WATCHED_CHANNEL_ID) return ack();
if (event.subtype && !HANDLED_SUBTYPES.has(event.subtype)) return ack();
if (event.bot_id) return ack();欲しくないイベントだからといってエラーステータスを返すのは、Slack の 3 回のリトライを 1 回消費して同じペイロードを受け取り直し、同じ結論に達するだけだ。変わるのは自分のエラー率と、エンドポイントの健全性についての Slack の評価だけである。
2xx 以外を返すのは、リトライによって違う結果に到達しうる唯一の状況、すなわちイベントを永続的に記録できなかった場合に限ること。それが次の節の話だ。(署名検証の失敗に対する 401 は別の軸である。認証されていないリクエストを拒否しているのであって、イベントを無視しているのではない。)
ack の前に意図を永続化する
ack してから waitUntil に渡す形は、落ちても構わない処理には十分だ。しかし必ず起こさなければならない副作用には足りない。waitUntil の処理はランタイムの破棄や時間予算によってキャンセルされうるうえ、その失敗は Slack からは見えないからだ。Slack はすでに 200 を受け取っており、そのイベントを二度と送ってこない。ack されたが記録されなかったイベントは、単に消える。
本番で使える形は、保証とレイテンシーを切り離す。同じ配送先が複数の経路から届きうる場合 — リトライされた配送と、ポーリングのスイープが同じ事実に着地する場合など — に必要になる、 事実の識別子キーを重ねたこの同じ台帳 + cron の形は、Workers での Events API も参照。
1. ack の前に、速い永続トランザクションを 1 つ
レスポンスを返す前に、単一のローカルトランザクションでレシートを記録し、同時に意図した副作用をアウトボックスの行としてキューに入れる。こうすると ack が依存するのは自分のデータベースへの書き込みだけになり、Slack Web API のレイテンシーにも LLM にも、そのほかどんなサードパーティにも依存しなくなる。
CREATE TABLE outbox (
id INTEGER PRIMARY KEY AUTOINCREMENT,
idempotency_key TEXT NOT NULL UNIQUE,
effect TEXT NOT NULL,
payload TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
attempts INTEGER NOT NULL DEFAULT 0,
claim_token TEXT,
claim_expires_at INTEGER,
created_at INTEGER NOT NULL,
completed_at INTEGER
);
CREATE INDEX outbox_claimable ON outbox (status, claim_expires_at, id);async function recordIntent(payload: SlackEnvelope, env: Env): Promise<boolean> {
const now = Date.now();
const batch = await env.DB.batch<{ changes: number }>([
env.DB.prepare(
`INSERT INTO slack_event_receipts (event_id, event_type, received_at)
VALUES (?, ?, ?)
ON CONFLICT (event_id) DO NOTHING`,
).bind(payload.event_id, payload.event.type, now),
// Must sit immediately after the receipt insert to report its row count.
env.DB.prepare("SELECT changes() AS changes"),
// Unconditional and idempotent: on a redelivery the UNIQUE key absorbs it.
env.DB.prepare(
`INSERT INTO outbox (idempotency_key, effect, payload, created_at)
VALUES (?, ?, ?, ?)
ON CONFLICT (idempotency_key) DO NOTHING`,
).bind(`${payload.event_id}:notify`, "notify", JSON.stringify(payload.event), now),
]);
return batch[1].results[0].changes === 1;
}アウトボックスへの INSERT を無条件に実行しているのは、それ自身の UNIQUE キーによって冪等になっているからだ。再配送では何も新しく積まれない。レシート側の changes() の判定は、もっと狭い用途に使う。この配送が即時試行も行うべきかどうかを決めるためである。
2. トランザクションが失敗したら ack しない
意図を永続化せずに ack してはいけない
ack 前のトランザクションが throw した場合、あるいは 3 秒のウィンドウ内に終わらない場合は、リトライ可能な 2xx 以外を返すこと。エラー応答が正解になる唯一のケースがこれだ。ここで 200 を返すと、Slack にはイベントを処理済みだと伝わって二度と送られてこないのに、それを実行しなければならないという記録はどこにも残っていない。イベントは黙って失われる。5xx はもう 1 回の配送を買い戻してくれる。
const ACK_BUDGET_MS = 2_000; // Headroom inside Slack's 3 seconds.
const budget = new Promise<never>((_, reject) =>
setTimeout(() => reject(new Error("pre-ack budget exceeded")), ACK_BUDGET_MS),
);
let isFirstDelivery: boolean;
try {
isFirstDelivery = await Promise.race([recordIntent(payload, env), budget]);
} catch (err) {
console.error("no durable intent recorded", err);
// Do NOT ack: only a retry can still save this event.
return new Response("retry", { status: 500 });
}このレースに負けても D1 への書き込みがキャンセルされるわけではない。500 を返したあとに commit されることは十分にありうる。それでも安全なのだが、安全である理由ははっきりしている。2 つの INSERT がどちらも冪等なので、Slack のリトライが同じトランザクションを走らせ直すと、レシートはすでに存在していて(changes() は 0 を返す)ack され、すでに積まれているアウトボックスの行はスイープに委ねられる。冪等でない ack 前の経路では、そもそもこのタイムアウトを許容できない。
3. 即時配送はあくまでベストエフォート
意図が永続化できたら、実際の配送は一般的なケースを速くするための最適化として ctx.waitUntil() に渡す。ここではあらゆるエラーを飲み込む。伝えるべき呼び出し元はもう残っておらず、バックグラウンドの Promise から throw しても、その実行をエラー扱いにする以外には何も達成しないからだ。
if (isFirstDelivery) {
ctx.waitUntil(
sweepOutbox(env).catch((err) => {
// Best effort. The cron sweep owns the guarantee; log and move on.
console.error("immediate delivery failed", err);
}),
);
}
return new Response(null, { status: 200 });ここで呼んでいるのは、cron のティックが呼ぶのと同じクレーム&配送ループ(後述の sweepOutbox)である点に注意してほしい。まったく同一のコードパスを、時計ではなくリクエストで起動しているだけだ。同期を保たなければならない別実装の「高速パス」は存在しないし、2 つのトリガーが衝突しないようにしているのは後述のリース機構である。
4. 契約を担うのは cron のスイープ
スケジュールされたティックが、即時試行が配送に失敗したもの、あるいはランタイムが途中でキャンセルしたものを、あらためてクレームし直す。このシステムについて考えるときは、自分に向かってはっきりこう言うとよい。即時配送は最適化であり、アウトボックスの cron スイープこそが契約である。 waitUntil の呼び出しを削除してすべての副作用をスイープだけに配送させてもよい、と思えないうちは、そのスイープはまだ正しくない。
5. クレームトークンによるリースと、フェンス付きの確定
アウトボックスには複数のワーカーが同時に到達しうる。即時試行、いま走っている cron のティック、そして時間を超過した前回のティックだ。リースは、それらが同じ仕事を二重に行わないようにする。
各試行はランダムなクレームトークンを発行し、自分が取った行に、期限を区切ってそれを刻む。
const LEASE_MS = 10 * 60 * 1000; // Comfortably longer than the worst-case delivery.
type OutboxRow = { id: number; effect: string; payload: string; attempts: number };
async function claimBatch(env: Env, token: string, limit: number): Promise<OutboxRow[]> {
const now = Date.now();
await env.DB.prepare(
`UPDATE outbox
SET claim_token = ?, claim_expires_at = ?
WHERE id IN (
SELECT id FROM outbox
WHERE status = 'pending'
AND (claim_expires_at IS NULL
OR claim_expires_at < ?
OR claim_token = ?)
ORDER BY id
LIMIT ?
)`,
).bind(token, now + LEASE_MS, now, token, limit).run();
const { results } = await env.DB.prepare(
"SELECT id, effect, payload, attempts FROM outbox WHERE claim_token = ? ORDER BY id",
).bind(token).all<OutboxRow>();
return results;
}性質は 3 つあり、それぞれ上の特定の句に由来する。
同時にクレームする者どうしは、互いに素な行をリースする。 適格性の判定と刻印が同一の
UPDATEに入っているため、データベースがそれらを直列化する。クレーム可能な id をSELECTしてから別のUPDATEを撃つ形にすると、重複排除の節で警告したのとまったく同じ time-of-check/time-of-use のギャップが復活する。同じトークンでのクレームのやり直しは冪等である。
claim_token = ?の選言が、そのトークンがすでに所有している行を適格なまま保つので、クラッシュ後にクレーム処理をやり直しても同じバッチを選び直してリースし直すことになる。ORDER BY id LIMIT ?と組み合わされているため、2 つめのより大きな集合へ黙って流れていくこともない。この選言がなければ、やり直しは自分の行を(リースがまだ有効なので)読み飛ばし、互いに素な別のバッチをリースしてしまう。つまり同じワーカーが 2 つのクレームを抱えることになる。行を解放するのは明示的なリリースではなく期限切れである。 バッチの途中で死んだワーカーは何も解放しない。リースが単に切れ、次のティックが再取得する。
そして確定は、そのトークンにフェンスされる。
async function completeClaimed(env: Env, id: number, token: string): Promise<boolean> {
const batch = await env.DB.batch<{ changes: number }>([
env.DB.prepare(
`UPDATE outbox
SET status = 'done', claim_token = NULL, claim_expires_at = NULL, completed_at = ?
WHERE id = ? AND claim_token = ?`,
).bind(Date.now(), id, token),
env.DB.prepare("SELECT changes() AS changes"),
]);
return batch[1].results[0].changes === 1;
}AND claim_token = ? がフェンスである。この試行がリースより長くかかり、別のワーカーがすでに新しいトークンでその行を再取得していたなら、この UPDATE は 0 行にしかマッチせず、完了を刻むことができない。フェンスがなければ、遅れて期限切れになった試行が、別のワーカーがいま実行中の仕事を確定してしまい、その 2 回目の配送はそのまま捨てられる。重複が損失に化けるということだ。フェンスによる拒否はログに残す価値がある。リースが実処理に対して短すぎたことを意味するからだ。
async function sweepOutbox(env: Env): Promise<void> {
const token = crypto.randomUUID();
for (const row of await claimBatch(env, token, 25)) {
try {
await performEffect(row, env);
} catch (err) {
await recordFailure(env, row.id, token, err); // Bump attempts, back off, park if exhausted.
continue;
}
if (!(await completeClaimed(env, row.id, token))) {
console.warn(`outbox ${row.id}: lease expired before ack; another worker owns it now`);
}
}
}6. 意図的に at-least-once を選ぶ
Slack への投稿とデータベース上の確定は別々のシステムであり、両者にまたがるトランザクションは存在しない。どちらかを先に置くしかなく、どちらの順序でも 2 つのあいだで中断されうる。
投稿してから確定 — 確定の書き込みが失敗すると、リースが切れ、スイープが行を再取得し、メッセージが 2 度投稿される。at-least-once。
確定してから投稿 — 投稿が失敗しても行はすでに完了扱いになっており、二度とリトライされない。副作用は黙って失われる。at-most-once。
この設計は前者を意図的に選んでいる。目に見える重複のほうが、見えない損失よりましな失敗だからだ。exactly-once であるかのようにほのめかすのではなく、自分たちのドキュメントにもそう書くこと。副作用が exactly-once だと信じた読み手は、そうでなくなった最初の瞬間に壊れるものを下流に作ってしまう。
窓を塞いだふりをせずに狭める手立ては 2 つある。選べる場面では、形からして冪等な副作用を優先すること。すでに ts を握っているメッセージの更新や、already_reacted を返すリアクションの追加は、再配送を無害に吸収する。一方 chat.postMessage には呼び出し側が指定できる冪等キーがなく、平然と 2 通目のメッセージを作る。そしてリースは最悪ケースの配送より十分に長く取り、フェンスがめったに拒否せずに済むようにしておくこと。
waitUntil の予算と 3 つの逃げ道
Cloudflare は waitUntil() の処理に対して、実行の終了時点から測って 30 秒という上限を文書化している。この予算はリクエスト内のすべての waitUntil() 呼び出しで共有される。失敗の仕方は正確に押さえておく価値がある。時間内に解決しなかった Promise は、reject されるのではなくそのままキャンセルされる。捕まえられるエラーはなく、走る finally もなく、後始末の機会もない。ctx.waitUntil() のリファレンスを参照。
その波及は、キャンセルされた処理が握っていたものすべてに及ぶ。実行ロック、アウトボックスのリース、「同期中」フラグ。どれも解放されない。解放するはずのコードが実行されないからだ。取り残されたロックは、それ自身の陳腐化ウィンドウが経過するまで以降のティックをすべて塞ぐ。上のアウトボックスのリースが finally ブロックでのクリアではなくタイムスタンプによる期限切れになっているのは、まさにこれが理由である。キャンセルを生き延びる解放手段は期限切れだけだ。
3 つの逃げ道を、望ましい順に挙げる。
予算内に収める。 1 回の実行あたりの処理量を区切る。クレームのバッチサイズを固定し、ページサイズに上限を設け、上流の結果セットを無制限にループしない。残りは次のティックに任せる。
リクエストの外へ移す。 キューか cron のスイープ。上のアウトボックスがまさにこのパターンだ。無関係な理由で Worker が破棄された場合まで生き延びるのは、この逃げ道だけである。
fetch()の中でインラインに await する。 直感に反するが、ある特定の形に対しては正しい。Cloudflare は HTTP でトリガーされた Worker には実行時間のハード上限がないと文書化している。クライアントが接続を保っているかぎり Worker は処理を続けられる。効いてくる制約は CPU 時間のほうで(有料プランはデフォルト 30 秒、最大 5 分、無料プランは 10 ms)、I/O 待ちの時間は CPU を消費しない。したがって、HTTP でトリガーされ、I/O 待ちが支配的で、呼び出し元が本当の結果を欲しがっている処理については、インラインに await するほうがwaitUntilより長く生き、しかも答えを返せる。
3 つめの逃げ道は Slack の ack 経路には使えない
インラインの await が成立するのは、クライアントが待ってくれるからだ。ここでのクライアントは Slack であり、3 秒で待つのをやめる。逃げ道 3 は自前の管理エンドポイント、手動のバックフィル、社内ツール向けであって、Events API・スラッシュコマンド・インタラクティビティのハンドラーに使うものではない。
スイープのサイズは cron 側の制限に合わせて決める。こちらはまた別の値だ。有料プランの scheduled() の実行は、cron の間隔が 1 時間未満なら CPU 時間 30 秒(1 時間以上の間隔なら 15 分)を得られ、全体は 15 分の実時間上限の下に置かれる(Workers の制限)。1 分ごとに走る上限付きのバッチは、途中で打ち切られる無制限のバッチと同じくらいちゃんとバックログを掃き出す。
つまずきどころ
ctx.waitUntil()のなかで起きた失敗は、呼び出し元からは見えない。 Slack はすでに200を受け取っている。バックグラウンドの Promise が throw した場合は、その Promise の内側で捕捉してログに残すしかない。それを載せる HTTP ステータスはもう残っていない。永続的な記録が ack より手前に存在していなければならないのは、まさにこのためだ。判断の基準は
event_idであって、X-Slack-Retry-Numではない。 リトライは最初の試行が失敗したという Slack の主張であって、処理が済んだ証拠ではない。最初の配送が意図の永続化まで到達していなかったのなら、成功しなければならないのはそのリトライのほうだ。いま自分がどちらのケースにいるかを知っているのは、レシートの行だけである。3 秒の予算には、署名検証も、JSON のパースも、そして永続化の書き込みも含まれる。 3 つともサードパーティのレイテンシーを含まないローカルな処理であり、だからこそこの予算には余裕がある。その状態を保ち、
returnより手前に余計なものが忍び込まないようにすること。スラッシュコマンドとインタラクティビティは ack のパターンこそ共有するが、
event_idを持たない。 3 秒ルールも、アウトボックスも、リースもそのまま通用する一方、レシートのキーにできるグローバルに一意なエンベロープ id は存在しない。これらのサーフェスにはアプリケーションレベルの操作 id が必要になるし、同一のスラッシュコマンド 2 回は重複ではなく 2 回の意図的な実行でありうる。