บทที่ 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 จบ

เลือกให้ถูกตัว

QuerySignalUpdate
แก้ 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 ก่อนรับ