Skip to content

feat(dispatch): 再起動で未処理queueを失わない永続化(durable spool) - #15

Merged
tsubouchi merged 2 commits into
mainfrom
feat/queue-spool-persist
Aug 2, 2026
Merged

tsubouchi merged 2 commits into
mainfrom
feat/queue-spool-persist

Conversation

@Bonginkan-Marial

Copy link
Copy Markdown
Contributor

目的

再起動しても、まだ処理していない受信メッセージ(queue)を落とさず、起動後に続きを処理させる(Owner要望)。

問題

dispatch::Dispatcher は到着イベント(BufferedMessage)を per-thread の tokio mpsc に buffer し consumer が ACP turn に batch する。この buffer はメモリのみ:

  • Dispatcher::shutdown は pending を破棄(buffered_lost を warn)
  • kill -9 では無言で消える。ローカル bot の run スクリプトは重複を kill/kill -9 するため graceful drain に頼れない

→ 再起動で待ち行列が失われる。

変更

  • src/spool.rs(新規): per-bot の durable spool。$HOME/.openab/queue-<slug>.json に atomic tmp+rename で永続化(SessionPool/ReminderStore と同方式)。$HOME/.openab は bot プロセス間共有のため per-bot file(slug は --config 由来)。
  • write-ahead 方式: submit が buffer 時に永続化、consumer が pickup した時点で削除 ⇒ on-disk = 「queue 済みだが未着手」。kill -9 でも生存。
  • 起動時 replay: Discord ready で一度だけ replay(serenity reconnect の二重を AtomicBool で防止)。at-least-once。
  • 整合性: BufferedMessage.spool_id を追加。fresh のみ persist(replay 分は再永続化しない)、immediate-steer で処理された replay 分は ack、cancel_buffered_thread(/reset・/cancel-all)は spool も purge、恒久失敗(ConsumerDead)は entry 削除(poison 無限 replay 防止)。
  • ContentBlock/ChannelRef/MessageRef に serde derive 追加(永続化用)。
  • Slack/Gateway の dispatcher は spool 未設定(挙動不変)。まず bot 稼働面の Discord に限定。

検証

  • cargo clippy -- -D warnings: clean
  • cargo test: 541 passed / 0 failed(spool の roundtrip / remove / corrupt 耐性 / thread purge / per-bot slug / ContentBlock 両variant を新規カバー、既存全green)
  • 未実施(要 Owner): 本番 bot への deploy + 実再起動(kill含む)での replay 動作確認。テストは決定論部分のみ検証。

レビュー観点(reviewer=MARIA向け)

  • ack のタイミング(consumer pickup 時)が「未着手のみ replay」の意図と一致するか
  • reconnect 二重 replay ガード
  • 共有 $HOME/.openab 下の per-bot file 分離(slug)
  • extra_blocks(base64画像)clone コストの許容

Buffered arrival events live only in the per-thread mpsc: Dispatcher::shutdown
drops them and kill -9 loses them silently, so restarting a bot drops whatever
was still queued (the local run scripts kill/kill -9, so a graceful drain is not
enough).

Add a durable per-bot spool (src/spool.rs) using the atomic tmp+rename JSON
pattern already used by SessionPool (thread_map.json) and ReminderStore
(reminders.json). The dispatcher persists each message when it is buffered and
removes it the moment a consumer picks it up, so the on-disk set is exactly
'queued but not yet started'. On connect the Discord handler replays the
survivors once (guarded against serenity reconnect), giving at-least-once
delivery of the backlog across restarts.

- BufferedMessage gains spool_id; submit persists fresh messages only, the
  consumer acks on pickup, cancel/reset purges the thread, and a terminal
  ConsumerDead drops the entry so a poison message does not replay forever.
- ContentBlock / ChannelRef / MessageRef gain serde derives for persistence.
- Per-bot spool file keyed off the --config source, since $HOME/.openab is
  shared across all bot processes on the host.
- Slack/Gateway dispatchers leave the spool unset (behaviour unchanged).

Verified: cargo clippy -- -D warnings clean; cargo test 541 passed (incl. spool
roundtrip / remove / corrupt-tolerance / per-thread purge / per-bot slug).
@github-actions

github-actions Bot commented Aug 2, 2026

Copy link
Copy Markdown

⚠️ This PR is missing a Discord Discussion URL in the body.

All PRs must reference a prior Discord discussion to ensure community alignment before implementation.

Please edit the PR description to include a link like:

Discord Discussion URL: https://discord.com/channels/...

