本章目標
理解 PostgreSQL outbox、internal/jobs runner 與 Temporal 的責任分工,能為一次副作用選擇可恢復、可重試且可觀測的邊界。本文是依現有程式碼整理的架構指南,不宣稱所有業務流程都已接通或已完成部署驗收。
前置條件
先讀技術架構與Repository 與資料存取,熟悉 tenant transaction、唯一鍵、JSON event payload 及 Go context。開發前先查 production event catalog、handler registry、config flag 和既有 worker wiring,不能自行新增第二份事件表。
三種工作形狀
| 形狀 | 使用 | 持久化事實 | 不該做 |
|---|---|---|---|
| outbox event | 單步、可重試、冪等副作用 | PostgreSQL event、attempt、lease、receipt | 在 API transaction 內直接呼叫外部系統 |
| Jobs runner | 固定週期或重建後可再次執行的排程 | job 自己的 durable source;runner 本身不保存業務狀態 | 把 runner 當可靠 queue |
| Temporal workflow | 等待、人審、重試、補償的長流程 | Temporal history + PostgreSQL domain projection | 用 workflow 當業務查詢資料庫 |
PostgreSQL 是業務資料、查詢投影和 outbox 的唯一事實來源。長流程交給 orchestrator 編排(目前只有 JML 入職/離職兩個定義),但狀態和列表查詢一律讀本地投影;表單簽核是同步的 PostgreSQL 狀態機,不經 Temporal。
Outbox 發佈與投遞
業務 service 透過本域 repository 的 unit of work(WithinTenant)開 tenant transaction,在同一個 callback 內寫入事實,並呼叫 tx 內嵌的窄 outbox.Publisher(repository.OutboxTx);transaction commit 後事件才可見,真正呼叫 adapter 在 transaction 外完成。這避免「資料成功但事件遺失」,也避免外部副作用反過來讓 domain transaction 長時間持鎖。WithTenantTx 屬於 internal/platform/postgres,由 repository 層的 unit of work 呼叫;depguard 規則 service-cannot-import-platform 禁止 internal/service/** import internal/platform/**。
internal/outbox/catalog.go 是 production event type 唯一清單;internal/outbox/handlers.go 的 NewHandlerRegistry 要求 catalog 與 handler 精確 1:1。事件 type 使用 exact key lookup,不能用 prefix、substring 或 wildcard 推測 handler。
// 節錄自 internal/service/hr/employees.go(省略驗證、稽核與錯誤映射)
err = governor.employees.WithinTenant(ctx, tenantID, func(ctx context.Context, tx repository.HREmployeeGovernanceTx) error {
created, err := tx.InsertEmployee(ctx, employee)
if err != nil {
return err
}
eventID, err := governor.newID()
if err != nil {
return err
}
// HREmployeeGovernanceTx 內嵌 repository.OutboxTx,也就是窄 outbox.Publisher
return tx.Publish(ctx, outbox.Event{
Type: outbox.EventTypeHREmployeeCreated,
AggregateID: &created.ID,
IdempotencyKey: string(outbox.EventTypeHREmployeeCreated) + ":" + eventID.String(),
Payload: map[string]any{"employee_id": created.ID.String(), "account_id": created.AccountID},
})
})Publish 先經 outbox.Preparer 正規化:type 不在 catalog 回 ErrUnknownEventType;IdempotencyKey 空白或 Payload 為 nil 回 ErrInvalidEvent,整筆 transaction 跟著失敗;PayloadVersion、MaxAttempts 省略時預設 1、5。接著在同一 transaction 依序寫 outbox_events、upsert 該租戶的 outbox_dispatch_queue,再以 pg_notify 喚醒 worker(commit 後才送達);(tenant_id, event_type, idempotency_key) 有唯一索引。
投遞、重試與失敗分類
worker 以專用連線 LISTEN nexus_outbox_wake 接收喚醒,並以 30 秒輪詢兜底;啟動時先同步掃描一次,補上 LISTEN 建立前的空窗。dispatcher 每輪先 claim 一個到期的 tenant queue 列,再 bounded claim 該租戶的事件(worker 設定每批 2 筆),逐筆執行 handler:
| handler 結果 | 完成狀態 |
|---|---|
| 回傳 nil | succeeded |
errors.Is(err, outbox.ErrPermanentFailure) | 立即 dead_lettered,不消耗重試次數 |
| 其他錯誤、逾時或 panic | failed,1 秒起倍增、上限 5 分鐘後重試;次數用盡轉 dead_lettered |
| type 沒有對應 handler | parked,metric label 為 unknown |
事件的 completion 以 event 的 claim_token 加 status = 'processing' 圍欄;tenant queue 列的 reschedule/delete 另以 queue 的 claim_token 加 generation 圍欄。過期 worker 晚到的 completion 不會套用,只記為 stale。其他現行參數:單次 attempt 逾時 200 秒、event lease 8 分鐘、queue lease 9 分鐘,都是 cmd/worker/assemble.go 的常數。
每個 handler 都要以 event idempotency key、effect receipt 或資料庫唯一約束防止重放。永久錯誤要回傳(或包裝)outbox.ErrPermanentFailure,不能靠重試「等它變好」;可恢復的 upstream timeout 才依 bounded retry policy 再試。事件 payload 中不要放密碼、token 或不必要的個資。
可觀測性與就緒檢查
worker 在 METRICS_ADDR(預設 127.0.0.1:9091)提供 GET /healthz、GET /readyz 與 GET /metrics。outbox 指標是 outbox_dispatch_total(label event_type、outcome)與 outbox_backlog(pending 加 failed 事件數,每次 scrape 即時查詢);job 指標是 worker_job_duration_seconds(label job、outcome)。job label 是封閉集合,目前 attendance_source_sync 與 attendance_digest 會記成 unknown。dispatch 啟用時,outbox_dispatch_queue、outbox_listen 與 objectstore 是 critical readiness check,啟動時任一失敗 worker 就拒絕啟動。
Jobs runner
internal/jobs/runner.go 是 process-local scheduling loop:第一次 invocation 先執行,再按 interval tick;同一 runner 透過 single-flight gate 避免 tick 與手動 Trigger 同時進入;每次 invocation 有 timeout,逾時的 callback 真正返回前不會准入下一輪;shutdown 用 Stop + Drain 等待已准入工作。
runner 不會代替 durable source。所有 runner 都在 cmd/worker/assemble.go 的 assembleDispatch 內建立,所以 OUTBOX_DISPATCH_ENABLED=false(預設 true)時 worker 不跑任何 job:
| job | 週期/逾時 | 建立條件 | 工作 |
|---|---|---|---|
outbox_dispatch | 30s/10m | 一律(另由 LISTEN 喚醒) | Dispatcher.Drain |
iam_reclaim | 1m/2m | 一律 | 收回到期 role binding 的授權並發 iam.grant_expired |
attachment_reaper | 1m/2m | 一律 | 逾期未完成的上傳標為失敗,並以刪除事件清 staging |
agent_run_reaper | 1m/2m | 一律 | 收斂被放棄的 queued/running agent run |
attendance_maintenance | 1m/45s | 一律 | 對帳考勤 dead-letter,每日檢查一次行事曆年度覆蓋 |
form_overdue_reminder | 1m/2m | 一律 | 表單逾期提醒 |
attendance_source_sync | 5s/5m | EHRMS_BASE_URL 與 EHRMS_API_KEY 皆非空 | 消費已受理的考勤同步批次 |
ehrms_sync | 1m/30m | 同上 | 認領各租戶到期的 30 分鐘來源同步(見平台適配層) |
attendance_digest | 1m/5m | ATTENDANCE_DIGEST_ENABLED=true | 每日考勤摘要 |
重啟後由 durable table、cursor、lease 或 idempotency receipt 決定要不要再跑;runner interval 只是檢查頻率,例如來源同步真正的 30 分鐘節奏記在 PostgreSQL 的租戶政策列。
runner, err := jobs.NewRunner(jobs.Options{
Name: "attendance_source_sync", Interval: 5 * time.Second,
Timeout: 5 * time.Minute, Run: syncRun, Metrics: runnerMetrics,
Logger: jobLogger{logger: logger}, // 實作 jobs.Logger 的 LogRunFailure,不是 *slog.Logger
})Run 的單輪應可在 timeout/取消下安全結束;批次工作要逐筆分類失敗,不讓一筆髒資料鎖死所有租戶。固定 interval、超時與 retry budget 要放在 wiring/config,而非埋在 domain 計算裡。
Temporal 分界
internal/platform/temporal/client.go 的 client 只實作 workflowservice.Orchestrator 的 Start/Signal 兩個寫方法(Start 讀 execution metadata 只為產生啟動回執);port 刻意沒有讀方法,由 tests/contract/workflow_orchestrator_port_test.go 守住,它沒有被當成業務 read model。TEMPORAL_ENABLED=true 時,StartWorker 在 worker process 註冊 JML workflow 和 activities,透過狹窄 StageWriter/RunWriter port 轉譯 Temporal 詞彙,業務規則仍留在 service。
JML 啟動與動作鏈
JML run 不由 API 直接啟動,而是經過兩段 outbox:
- HR 員工交易發
hr.employee.created(對應入職),或在狀態轉為 offboarded 時發hr.employee.status_changed(對應離職)。 - worker 的 HR handler 只在
WORKFLOW_JML_AUTO_START_ENABLED=true(預設false)時,經 run service 建立 run 列,並在同一交易發workflow.run.start_requested;flag 關閉時只記 info 並略過。 - worker 的
workflow.run.start_requestedhandler 呼叫Orchestrator.Start(Temporal 端為SignalWithStartWorkflow,workflow id 是jml-加 run id),成功後寫回 execution;重試用盡或永久失敗時 run 進start_failed。 - 審批動作
POST /v1/workflow-runs/:workflowRunId/approve、reject、return、withdraw在 API 交易內寫workflow_actions動作證據,交易外只呼叫Orchestrator.Signal(以signal_sent_at判斷是否需要補送);列表與明細由GET /v1/workflow-runs讀 PostgreSQL 投影。
目前 JML workflow 的實際控制流是逐關 ActivateStage → 等 workflow.stage.decision signal → RecordStageDecision → approve 湊滿本關才前進、reject 與 withdraw 走 compensation、return 讓 run 標為 returned 並留在原關;notify 關卡在活化交易內直接結掉。第一關無法活化時,以持續重試的投影 activity 把 run 寫成 start_failed;後續關卡活化受阻時以 durable timer 退避重試(1 分鐘起倍增、上限 1 小時),等待期間仍接受撤回。activity 把 permanent error、blocked 結果與 retryable transport failure 分開。這是已讀取的 workflow 實作,不等於所有 HR 流程都已採 Temporal。
正在繪製架構圖…
查看圖表原始碼
flowchart LR
API["Gin command"] --> DB[("PostgreSQL facts / projection / outbox")]
DB -- "NOTIFY 或 30s poll" --> W["worker dispatcher"]
W --> A["platform adapter"]
W -- "start_requested 後呼叫 Start" --> T["Orchestrator:Temporal 或 in-process fake"]
API -- "審批動作只送 Signal" --> T
T --> TW["Temporal activities(worker)"]
TW --> DB
UI["Query"] --> DB跨域副作用不要在同一個 tenant transaction(WithinTenant callback)內同步呼叫另一個 service。需要通知、HR 狀態收斂或外部 adapter 時,事件在 transaction 內投遞,consumer 以自己的 transaction 寫自己的 domain。
操作、驗證與常見錯誤
- 列出本次功能的「已發生事實、單步副作用、長流程等待」三欄。
- 對照
ProductionCatalog、handler registrations、cmd/worker/assemble.go和對應config;確認 flag off 時 baseline 不變。 - 為 handler 測重放、過期 lease、永久錯誤(
outbox.ErrPermanentFailure)、retry exhaustion、tenant RLS 和 shutdown drain。 - 為 Temporal 測 start idempotence、signal routing、activity retry/non-retry、projection failure recovery;查詢則直接讀 PostgreSQL。
# 工作目錄:nexus-pro-be-plus;依來源核實,本輪未跑業務 worker/Temporal
go test ./internal/outbox ./internal/jobs ./internal/platform/temporal/... -count=1
go test ./tests/contract -run 'Outbox|Orchestrator|Temporal' -count=1
# 連真 Temporal frontend 的 adapter 測試在 tests/integration/postgres/workflow_temporal_test.go,
# 屬於 make integration,需要實際依賴環境預期結果: event catalog 與 handlers 精確對齊;重放不重複副作用;一筆 event 失敗不阻擋可獨立處理的事件;Temporal 重啟後由自身 history 重放恢復,activity 寫投影保持冪等;UI 不向 Temporal 查列表。
常見錯誤:把外部 HTTP 放 API transaction、在 service import internal/platform/postgres 自開 transaction、發佈時漏掉 IdempotencyKey 或 Payload、把所有事情包成 workflow、在 API 直接呼叫 Orchestrator.Start、用 in-memory channel 當 outbox、以 event prefix 路由、沒有 dead-letter 證據、用 runner interval 取代 durable cursor、或把「TEMPORAL_ENABLED=true」當成流程已驗收(預設 false 時跑的是 in-process fake)。
相關文件
外部服務整合見平台適配層;JML run 的 API 與畫面見工作流與審批;資料交易見Repository 與資料存取;交付門見部署與交付驗證。