zudo-slack-wisdom
GitHub リポジトリ

検索したい単語を入力

いつでも検索バーを開ける

Cron による定期投稿

外部データを Slack へミラーするスケジュール実行の Worker -- at-least-once なティック、多重実行の抑制、ローカル時刻のウィンドウ、冪等な投稿、ティックごとの予算。

概要

受信 webhook を捌いているのと同じ Worker を、時計で動かすこともできる。Cron Trigger は、wrangler.toml で定義したスケジュールに従って scheduled() ハンドラーを起動する。HTTP リクエストとは無関係に動く経路だ。

[triggers]
crons = ["*/5 * * * *"]
export default {
  async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
    return await handleRequest(request, env, ctx);
  },

  async scheduled(controller: ScheduledController, env: Env, ctx: ExecutionContext): Promise<void> {
    ctx.waitUntil(syncToSlack(env));
  },
} satisfies ExportedHandler<Env>;

よくある形は、外部データベースの行を Slack チャンネルへステータスボードとしてミラーするというものだ。その既定のサーフェスはメッセージ 1 通で、一度投稿したらあとはその場で編集し続ける。初回は chat.postMessage、以降のティックはすべて chat.update になる。その場で編集する方式は、読み取り専用性のために設定を一切必要としない唯一の Slack サーフェスでもある。「更新できるのは認証したユーザーが投稿したメッセージだけ」なので、設定を間違えうる共有ダイアログも ACL も存在せず、しかも編集はチャンネルに通知を出さない。

async function syncToSlack(env: Env): Promise<void> {
  const rows = await fetchExternalData(env);
  const blocks = buildStatusBlocks(rows);

  const messageTs = await env.KV.get("status-board:message-ts");
  if (!messageTs) {
    const res = await callSlackApi<{ ts: string }>(
      "chat.postMessage",
      { channel: env.STATUS_CHANNEL_ID, blocks },
      env.SLACK_BOT_TOKEN,
    );
    await env.KV.put("status-board:message-ts", res.ts);
  } else {
    await callSlackApi(
      "chat.update",
      { channel: env.STATUS_CHANNEL_ID, ts: messageTs, blocks },
      env.SLACK_BOT_TOKEN,
    );
  }
}

webhook ハンドラーの場合と同じく、ここでも ctx.waitUntil() が効いてくる。Workers ランタイムは関数が返った時点で scheduled() の呼び出しは終わったとみなすため、実際の Slack 呼び出しを waitUntil() で登録しておかないと、完了する前にアイソレートが破棄されかねない。

この節から先はすべて、ひとつの問いをめぐる話だ。ティックがきっかり 1 回だけ走るとは限らないとき、何が起きるのか。Cloudflare は同じティックを複数回配送しうるし、2 つの実行は重なりうるし、Slack とデータベースを一緒にコミットすることはできない。続く 3 つの節はひと続きのスタックで、弱い層から順に並んでいる。ティックごとのクレームは重複した配送をつぶし、リースは重なった実行を抑制し、アイテムごとのマーカーだけが Slack への重複投稿 を実際に防ぐ。最後の層だけに頼れるくらい小さなジョブなら、それは擁護できる設計だ。だが最後の層を省くジョブは擁護できない。

Cron の配送は at-least-once

Cron Trigger は exactly-once ではない。ランタイムは同じティック、つまり同じ cron 式と同じ scheduledTime に対して scheduled() を複数回呼び出しうる。プラットフォーム自身の API の中で最もはっきりした手がかりが ScheduledController.noRetry() だ。再配送を抑止することだけを目的にしたメソッドは、再配送するシステムでなければ存在する意味がない。

重複をつぶすには、ティックの同一性をキーにしたクレームを、処理が始まる前に書く。

const TICK_CLAIM_TTL_SECONDS = 900; // Comfortably longer than the cron interval.

// `scheduledTime` is unique per tick, so a claim for one tick can never block
// the next one -- nothing has to clear it, the TTL does that.
function tickClaimKey(controller: ScheduledController): string {
  return `tick-claim:${controller.cron}-${controller.scheduledTime}`;
}

async function claimTick(env: Env, key: string): Promise<boolean> {
  if (await env.KV.get(key)) return false;
  await env.KV.put(key, String(Date.now()), { expirationTtl: TICK_CLAIM_TTL_SECONDS });
  return true;
}

export default {
  async scheduled(controller: ScheduledController, env: Env, ctx: ExecutionContext): Promise<void> {
    ctx.waitUntil(
      (async () => {
        const key = tickClaimKey(controller);
        if (!(await claimTick(env, key))) {
          console.log(`[cron] duplicate delivery of ${key}, skipping`);
          return;
        }
        await syncToSlack(env, controller);
      })(),
    );
  },
} satisfies ExportedHandler<Env>;

調整すべきつまみは TTL だけで、cron の間隔より十分に長くしたい。遅れて届いた再配送でもクレームを見つけられる程度には長く、キーが永遠に溜まり続けない程度には短く、というバランスだ。scheduledTime はティックごとに異なるので、TTL を長めにしても正しさの面では何も失わない。次のティックのキーは別物だからだ。

管理者向けの手動トリガーも生やしている Worker は、同じパイプラインを HTTP リクエストから走らせることになる。そこには、つぶすべきティックが存在しない。この経路には呼び出しごとに一意なキーを与えて、クレームを偶発的なロックではなく意図的な no-op にしておく。

// A manual run has no scheduled tick to deduplicate against, so its claim key
// is unique per invocation -- the claim always succeeds, by design.
function manualClaimKey(): string {
  return `tick-claim:manual-${crypto.randomUUID()}`;
}

この層はベストエフォートであり、権威ある仕組みではない

KV.get() に続けて KV.put() を呼ぶ流れはアトミックではないし、KV は結果整合性だ。同じティックの 2 つの配送が近い時刻に着地すれば、両方とも「まだクレームされていない」と読んで両方とも進んでしまいうる。論理的な入口が 1 つしかなくても、これは起こる。だからこれは重複排除ではなく重複の抑制として扱う。よくあるケースを安価に取り除くだけで、まれなケースには何もしない。ティック単位の重複排除が権威あるものでなければならないなら、KV のクレームを D1 への一意制約付き INSERT に置き換え、制約違反そのものを「すでにクレーム済み」のシグナルにする。

この層がカバーしないものにも注意しておきたい。これがつぶすのは 1 つのティックの重複配送だ。個々の行が 2 回投稿されるかどうかについては何も言っていない。そちらはアイテムごとの冪等性という別の仕組みで、順序のルールも異なる。後述の「冪等な投稿: 成功のあとにマークする」で扱う。

多重実行の抑制: D1 のハートビートリース

Cron Trigger は排他制御を保証しない。1 ティック分の処理が起動間隔より長くかかることがときどきあると(上流 API が遅い、バッチが大きいなど)、2 つの呼び出しが同時に走りうる。ミラージョブなら、2 つの実行が同じ ts に対して chat.update を競うことで Web API の呼び出しが 1 回無駄になるし、2 つの実行がソースデータの別々のスナップショットを読んでいれば作りかけのボードを公開してしまうこともある。行ごとに通知を投稿するジョブでは、重なりの影響はもっと悪い。誰かのチャンネルに重複したメッセージが出るということだからだ。

訂正: 短い TTL の KV リースでは足りない

このページの以前の版では、実行ロックとして短い TTL の KV リースを勧めていた。あの助言が成り立つのは、入口が cron ティックだけで、最悪ケースが「1 ティック飛ばす」で済むジョブに限られる。実運用のデプロイはほぼ必ず 2 つ目の入口が生えてくる。管理者向けの手動トリガー、webhook 起点の同期といったものだ。そうなると KV の結果整合性(キャッシュの下限がおよそ 60 秒)のせいでリースは成立しなくなる。別々の入口から始まった 2 つの実行が、両方とも「ロックなし」と読んで両方とも進みうるからだ。上で KV のリースを残しているのはティックごとのクレームとしてだけで、そこでは抑制に失敗しても重複するのは処理であって投稿ではない。

現実に持ちこたえるパターンは、D1 のシングルトン行をアトミックな条件付き UPDATE でクレームし、TTL ではなくハートビートで生かし続けるというものだ。

