搜尋開發指南

試試「多租戶」、「啟動」或「outbox」。
搜尋僅使用本站文件,不會傳送至外部服務。

文件目錄

NEXUSPRO / HANDBOOK

Outbox、Jobs 與 Temporal

分辨單步冪等副作用、process-local 排程與長流程編排,避免把查詢或跨域交易放錯邊界。

內容核對 2026-09-28·圖文指南

本章目標

理解 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。

go
// 節錄自 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 結果完成狀態
回傳 nilsucceeded
errors.Is(err, outbox.ErrPermanentFailure)立即 dead_lettered,不消耗重試次數
其他錯誤、逾時或 panicfailed,1 秒起倍增、上限 5 分鐘後重試;次數用盡轉 dead_lettered
type 沒有對應 handlerparked,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_dispatch30s/10m一律(另由 LISTEN 喚醒)Dispatcher.Drain
iam_reclaim1m/2m一律收回到期 role binding 的授權並發 iam.grant_expired
attachment_reaper1m/2m一律逾期未完成的上傳標為失敗,並以刪除事件清 staging
agent_run_reaper1m/2m一律收斂被放棄的 queued/running agent run
attendance_maintenance1m/45s一律對帳考勤 dead-letter,每日檢查一次行事曆年度覆蓋
form_overdue_reminder1m/2m一律表單逾期提醒
attendance_source_sync5s/5mEHRMS_BASE_URL 與 EHRMS_API_KEY 皆非空消費已受理的考勤同步批次
ehrms_sync1m/30m同上認領各租戶到期的 30 分鐘來源同步(見平台適配層)
attendance_digest1m/5mATTENDANCE_DIGEST_ENABLED=true每日考勤摘要

重啟後由 durable table、cursor、lease 或 idempotency receipt 決定要不要再跑;runner interval 只是檢查頻率,例如來源同步真正的 30 分鐘節奏記在 PostgreSQL 的租戶政策列。

go
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:

  1. HR 員工交易發 hr.employee.created(對應入職),或在狀態轉為 offboarded 時發 hr.employee.status_changed(對應離職)。
  2. worker 的 HR handler 只在 WORKFLOW_JML_AUTO_START_ENABLED=true(預設 false)時,經 run service 建立 run 列,並在同一交易發 workflow.run.start_requested;flag 關閉時只記 info 並略過。
  3. worker 的 workflow.run.start_requested handler 呼叫 Orchestrator.Start(Temporal 端為 SignalWithStartWorkflow,workflow id 是 jml- 加 run id),成功後寫回 execution;重試用盡或永久失敗時 run 進 start_failed。
  4. 審批動作 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。

架構流程圖 · Mermaid

正在繪製架構圖…

查看圖表原始碼
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。

操作、驗證與常見錯誤

  1. 列出本次功能的「已發生事實、單步副作用、長流程等待」三欄。
  2. 對照 ProductionCatalog、handler registrations、cmd/worker/assemble.go 和對應 config;確認 flag off 時 baseline 不變。
  3. 為 handler 測重放、過期 lease、永久錯誤(outbox.ErrPermanentFailure)、retry exhaustion、tenant RLS 和 shutdown drain。
  4. 為 Temporal 測 start idempotence、signal routing、activity retry/non-retry、projection failure recovery;查詢則直接讀 PostgreSQL。
bash
# 工作目錄: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 與資料存取;交付門見部署與交付驗證。

內容來源與核實範圍

以下路徑相對於所列業務倉庫;核實層級為「已讀原始碼」,不是本輪業務測試或部署驗收。

  • nexus-pro-be-plus/internal/outbox/catalog.go
  • nexus-pro-be-plus/internal/outbox/handlers.go
  • nexus-pro-be-plus/internal/outbox/dispatcher.go
  • nexus-pro-be-plus/internal/outbox/prepare.go
  • nexus-pro-be-plus/internal/outbox/wake.go
  • nexus-pro-be-plus/internal/repository/outbox.go
  • nexus-pro-be-plus/internal/repository/postgres/outbox.go
  • nexus-pro-be-plus/db/queries/outbox_dispatch.sql
  • nexus-pro-be-plus/internal/service/hr/employees.go
  • nexus-pro-be-plus/internal/jobs/runner.go
  • nexus-pro-be-plus/internal/platform/temporal/client.go
  • nexus-pro-be-plus/internal/platform/temporal/worker.go
  • nexus-pro-be-plus/internal/platform/temporal/workflows/jml.go
  • nexus-pro-be-plus/internal/service/workflow/run_ports.go
  • nexus-pro-be-plus/internal/wiring/workflow.go
  • nexus-pro-be-plus/internal/worker/hr.go
  • nexus-pro-be-plus/internal/worker/workflow.go
  • nexus-pro-be-plus/cmd/worker/assemble.go
  • nexus-pro-be-plus/cmd/worker/ehrms_sync.go
  • nexus-pro-be-plus/cmd/worker/observability.go
  • nexus-pro-be-plus/.golangci.yml
NexusPro 開發指南以程式碼為準 · 以驗證為據