本日程式碼:repo tag day-09
昨天的 Worker 其實已經可以同時跑好幾個:它自己沒有狀態,八個 goroutine 對同一個 Worker 呼叫 RunOnce,五十筆 intent 還是每一筆剛好送一次。「多個 worker 不會把同一筆錢送兩次」這件事,queue 的 lease 加上 worker 每一步冪等早就給了,今天不用再證明一次。今天要處理的是剩下的三件事:工作怎麼分給 N 個 worker、收工的時候手上還沒送完的交易怎麼辦、以及 N 個 worker 不等於 N 條連線打在 RPC 上。
我們今天要做一個 Pool,把昨天的一個 worker 變成 N 個 goroutine 共用同一條 queue,再加一個擋在 worker 前面的 Throttle。沿用的公開設計有三個:
golang.org/x/time/rate 的 token bucket(這裡自己寫了四十行的版本,backend 到今天還是零外部依賴)repo 中 internal/relayer 的 Example_poolDrain 會排六份 job、起兩個 worker,然後在兩筆交易正在送的時候發出停止訊號。輸出大概長這樣:
enqueue 6 job(s); pool of 2 workers, 2 sends in flight at most
send [pi_0001 pi_0002] in flight
stop requested: stop leasing, let the in-flight sends finish
pool stopped: sent 2, retry 0, abandoned 0, panics 0
queue 4 job(s) left, next lease ok=true
pi_0001 confirming v4 tx=0x0001
pi_0002 confirming v4 tx=0x0002
pi_0003 authorized v2 tx=
pi_0004 authorized v2 tx=
pi_0005 authorized v2 tx=
pi_0006 authorized v2 tx=
Run 才回來。next lease ok=true),下一個起來的程序直接接著做。Stats 把這次跑的結果分成 sent、retry、abandoned、panics 幾個數字,abandoned 是收工逾時被放棄的那幾份,這篇下面會講到。講到 Go 的 worker pool,我第一個想到的是 Go by Example 那一頁:開一條 jobs channel,起 N 個 goroutine 用 range 從裡面拿。套到我們這裡最直覺的寫法是:一個 dispatcher 從 queue 一次 Lease 十份,塞進 channel,worker 再從 channel 拿。我沒有這樣做:
queue 本身就是那條 channel。job 一被 Lease,對 queue 來說它就是「已領走」:lease 的時鐘開始走、別的程序看不到它。它如果只是躺在某個 goroutine 的 channel buffer 裡排隊,等於用掉了 lease 卻沒人在做;dispatcher 一死,整批要等 lease 過期才回得來。每個 worker 自己領、手上永遠只有一份,慢的 worker 自然少領、快的多領,不需要有人在中間分配。代價是 queue 空的時候 N 個 worker 會一起輪詢,接 SQS 這種 long polling 的 queue 沒差,接資料庫 table 就是每個 Idle 多 N 次 SELECT,我覺得可以接受。
另一個順便一起做的決定是 panic。Go 的 panic 會帶走整個程序,一份 job 的 Sender 有 bug,不該讓另外 N-1 個 worker 陪葬。pool 在每一圈外面 recover,接住之後什麼都不做:不 Ack、不 Nack,job 留在 lease 裡,過期後自然回到 queue,重來的 worker 照 intent 現在的狀態處理。這裡絕對不能 Ack,因為 panic 的那一步不知道做到哪,Ack 等於把這份工作當做完。
Run 收到停止訊號之後做的事,跟 http.Server.Shutdown 是同一個形狀:先關入口、等在跑的做完、超時才放棄。
為什麼不乾脆把 ctx 傳進去、一取消全部停?因為一筆被取消到一半的 Send,回來的結果是「不知道送出去了沒」。照昨天的規則,這筆 intent 會停在 settling 沒有 tx hash,重來的 worker 不敢重送,五分鐘後送審。一次部署就製造一批 needs_review,比多等幾秒貴太多。所以 worker 處理 job 用的是另一個 workCtx(context.WithoutCancel 切出來的),外面的 ctx 結束只代表「不要再領新的」。
那為什麼還要有 DrainTimeout?因為 grace period 過後程序反正會被 KILL(Kubernetes 預設給 30 秒)。與其被 KILL 在不知道哪一行,不如自己挑在哪一步停:時間到就取消 workCtx,被打斷的 Send 以 retry 收場、job Nack 回 queue、intent 停在 settling,然後把放棄了幾份記進 Stats.Abandoned。放棄是有意識的、有數字的,跟被 KILL 是兩回事。DrainTimeout 要比一次 Send 的正常耗時長、比部署系統的 grace period 短,這裡設 20 秒。
「不把 RPC 打爆」聽起來是 Sender 的事:在 Send 裡面包一個 semaphore,同時最多幾筆在送。我一開始也是這樣想的,但它擋錯位置了:
昨天的順序是固定的:先 settling、再 hold、再廣播。名額如果擋在 Send 裡,被擋住的 job 已經是 settling 加 hold,而重來的 worker 在 settling 那一格是不重送的,一個單純的限流等待就會把一筆好好的 intent 送進 needs_review。所以名額要在任何副作用之前拿:
// 先跟 limiter 要一個送出的名額,拿到之前一個 byte 都不寫:這時放手最便宜,
// intent 還在 authorized、帳上沒有 hold,job Nack 回 queue 就好。
// 等待的上限是這份 lease 剩下的時間:lease 過期後 job 會被別人領走,再等下去也輪不到我們做。
actx, cancel := context.WithTimeout(ctx, d.LeaseUntil.Sub(w.now()))
err := w.limiter.Acquire(actx)
cancel()
if err != nil {
return Report{Outcome: OutcomeRetry, Detail: "throttled: " + err.Error()}, nil
}
defer w.limiter.Release()
被擋住的 job 原封不動回 queue:intent 還在 authorized、帳上沒有 hold、Send 沒被叫,等於這次交付沒發生過。代價是等名額的時候占著一份 lease,所以等待的上限就是 lease 剩下的時間,超過就放手,反正 job 會被別人領走。
Throttle 有兩個旋鈕,因為 RPC 怕的是兩件不同的事:節點怕同時太多條連線,服務商限的是每秒多少請求(Alchemy 算的是 compute units per second,超過就回 429)。MaxInFlight 是 Effective Go 那個 buffered channel semaphore;PerSecond 與 Burst 是 token bucket,拿不到 token 就照缺口算要等多久。Example_throttle 設每秒 2 筆、burst 1,輸出長這樣:
acquire #1 t+0s no wait
acquire #2 t+0s waited 500ms
acquire #3 t+500ms waited 500ms
acquire #4 t+1s waited 500ms
N 個 worker 共用同一個 Throttle,所以名額限的是整個 pool 對外的總量:八個 worker、三個名額,是很正常的組合,worker 數與 RPC 連線數就是兩個數字。
今天都在定義 N 個 worker 怎麼共存:每個 worker 自己領、手上只有一份,queue 就是 channel,不放 dispatcher;收工分兩段,先停止領新的、讓在送的送完,逾時才有意識地放棄並記下數字;限流擋在 settling 之前,被擋住的 job 原封不動回 queue,名額限的是整個 pool 對 RPC 的總量。這三個決定都沒有動到昨天的順序,改的只有「什麼時候可以放手」:放手要在副作用之前,不然就得等到副作用結束。
明天我打算讓 relayer 真的開始發送交易。我覺得第一個碰到的問題一定是同一個錢包送出的交易要怎麼排隊,畢竟我們要顧及到 EVM 的 nonce、Solana 的 blockhash、TON 的 seqno、SUI 的 object version 這四條鏈的四套做法。
明天見。