CREATE TABLE sync_lease (
  id           INTEGER PRIMARY KEY CHECK (id = 1), -- Singleton row.
  run_id       TEXT,
  holder       TEXT,
  started_at   INTEGER,
  heartbeat_at INTEGER
);

INSERT OR IGNORE INTO sync_lease (id, run_id) VALUES (1, NULL);
const HEARTBEAT_INTERVAL_MS = 2 * 60 * 1000;
const STALE_AFTER_MS = 10 * 60 * 1000;

// One atomic conditional UPDATE does the whole acquisition: the row is claimable
// only when it is free, or when its last heartbeat is older than the staleness
// window. `meta.changes === 0` means we lost the race -- there is no separate
// read to go stale between check and write.
async function acquireLease(env: Env, holder: string): Promise<string | null> {
  const runId = crypto.randomUUID();
  const now = Date.now();

  const res = await env.DB.prepare(
    `UPDATE sync_lease
        SET run_id = ?, holder = ?, started_at = ?, heartbeat_at = ?
      WHERE id = 1
        AND (run_id IS NULL OR heartbeat_at < ?)`,
  )
    .bind(runId, holder, now, now, now - STALE_AFTER_MS)
    .run();

  return res.meta.changes === 1 ? runId : null;
}

// Fenced by run_id: a run that was already taken over cannot renew or release
// a lease that now belongs to somebody else.
async function renewLease(env: Env, runId: string): Promise<boolean> {
  const res = await env.DB.prepare(
    `UPDATE sync_lease SET heartbeat_at = ? WHERE id = 1 AND run_id = ?`,
  )
    .bind(Date.now(), runId)
    .run();
  return res.meta.changes === 1;
}

async function releaseLease(env: Env, runId: string): Promise<void> {
  await env.DB.prepare(
    `UPDATE sync_lease SET run_id = NULL, holder = NULL WHERE id = 1 AND run_id = ?`,
  )
    .bind(runId)
    .run();
}

運用の心地よさを決めているのはハートビートだ。TTL 方式では、現実的に最も遅いティックを前もって当てにいく必要があり、どちらかの方向に必ず外す。短すぎれば、正当に時間のかかっている実行が自分のロックを失う。長すぎれば、クラッシュした実行が期限切れまですべてのティックを塞ぐ。2 分ごとに更新し、最後のハートビートから 10 分経ってはじめて stale とみなすようにすれば、この推測は丸ごと消える。遅くても生きている実行は更新を続けて必要なだけロックを持ち続けるし、退避させられた Worker は更新をやめてリースが自ら解ける。

更新はそのまま所有権チェックを兼ねるので、実行側が処理の合間に問うべきことは 1 つで済む。

class Lease {
  private lastBeat = Date.now();

  constructor(
    private readonly env: Env,
    readonly runId: string,
  ) {}

  // Answers the only question a work loop cares about: do we still own this run?
  // Renews when the heartbeat interval has elapsed, and otherwise just reads.
  async stillOurs(): Promise<boolean> {
    if (Date.now() - this.lastBeat < HEARTBEAT_INTERVAL_MS) {
      const row = await this.env.DB.prepare(
        `SELECT run_id FROM sync_lease WHERE id = 1`,
      ).first<{ run_id: string | null }>();
      return row?.run_id === this.runId;
    }

    const renewed = await renewLease(this.env, this.runId);
    if (renewed) this.lastBeat = Date.now();
    return renewed;
  }
}

for (const item of deliverables) {
  // Check ownership immediately before anything user-visible, not once at the
  // top of the run. A takeover between two items must not produce a second post.
  if (!(await lease.stillOurs())) {
    console.warn("[cron] lease lost mid-run, stopping before the next post");
    break;
  }
  await postAndStamp(env, lease.runId, item);
}

2 つの入口は、競合に負けたときに求める振る舞いが違う。cron ティックはきれいにスキップすればいい。次のティックがその処理を拾ってくれる。手動トリガーには待っている呼び出し元がいるので、いま誰がリースを持っているのかを伝える。

const runId = await acquireLease(env, "manual");
if (!runId) {
  const holder = await env.DB.prepare(
    `SELECT holder, started_at, heartbeat_at FROM sync_lease WHERE id = 1`,
  ).first();
  return Response.json({ error: "sync_in_progress", holder }, { status: 409 });
}

これは多重実行の抑制であって、at-most-one の保証ではない

このリースを排他制御として読んではいけない。奪取は時間ベースであり、ハートビートを止めた保持者が処理を止めた保持者とは限らない。遅い fetch で止まっていたアイソレートは、リースが再割り当てされたあとに目を覚まし、自分が置き換えられたことを知らないまま Slack を呼びうる。ここから 3 つの帰結が出てくる。3 つとも設計を支える要素だ。

  1. リースを失っていないかのチェックは、副作用のたびに直前で行う。 実行の開始時に 1 回だけではだめだ。

  2. データベースへの書き込みは run_id でフェンスする。 置き換えられた実行が、新しい保持者がすでに進めた状態を上書きできないようにする。

  3. 成果物ごとの冪等性を維持する。 リースは重複をまれにするだけだ。二重に公開されることを不可能にするのは、次の節のマーカーだけだ。

このリースを正当化しているのはコストの非対称性であり、これはある本番リファレンス統合でも裏づけられている。飛ばしたティックは目に見えず、勝手に回復する。一方 Slack への重複投稿はユーザーの目に触れ、見なかったことにはできない。この設計は複雑さの点では Durable Object の下、健全さの点では KV リースの上に位置する。ワークロードが抑制ではなく本物の直列化を必要とするなら、Durable Object に手を伸ばすとよい。

冪等な投稿: 成功のあとにマークする

Slack とストレージをアトミックにコミットすることはできない。chat.postMessage と D1 の UPDATE にまたがるトランザクションは存在しないので、何らかの順序を選び、何らかの失敗ウィンドウを受け入れるしかない。選ぶべきは、失敗モードが回復可能なほうだ。

通知ループはマーカーの不在から駆動し、マーカーは Slack が投稿を受け付けたあとにだけ書く。

// The query is the queue: anything without a marker is still owed a post, so a
// crashed run resumes exactly where it stopped on the next tick.
const pending = await env.DB.prepare(
  `SELECT id, payload FROM deliverables
    WHERE notified_at IS NULL AND due_date = ?
    ORDER BY id
    LIMIT ?`,
)
  .bind(localDate, BATCH_LIMIT)
  .all<{ id: string; payload: string }>();

for (const row of pending.results) {
  await postToSlack(env, row.payload);

  // Stamp only once Slack accepted, and repeat the NULL check in the WHERE
  // clause so a concurrent run holding the same row cannot double-stamp it.
  await env.DB.prepare(
    `UPDATE deliverables SET notified_at = ? WHERE id = ? AND notified_at IS NULL`,
  )
    .bind(Date.now(), row.id)
    .run();
}

残るウィンドウは「Slack が投稿を受け付けた」から「マーカーが書かれた」までの隙間だ。ここでクラッシュすると、次のティックは notified_at IS NULL を見てそのアイテムをもう一度投稿する。重複はきっかり 1 つ。このウィンドウは狭められる(2 つの操作を隣り合わせに保ち、マーカーの書き込みをループの最後にまとめたりしない)が、消すことはできない。これが「ベストエフォートで 1 回」であり、意図的なトレードだ。重複は目に見え、説明でき、削除できる。逆の順序は、投稿されなかったアイテムを配信済みとしてマークする形で失敗するが、届かなかったメッセージには誰も気づかない。

2 つの順序、2 つの失敗モード

このページに出てくる 2 つの順序ルールはどちらも正しく、そしてどちらも相手の場所では間違いになる。

処理の前にクレームするのは、重複した処理をスキップしたいときだ。ティックごとのクレームを先に書くことで、再配送はそれを見つけて何もしない。

成功のあとにマークするのは、ユーザーの目に触れる投稿を絶対に取りこぼしたくないときだ。アイテムごとのマーカーは、Slack が受け付けたあとにだけ書かれる。