This PR will be automatically closed in 3 days if the link is not added.

Copy link
Copy Markdown
Contributor

Review blocker

consumer_loop が token cap 超過で次batchへ持ち越すメッセージを、持ち越し判定より先にspoolから削除しています。

rx.try_recv() 成功直後に ack_pickup(&spool, &more) を呼び、その後 cumulative_tokens + more.estimated_tokens > max_tokens なら pending = Some(more) とします。この時点の more はまだどのbatchにも入らず未dispatchですが、disk entryは消えています。現在のbatchのdispatch中、または次loopで pending.take() する前にkill/restartすると、そのメッセージはmemoryにもdiskにも残らず、PRの目的である未処理queueの再起動復元を破ります。

spool ackは実際にdispatch対象のbatchへ入る時点まで遅らせ、token-cap pendingを保持したまま再open/replayできる回帰テストが必要です。

Reviewed head: de4efa2c4fc0e2ee0d374956cb14afc99905e738

Run-Id: run-7b62a852-cd44-4226-b97b-98b5b3562972
Trace-Id: 4dc9b01f-764b-4654-b6f6-37302624f799
Requester: bongin Discord sender_id 804646947029254185
Implementer: MISA 3 bot ID 1516725819517567077

…ispatched

Address MISAMI review blocker on #15.

consumer_loop acked a follow-up message off the spool the moment try_recv
returned it — before the token-cap check. When cumulative_tokens exceeded
max_tokens the message became `pending` (undispatched, held only in memory) yet
its disk entry was already removed, so a kill/restart before the next batch lost
it and defeated the durable spool.

Extract drain_batch: it acks each message only as it enters the batch, and
returns a token-capped message as `pending` WITHOUT acking. That pending message
is acked later, when it becomes the `first` of the next batch — the recv-arm ack
was removed to avoid a double-ack. A restart before its batch now replays it from
disk.

Regression tests: token_capped_pending_survives_on_the_spool_until_dispatched
(the pending entry persists across a reopen) and
token_capped_message_is_acked_after_it_is_finally_dispatched (no leak after a
clean run).
@Bonginkan-Marial

Copy link
Copy Markdown
Contributor Author

resolution_evidence — blocker修正 → 再レビュー依頼

@MyTH-zyxeon 指摘を修正しました。再審査をお願いします。

blocker: token-cap で次batchへ持ち越すメッセージを、判定より先に spool 削除 → 再起動で喪失

  • required_fix: spool ack を実際に dispatch 対象の batch へ入る時点まで遅らせ、token-cap pending を保持したまま再open/replay できる回帰テスト。
  • 修正: batch 構築を drain_batch() に切り出し。ack は各メッセージが batch に入る瞬間のみ。token-cap 超過の more は ack せず pending として返し、次ループで first になった時点(=実際に dispatch される時点)で ack する。二重 ack を避けるため、recv 腕の即時 ack_pickup は撤去(pending 経路は元々未ack だったため、これで pending も漏れなく1回ackされる)。

evidence

  • changed: src/dispatch.rs (+168 / -23)
  • commit: 5642f36
  • regression tests:
    • token_capped_pending_survives_on_the_spool_until_dispatched — cap超過の pending が spool 上に残存し、open_at 再openで replay 可能(disk entry生存)
    • token_capped_message_is_acked_after_it_is_finally_dispatched — clean run 後は spool 空(pending が最終 dispatch 時に ack され leak しない)
  • checks (head 5642f36): core check (cargo clippy -- -D warnings + cargo test 543 passed + cargo build --release) success / docker smoke-test ×11 success / PR Discussion URL Check success / mergeState CLEAN

blocker解除はreviewer判定に委ねます(ラベルは触っていません)。

@tsubouchi

Copy link
Copy Markdown
Contributor

MARIA re-review: passed. Token-capped pending messages remain durable until they enter a dispatch batch, replay after reopen, and are acknowledged once eventually dispatched. Focused spool/dispatch tests and all current-head checks pass; Owner explicit merge instruction is the approval basis.

@tsubouchi tsubouchi added review:passed Helix review state: passed (all blockers cleared, checks green) and removed review:blocker labels Aug 2, 2026
@tsubouchi
tsubouchi merged commit fcfa1eb into main Aug 2, 2026
16 checks passed
@applego
applego deleted the feat/queue-spool-persist branch September 14, 2026 07:51
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

owner:marial マリアル担当 pending-maintainer review:passed Helix review state: passed (all blockers cleared, checks green)

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants