บทที่ 7 · Part 2 — Building
Signal, Query, Update
ยืนยัน OTP, ขอ approval จากทีม compliance, อ่านสถานะแบบ real-time และ Update-with-Start
Workflow ไม่ได้เป็นแค่สคริปต์ที่รันแล้วจบ มันคุยกับโลกภายนอกได้ระหว่างทาง — รอ OTP จากผู้ใช้, รับ callback จากธนาคาร, ให้ compliance กดอนุมัติ, และตอบคำถามว่า "ตอนนี้ธุรกรรมถึงไหนแล้ว"
จบบทนี้คุณจะ
- เลือกได้ว่าเคสไหนใช้ Signal, Query หรือ Update
- เขียน flow ยืนยัน OTP ที่มี timeout จริง
- ใช้ Update-with-Start เพื่อลด round trip
- รู้กับดักของ handler ที่ยังทำงานค้างตอน workflow จบ
เลือกให้ถูกตัว
| Query | Signal | Update | |
|---|---|---|---|
| แก้ state ได้ | ไม่ | ได้ | ได้ |
| คืนค่ากลับ | ได้ | ไม่ | ได้ |
| รอ/บล็อกได้ | ไม่ | ได้ | ได้ |
| บันทึกใน history | ไม่ | บันทึก | บันทึก |
| มี validator | ไม่ | ไม่ | มี |
| ใช้กับ | อ่านสถานะปัจจุบัน | แจ้งเหตุการณ์แบบ fire-and-forget | สั่งงานที่ต้องรู้ผลทันที |
จำสั้นๆ: Query เพื่อแอบดู · Signal เพื่อผลัก · Update เพื่อสั่งแล้วรอคำตอบ
Signal — ยืนยัน OTP
เคสจริง: เติมเงินยอดสูงต้องให้ผู้ใช้กรอก OTP ภายใน 5 นาที ถ้าไม่ทันให้ยกเลิกรายการ
type OTPSubmission struct {
Code string
Submitted time.Time
}
func TopUpWithOTPWorkflow(ctx workflow.Context, req TopUpRequest) (TopUpResult, error) {
logger := workflow.GetLogger(ctx)
actCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: 30 * time.Second,
})
var a *activities.Activities
if err := workflow.ExecuteActivity(actCtx, a.SendOTP, req.WalletID).Get(actCtx, nil); err != nil {
return TopUpResult{}, err
}
otpCh := workflow.GetSignalChannel(ctx, "otp-submitted")
// รอ OTP สูงสุด 5 นาที — Selector คือวิธี "รออันไหนมาก่อน" ที่ replay-safe
var submission OTPSubmission
received := false
timerCtx, cancelTimer := workflow.WithCancel(ctx)
timeout := workflow.NewTimer(timerCtx, 5*time.Minute)
selector := workflow.NewSelector(ctx)
selector.AddReceive(otpCh, func(c workflow.ReceiveChannel, _ bool) {
c.Receive(ctx, &submission)
received = true
cancelTimer() // ไม่ต้องรอ timer แล้ว
})
selector.AddFuture(timeout, func(f workflow.Future) {
if f.Get(ctx, nil) == nil {
logger.Info("otp timed out", "txnID", req.TransactionID)
}
})
selector.Select(ctx)
if !received {
return TopUpResult{}, temporal.NewNonRetryableApplicationError(
"OTP timeout", "OTPTimeout", nil)
}
// ตรวจ OTP ใน activity เพราะต้องคุยกับระบบภายนอก
var valid bool
err := workflow.ExecuteActivity(actCtx, a.VerifyOTP, req.WalletID, submission.Code).
Get(actCtx, &valid)
if err != nil {
return TopUpResult{}, err
}
if !valid {
return TopUpResult{}, temporal.NewNonRetryableApplicationError(
"OTP incorrect", "OTPInvalid", nil)
}
return continueTopUp(ctx, req)
}
ส่ง signal จากฝั่ง client:
err := c.SignalWorkflow(context.Background(),
"topup-"+txnID, "", // workflow ID, run ID (ว่าง = run ปัจจุบัน)
"otp-submitted",
OTPSubmission{Code: "123456", Submitted: time.Now()},
)
หรือจาก CLI ตอน debug:
temporal workflow signal \
--workflow-id topup-TXN-001 \
--name otp-submitted \
--input '{"Code":"123456"}'
ทำไมต้อง cancelTimer
ถ้าไม่ยกเลิก timer มันจะค้างอยู่ใน history จนครบ 5 นาที
ทำให้ workflow ปิดตัวไม่ได้ทันทีและกิน resource โดยไม่จำเป็น
pattern WithCancel + cancelTimer() นี้ควรใช้ทุกครั้งที่ตั้ง timeout คู่กับ signal
รับหลาย signal พร้อมกัน
ถ้าต้องฟังหลาย channel ตลอดอายุ workflow ให้แยกไป goroutine ของ Temporal:
func WalletSessionWorkflow(ctx workflow.Context, walletID string) error {
var pending []Instruction
closed := false
workflow.Go(ctx, func(gctx workflow.Context) {
addCh := workflow.GetSignalChannel(gctx, "add-instruction")
closeCh := workflow.GetSignalChannel(gctx, "close-session")
for {
sel := workflow.NewSelector(gctx)
sel.AddReceive(addCh, func(c workflow.ReceiveChannel, _ bool) {
var ins Instruction
c.Receive(gctx, &ins)
pending = append(pending, ins)
})
sel.AddReceive(closeCh, func(c workflow.ReceiveChannel, _ bool) {
c.Receive(gctx, nil)
closed = true
})
sel.Select(gctx)
if closed {
return
}
}
})
// รอจนกว่าจะมีงานเข้ามา หรือ session ถูกปิด
return workflow.Await(ctx, func() bool { return len(pending) > 0 || closed })
}
Signal อาจมาถึงก่อนที่ workflow จะพร้อมรับ
Temporal เก็บ signal ที่ส่งเข้ามาไว้ใน channel buffer ให้เสมอ
จึงไม่หายแม้ workflow ยังไม่ได้ Receive แต่ ถ้า workflow จบไปแล้ว signal จะหาย
เรื่องนี้สำคัญมากตอนทำ Continue-As-New (ดูบทที่ 8)
Query — อ่านสถานะแบบ real-time
Query ใช้ตอบคำถามจากหน้าจอลูกค้าหรือ dashboard ของ ops โดยไม่แตะ history
type TopUpStatus struct {
Stage string // "checking" | "awaiting_otp" | "charging_bank" | "crediting" | "done"
BankRef string
Attempts int
}
func TopUpWorkflow(ctx workflow.Context, req TopUpRequest) (TopUpResult, error) {
status := TopUpStatus{Stage: "checking"}
err := workflow.SetQueryHandler(ctx, "status", func() (TopUpStatus, error) {
return status, nil
})
if err != nil {
return TopUpResult{}, err
}
// อัปเดต status ระหว่างทาง
status.Stage = "awaiting_otp"
// ...
status.Stage = "charging_bank"
// ...
status.Stage = "done"
return TopUpResult{}, nil
}
temporal workflow query --workflow-id topup-TXN-001 --name status
ข้อห้ามของ Query handler
- ห้ามแก้ state — จะทำให้ replay ไม่ตรงกัน
- ห้ามบล็อก — เรียก activity,
workflow.Sleep,workflow.Awaitไม่ได้ทั้งหมด
query handler ต้องอ่านตัวแปรแล้ว return ทันทีเท่านั้น ถ้าอยากทำอะไรที่มี side effect ให้ใช้ Signal หรือ Update
Query ยังทำงานได้กับ workflow ที่จบไปแล้ว ตราบใดที่ history ยังอยู่ใน retention ซึ่งมีประโยชน์มากตอนสืบสวนเคสย้อนหลัง
Update — สั่งแล้วรอผล
Update คือ Signal + Query รวมกัน พร้อม validator ที่ปฏิเสธคำสั่งได้ก่อนบันทึกลง history
เคสจริง: ทีม compliance ขออนุมัติปรับ limit ให้ลูกค้าระหว่างที่ธุรกรรมค้างอยู่
func TopUpWorkflow(ctx workflow.Context, req TopUpRequest) (TopUpResult, error) {
approvedLimit := int64(0)
err := workflow.SetUpdateHandlerWithOptions(
ctx,
"approve-override",
// handler — แก้ state ได้ และคืนค่ากลับได้
func(ctx workflow.Context, in OverrideRequest) (OverrideResult, error) {
approvedLimit = in.NewLimitSatang
return OverrideResult{
Approved: true,
EffectiveAt: workflow.Now(ctx),
}, nil
},
workflow.UpdateHandlerOptions{
// validator — อ่านอย่างเดียว ห้ามแก้ state ห้ามบล็อก
Validator: func(ctx workflow.Context, in OverrideRequest) error {
if in.NewLimitSatang <= 0 {
return fmt.Errorf("limit must be positive")
}
if in.NewLimitSatang > 50_000_00 {
return fmt.Errorf("override exceeds maximum allowed")
}
if in.ApproverID == "" {
return fmt.Errorf("approver required")
}
return nil
},
},
)
if err != nil {
return TopUpResult{}, err
}
// ... ตรรกะที่ใช้ approvedLimit ...
_ = approvedLimit
return TopUpResult{}, nil
}
ฝั่ง client:
handle, err := c.UpdateWorkflow(context.Background(), client.UpdateWorkflowOptions{
WorkflowID: "topup-" + txnID,
UpdateName: "approve-override",
Args: []interface{}{OverrideRequest{NewLimitSatang: 2_000_00, ApproverID: "ops-42"}},
WaitForStage: client.WorkflowUpdateStageCompleted,
})
if err != nil {
return err // validator ปฏิเสธ หรือส่งไม่ถึง
}
var result OverrideResult
if err := handle.Get(context.Background(), &result); err != nil {
return err
}
WaitForStage เลือกได้ 2 แบบ:
WorkflowUpdateStageAccepted— คืนทันทีที่ validator ผ่าน (ยังไม่รู้ผลลัพธ์)WorkflowUpdateStageCompleted— รอจนกว่า handler จะทำงานเสร็จและได้ผลลัพธ์
Validator ต้องอ่านอย่างเดียวจริงๆ
เหมือน query handler: ห้ามแก้ state, ห้ามเรียก activity, ห้าม sleep ข้อดีของ validator คือ update ที่ถูกปฏิเสธ จะไม่ถูกบันทึกลง history เลย จึงกันไม่ให้ history โตจากคำสั่งที่ผิดตั้งแต่แรก
จาก CLI:
temporal workflow update execute \
--workflow-id topup-TXN-001 \
--name approve-override \
--input '{"NewLimitSatang":200000,"ApproverID":"ops-42"}'
`temporal workflow update` เป็นกลุ่มคำสั่ง ไม่ใช่คำสั่งเดียว
ต้องใช้ subcommand เสมอ: execute, start, result, describe
และ --wait-for-stage ของ update start รับค่าได้เพียง accepted เท่านั้น
Update-with-Start
ปัญหาที่พบบ่อย: API ได้ request เข้ามา ต้อง "สร้าง wallet session ถ้ายังไม่มี แล้วสั่งเติมเงินทันที" ถ้าทำ 2 ขั้น (start แล้วค่อย update) จะมี race และเสีย round trip เพิ่ม
Update-with-Start ทำ 2 อย่างนี้ในคำสั่งเดียวแบบ atomic
startOp := c.NewWithStartWorkflowOperation(
client.StartWorkflowOptions{
ID: "wallet-" + walletID,
TaskQueue: "wallet-core",
// จำเป็นต้องระบุ — บอกว่าถ้ามีตัวที่รันอยู่แล้วให้ใช้ตัวเดิม
WorkflowIDConflictPolicy: enumspb.WORKFLOW_ID_CONFLICT_POLICY_USE_EXISTING,
},
workflows.WalletAccountWorkflow, walletID,
)
handle, err := c.UpdateWithStartWorkflow(context.Background(),
client.UpdateWithStartWorkflowOptions{
StartWorkflowOperation: startOp,
UpdateOptions: client.UpdateWorkflowOptions{
UpdateName: "top-up",
Args: []interface{}{TopUpCommand{TransactionID: txnID, AmountSatang: 100000}},
WaitForStage: client.WorkflowUpdateStageCompleted,
},
})
if err != nil {
return err
}
var res TopUpAccepted
err = handle.Get(context.Background(), &res)
ใช้คู่กับ Use Existing เสมอ
เอกสาร Temporal แนะนำให้เลือก WORKFLOW_ID_CONFLICT_POLICY_USE_EXISTING
และเขียน update handler ให้ idempotent เพื่อให้ client retry ได้อย่างปลอดภัย
ในบริบทการเงินแปลว่า handler ต้องเช็คว่า transaction ID นี้เคยรับไปแล้วหรือยัง
Signal-with-Start
รูปแบบเก่ากว่าแต่ยังใช้ได้ดีเมื่อไม่ต้องการผลลัพธ์กลับ:
_, err := c.SignalWithStartWorkflow(context.Background(),
"wallet-"+walletID,
"bank-callback", callbackPayload,
client.StartWorkflowOptions{ID: "wallet-" + walletID, TaskQueue: "wallet-core"},
workflows.WalletAccountWorkflow, walletID,
)
เหมาะกับการรับ webhook จากธนาคาร ที่อาจมาถึงก่อนที่ workflow ฝั่งเราจะถูกสร้างด้วยซ้ำ
กับดัก: handler ที่ยังทำงานค้าง
Signal และ Update handler สามารถเรียก activity ได้ ซึ่งแปลว่ามันอาจยังทำงานไม่เสร็จ ตอนที่ workflow หลักกำลังจะ return ผลคือ error TMPRL1102 (Unfinished handlers)
// ก่อนจบ workflow ต้องรอ handler ทั้งหมดให้เสร็จก่อน
if err := workflow.Await(ctx, func() bool {
return workflow.AllHandlersFinished(ctx)
}); err != nil {
return err
}
return result, nil
กฎเดียวกันนี้ใช้กับ Continue-As-New ด้วย — ต้องรอ handler เสร็จก่อนเสมอ
สรุปการเลือกใช้ในระบบ wallet
- Query — หน้าจอลูกค้าถามว่า "เติมเงินถึงไหนแล้ว", dashboard ของ ops
- Signal — ผู้ใช้ส่ง OTP, ธนาคารยิง webhook กลับมา, ops สั่งให้เดินหน้าต่อ
- Update — ขออนุมัติ override แล้วต้องรู้ผลทันที, สั่งงานที่ต้อง validate ก่อนรับ