この 2 つを入れ替えると、それぞれが防ごうとしていた失敗がそのまま起きる。処理のあとに書いたクレームは重複した実行を通してしまうし、投稿の前に書いたマーカーは、記録するはずだった投稿を黙って飲みこんでしまう。

マーカーの列は成果物ごとに 1 つにする。実行ごとでも日ごとでもない。「今日は同期済み」というフラグ 1 つでは、1 件の失敗がその後ろにいる他のすべてのアイテムの同日中の再試行を封じてしまう。成果物ごとのマーカーなら、各アイテムがそれぞれのペースで失敗し、回復できる。

競合の激しいキュー、つまり同じテーブルに対して多数のワーカーが並行して動くような場合には、素のマーカーをリース付きのクレームへ格上げする。各ワーカーがランダムなトークンを書き、処理はそのトークンで承認し、クラッシュしたクレーム者の行はリース期限が切れたときに再びクレーム可能になる。永遠に詰まったままにはならない。

// Claim: take the row only if nobody holds it, or the previous claim expired.
const token = crypto.randomUUID();
const now = Date.now();

const claimed = await env.DB.prepare(
  `UPDATE deliverables
      SET claim_token = ?, claim_expires_at = ?
    WHERE id = ? AND notified_at IS NULL
      AND (claim_token IS NULL OR claim_expires_at < ?)`,
)
  .bind(token, now + CLAIM_LEASE_MS, row.id, now)
  .run();

if (claimed.meta.changes === 1) {
  await postToSlack(env, row.payload);

  // Acknowledge with the same token: a claim that expired mid-post and was
  // taken over by another worker cannot stamp the row out from under it.
  await env.DB.prepare(
    `UPDATE deliverables
        SET notified_at = ?, claim_token = NULL
      WHERE id = ? AND claim_token = ? AND notified_at IS NULL`,
  )
    .bind(Date.now(), row.id, token)
    .run();
}

ローカル時刻でのスケジューリング

Cron Trigger は UTC で発火する。ローカル時刻で表現された投稿ウィンドウ(「JST の 08:00 から 21:00 のあいだ 15 分ごと」)は手で変換するしかなく、その変換には罠がある。開始時刻が UTC オフセットより手前にあるローカルウィンドウは UTC の深夜をまたぎ、連続しない UTC 時刻の集合として着地するのだ。

さらにもうひとつの制約が重なる。cron の分フィールドは、その式がマッチするすべての時刻に適用される。したがって 1 つの式で「23:00 から 11:59 まで 15 分ごと、そして 12:00 ちょうど」と言うことはできない。このウィンドウにはルールが 2 つ必要だ。時刻レンジのルールと、締めのティックのルールである。

[triggers]
crons = [
  "*/15 0-11,23 * * *", # 08:00-20:45 local; the window wraps past UTC midnight.
  "0 12 * * *",         # 21:00 local exactly -- the closing tick.
]

このコメントがマッピングの唯一のドキュメントであり、だからこそ腐る。ウィンドウの変更、サマータイムの前提、ルールの追加、どれもティックを意図したローカル時刻の外へ静かに押し出しうるのに、何も落ちない。ジョブはただ 07:45 に投稿するだけだ。このマッピングは、Worker がデプロイに使うのと同じ設定を読むテストにしてしまう。

// Expand every deployed cron expression into its individual ticks, convert each
// to local time, and assert it lands inside the window. A schedule edit that
// drifts outside the intended hours then fails CI instead of surprising users.
import { CRON_EXPRESSIONS, POSTING_WINDOW } from "../src/config";
import { expandCron, toLocalMinutes } from "./helpers/cron";

test("every cron tick falls inside the local posting window", () => {
  const ticks = CRON_EXPRESSIONS.flatMap((expression) => expandCron(expression));
  expect(ticks.length).toBeGreaterThan(0);

  for (const utcTick of ticks) {
    const local = toLocalMinutes(utcTick, POSTING_WINDOW.timeZone);
    expect(local).toBeGreaterThanOrEqual(POSTING_WINDOW.startMinutes);
    expect(local).toBeLessThanOrEqual(POSTING_WINDOW.endMinutes);
  }
});

このテストに価値があるのは、CRON_EXPRESSIONS がデプロイで実際に使われている配列である場合だけだ。wrangler.tomlcrons をそこから生成するか、両者が一致することをアサートする。スケジュールを手でコピーした複製に対するテストは何も証明しない。

日付は Date.now() ではなく scheduledTime から導く

ハンドラーの内側では、日付とウィンドウの計算はすべて controller.scheduledTime、つまりそのティックが走るはずだった時刻から導き、Date.now() からは決して導かない。コールドスタートで遅れたティックや、数分遅れて再配送されたティックも、スケジュールされたローカル日付のもとで動かなければならない。Date.now() はときどきローカルの深夜をまたいでしまっていて、実行に翌日分のボードを計算させてしまう。Date.now() は実行中の経過時間の計測、つまりハートビート、デッドライン、リトライのバックオフのために取っておく。

const localDate = toLocalDate(controller.scheduledTime, POSTING_WINDOW.timeZone);

呼び出し種別ごとのリトライ予算

ティック全体に単一のリトライポリシーを当てると、両方向に間違う。API のラッパーは 1 つに保ったまま、呼び出しの種別でパラメータ化する。maxRetriesmaxRetryAfterMs のクランプ、そして 5xx をそもそもリトライするかどうかだ。

interface RetryPolicy {
  maxRetries: number;
  maxRetryAfterMs: number;
  retryServerErrors: boolean;
}

// Per-item posting: fail fast. A long Retry-After sleep here stalls every
// remaining item in the loop, and the un-set marker retries this one next tick.
const POSTING_POLICY: RetryPolicy = {
  maxRetries: 2,
  maxRetryAfterMs: 2_000,
  retryServerErrors: false,
};

// Once-per-tick reads and sweeps: be patient. A throttled non-Marketplace app
// can legitimately be told to wait ~60s, and there is no item loop to stall.
const SWEEP_POLICY: RetryPolicy = {
  maxRetries: 5,
  maxRetryAfterMs: 60_000,
  retryServerErrors: true,
};

async function callSlackApiWithPolicy<T>(
  method: string,
  body: unknown,
  token: string,
  policy: RetryPolicy,
): Promise<T> {
  for (let attempt = 0; ; attempt++) {
    const res = await fetch(`https://slack.com/api/${method}`, {
      method: "POST",
      headers: {
        Authorization: `Bearer ${token}`,
        "Content-Type": "application/json; charset=utf-8",
      },
      body: JSON.stringify(body),
    });

    if (res.status === 429) {
      const waitMs = Number(res.headers.get("retry-after") ?? 1) * 1000;
      if (attempt >= policy.maxRetries || waitMs > policy.maxRetryAfterMs) {
        throw new SlackRateLimitError(method, waitMs);
      }
      await sleep(waitMs);
      continue;
    }

    if (res.status >= 500) {
      if (!policy.retryServerErrors || attempt >= policy.maxRetries) {
        throw new SlackServerError(method, res.status);
      }
      await sleep(backoffMs(attempt));
      continue;
    }

    return (await res.json()) as T;
  }
}

この非対称性こそが要点だ。投稿の経路では、辛抱強さはむしろ害になる。30 秒の Retry-After を眠って待てば、ループに残っている他のすべてのアイテムが止まる。しかもより安価な回復手段がすでにある。そのアイテムのマーカーは立っていないので、次のティックがただで再試行してくれるのだ。そこで 5xx をリトライせず即座に throw しているのも同じ理由による。一時的な Slack のエラーは、バッチを止めてまで待つ価値がない。一方、ティックに 1 回だけの読み取りやスイープの経路には、その呼び出しの後ろに並んでいるキューがない。だから待つのが最も安価な選択肢になるし、スロットルされた非 Marketplace アプリに対する長い Retry-After は、エラーではなく従うべき正常で正しい指示だ。

読み取り側の予算と公平性

上流 API をページングする読み取りスイープには自然な停止点がなく、終わるまで走り続けるティックは、いずれリースのハートビートを追い越す。スイープにはティックごとの明示的な予算を与える。ページ数、呼び出し数、そしてリクエストの途中ではなく処理の合間にチェックする緩い実時間デッドラインだ。この同じバルク+末尾の予算の形を、非 Marketplace の読み取り上限のもとで動く conversations.history のスイープに対して詳しく扱っているのがチャンネル履歴の読み取りであり、以下の二段階のカバレッジパターンはそれを踏まえている。

interface SweepBudget {
  maxPages: number;
  maxCalls: number;
  deadlineAt: number; // Soft: checked between work units, never mid-request.
}

function newSweepBudget(): SweepBudget {
  return { maxPages: 20, maxCalls: 60, deadlineAt: Date.now() + 20_000 };
}

function budgetExhausted(budget: SweepBudget): boolean {
  return budget.maxPages <= 0 || budget.maxCalls <= 0 || Date.now() >= budget.deadlineAt;
}

予算の枯渇はエラーではなく、きれいな打ち切りだ。 スイープはその場で止まり、打ち切ったことを記録し、次のティックが同じローテーションの続きから再開する。枯渇で例外を投げると、定常状態の正常な状況がアラートに化け、チームはアラートを読まなくなる。

カバレッジは 2 つのフェーズに分けると考えやすい。バルクのページングが直近のウィンドウを安価にさばき、末尾ローテーションによるアイテムごとの取得が、バルクのウィンドウから漏れたものを拾う。最終ポーリングが古い順、未ポーリングのものはすべての先頭に置く。

async function sweep(env: Env, budget: SweepBudget): Promise<SweepStats> {
  const stats = { pages: 0, tailFetched: 0, truncated: false };

  // Phase 1 -- bulk: page the recent window until the budget runs out.
  let cursor: string | undefined;
  do {
    if (budgetExhausted(budget)) {
      stats.truncated = true;
      break;
    }
    const page = await fetchRecentPage(env, cursor, budget);
    await upsertRows(env, page.rows);

    // Stamp last-polled on SUCCESS only: a page that failed did not cover its
    // rows, and must not leave them looking covered.
    await stampLastPolled(env, page.rows.map((row) => row.id));
    cursor = page.nextCursor;
    stats.pages++;
  } while (cursor);

  // Phase 2 -- tail rotation: per-item fetches over what bulk missed. NULL
  // (never polled) sorts ahead of everything, then oldest last-polled first.
  const tail = await env.DB.prepare(
    `SELECT id FROM tracked_items
      ORDER BY last_polled_at IS NOT NULL, last_polled_at ASC
      LIMIT ?`,
  )
    .bind(TAIL_BATCH)
    .all<{ id: string }>();

  for (const row of tail.results) {
    if (budgetExhausted(budget)) {
      stats.truncated = true;
      break;
    }
    try {
      await upsertRows(env, [await fetchOne(env, row.id, budget)]);
    } catch (err) {
      console.warn(`[cron] tail fetch failed for ${row.id}`, err);
    } finally {
      // Stamp EVERY attempt here, success or failure. An item that fails
      // permanently and keeps its old stamp parks at the head of the rotation
      // forever and starves everything behind it.
      await stampLastPolled(env, [row.id]);
    }
    stats.tailFetched++;
  }

  return stats;
}

2 つのフェーズのあいだの非対称性は意図的なもので、そして間違えやすい。バルクが成功時にしかスタンプを打たないのは、そこでのスタンプがカバレッジの主張であり、失敗したページは何もカバーしていないからだ。末尾が失敗も含めてすべての試行にスタンプを打つのは、そこでのスタンプが公平性のカーソルだからだ。毎回必ず取得に失敗するアイテムは、そうしなければ最も古いスタンプを永遠に持ち続け、毎ティック真っ先に選ばれ、後ろにいる健全なアイテムの順番が回ってくる前に末尾の予算を食いつぶす。ここでは公平性がリトライの熱心さに勝る。失敗しているアイテムも、他のすべてと同じように次のローテーションでまた回ってくる。

予算のカウンターは毎ティック記録する。スイープが追いつかなくなったことを最も早く教えてくれるのがこの値だからだ。

console.log(
  `[cron] sweep pages=${stats.pages} tail=${stats.tailFetched} truncated=${stats.truncated}`,
);

1 回の打ち切りは正常だ。ティックの打ち切りが長く続いている場合は、ローテーションが一周しきれておらず、末尾が事実上監視されていないということになる。アラートを出す価値があるのはそちらの状態であって、個々の打ち切りではない。

鮮度スタンプ

黙って古びていくステータスボードは、遅れているのが目に見えるボードより性質が悪い。メッセージがたまたま描画された時刻ではなく、データを実際に取得した時刻を毎回の投稿に刻んでおく。

function buildStatusBlocks(rows: StatusRow[]): unknown[] {
  const syncedAt = new Date().toISOString();
  return [
    // ... table/section blocks built from rows ...
    {
      type: "context",
      elements: [{ type: "mrkdwn", text: `Last synced: ${syncedAt}` }],
    },
  ];
}

こうしておけば、ティックが途中で失敗しても(リースを奪われた、上流の取得がエラーになった)、前回のメッセージは古いタイムスタンプを保ったままになり、黙って最新のように見えることがない。ボードが古くなったことに気づく必要があるほど重要なら、「鮮度スタンプが N ティック分動いていない」というアラートと組み合わせるとよい。

レートティアを踏まえたバッチ化

多数の行に触れるティックを、1 行につき 1 回の Web API 呼び出しに変えてはいけない。slackLists.items.updatecells 引数は複数の行と列を 1 回の呼び出しにまとめられる。これは最適化ではなく仕組みそのものだ。items.update は Tier 3(およそ毎分 50 回以上)で動くため、数百行の変更に対して素朴に 1 行ずつループすると 1 ティックで予算を使い切り、429 を踏み始める。代わりに、変更されたセルを(ドキュメント化された cells の上限まで)1 回の呼び出しにまとめる。

// Batch every changed cell from this tick into as few calls as the
// documented per-call cap allows, rather than one call per row.
const CELLS_PER_CALL = 100;

async function syncListChanges(
  env: Env,
  changedCells: Array<{ row_id: string; column_id: string; select: string[] }>,
): Promise<void> {
  for (let i = 0; i < changedCells.length; i += CELLS_PER_CALL) {
    const batch = changedCells.slice(i, i + CELLS_PER_CALL);
    await callSlackApiWithRetry(
      "slackLists.items.update",
      { list_id: env.STATUS_LIST_ID, cells: batch },
      env.SLACK_BOT_TOKEN,
    );
  }
}

429 に対して Retry-After を尊重する話は、cron の文脈でもその外側とまったく同じように当てはまる。fetch による Web API 呼び出しを参照。

つまずきどころ

  • ctx.waitUntil() なしで返る scheduled() ハンドラーは、同期の途中で殺されうる。 3 秒 ack と同じルールで、違いは駆動する HTTP レスポンスがないことだけだ。実処理はきちんとラップすること。

  • リースは多重実行を抑制するだけで、防ぎはしない。 置き換えられた実行がまだ生きていて、まだ Slack を呼んでいることはありうる。ユーザーの目に触れる副作用にはその直前でリース喪失のチェックが要るし、データベースへの書き込みには run_id によるフェンスが要る。そして 2 回目の投稿を実際に止める仕組みは、依然としてアイテムごとのマーカーだけだ。

  • Date.now() はそのティックの時刻ではない。 日付やウィンドウから導かれるものはすべて controller.scheduledTime から取る。Date.now() は実行中の経過時間を測るためだけのものだ。

  • 予算で打ち切られたティックは成功だ。 アラートを出すのは打ち切りが何ティックも連続して続いていることに対してであって、1 回の打ち切りに対してではない。後者は仕組みが設計どおり働いている姿だ。

  • cron ジョブにはエラーを報告する相手がいない。 webhook ハンドラーと違い、失敗を持ち帰らせる HTTP レスポンスがない。scheduled() ハンドラーの内側で大きな声でログを出す(ボードが重要ならアラートも出す)こと。静かな失敗は、鮮度スタンプが動かなくなるという形でしか現れない。

  • * * * * *(毎分)が Cron Trigger の最小粒度。 ティックが日常的に長引くなら、毎分より細かくスケジュールしようとするのではなく、バッチか 1 回あたりのタイムアウトを小さくすること。

Revision History

作成更新