Action แบบหน่วงเวลา / ตั้งเวลา
ภาพรวม
Workflow จำนวนมากต้องการ "รอ" ก่อนทำอะไรต่อ เช่น ส่งข้อความติดตามผลหลังลูกค้าซื้อครบ 3 วัน หรือรอดูว่าลูกค้าจะกดลิงก์ภายใน 24 ชั่วโมงหรือไม่ ถ้าไม่กดจึงค่อยส่งซ้ำ
ระบบเลือกที่จะไม่ใช้ delayed message ของ broker แต่บันทึกลงตาราง scheduled_action แล้วมี cron
มากวาดแถวที่ถึงเวลาทุกนาที เพราะวิธีนี้ตรวจสอบได้ แก้ไขได้ และรอดจากการ restart ของ worker
Business Flow
ขั้นตอนการสร้าง (schedule)
TriggerService.executeActionเมื่อพบว่าactionConfig.delayมีค่า จะเรียกScheduleActionแทนการยิง action ทันที- คำนวณ
execute_atตามโหมด- โหมดปกติ ใช้เวลาปัจจุบันบวก
delay.seconds - โหมด
waitForTrackingจะตั้งexecute_atเป็นเวลา หมดอายุ (ใช้expiresSecondsหรือdelaySeconds) พร้อมบันทึกexpires_atซึ่งมีความหมายว่า "รอจนหมดเวลา เว้นแต่จะมีคนคลิกก่อน"
- โหมดปกติ ใช้เวลาปัจจุบันบวก
- หาก
delay.reEvaluateConditions = trueจะเก็บsourceConfigไว้ด้วย เพื่อนำมา ประเมินเงื่อนไขซ้ำ ตอนถึงเวลาจริง เพราะสภาพของลูกค้าอาจเปลี่ยนไปแล้ว - ก่อนบันทึกจะตรวจ Redis ว่ามี tracking click ที่มาถึงก่อนหน้าแล้วหรือไม่ เพื่อกัน race condition
ขั้นตอนการทำงาน (process due actions)
รันทุกนาที ทั้งบน profile trigger-worker ผ่าน ScheduledActionService.ProcessDueActions
และบน profile cron-scheduler ผ่าน ScheduledActionScannerService.Run ซึ่งจะ publish เข้า queue แทน
recoverStaleProcessingกู้แถวที่ค้างอยู่ในสถานะprocessingเกิน 5 นาที ซึ่งเกิดจาก worker ตายกลางคันclaimBatchจองงานแบบ atomic ครั้งละ 500 แถว ด้วยFOR UPDATE SKIP LOCKEDทำให้หลาย replica ทำงานพร้อมกันได้โดยไม่ชนกัน- ประมวลผลทีละ action ผ่าน
processAction- โหลด
line_userปัจจุบัน - หากมี
sourceConfigเก็บไว้ จะเรียกEvaluateConditionsซ้ำ หากไม่ผ่านแล้วก็ยกเลิก action นั้น - หากผ่าน จะสร้าง synthetic
TriggerRuleขึ้นมาแล้วเรียกExecuteActionของ trigger engine - retry ได้สูงสุด 3 ครั้ง (
maxRetries) ก่อน mark เป็น failed
- โหลด
handleExpiredจัดการ action ที่เลยexpires_atแล้วแต่ยังไม่ถูกปลดล็อก ตามที่ config ไว้ เช่น เดินเส้นทาง branch "ไม่คลิก"- วน
claimBatchต่อไปจนไม่มีงานเหลือ
การปลดล็อกก่อนถึงเวลา
หากผู้ใช้กดลิงก์ระหว่างที่ระบบกำลังรออยู่ tracking consumer จะเรียก checkAndCompleteDelayWait
เพื่อให้ action ของ workflow นั้นทำงานทันที
ไฟล์และฟังก์ชันหลัก
internal/scheduledaction/scheduledaction.goScheduledActionService.ScheduleAction(ctx, params)— บันทึกลงตารางscheduled_actionProcessDueActions(ctx)— cron รายนาทีของ profiletrigger-workerclaimBatch()(ใช้FOR UPDATE SKIP LOCKED),processAction(),recoverStaleProcessing(),handleExpired(),updateStatus(),getLineUser()- ค่าคงที่
batchSize = 500,maxRetries = 3,staleProcessingMinute = 5 - interface
TriggerExecutorที่มีEvaluateConditionsและExecuteActionใช้ตัดวงจร circular dependency
internal/actionexecute/scheduler.go—ActionSchedulerService.ScheduleAction()ซึ่งเป็น เวอร์ชันเบาที่ใช้ที่ profilemainโดย insert อย่างเดียวและไม่มี croninternal/cronscheduler/scheduled_action_scanner.go—ScheduledActionScannerService.Run()ของ profilecron-schedulerที่ claim แล้ว publish เข้าscheduled_action_fireinternal/trigger/service_events.go—ProcessScheduledActionFire()และhandleScheduledActionExpired()cmd/worker/main.go—runTriggerWorker()ที่ cross-wire ทั้งสอง service เข้าหากัน- Queue:
scheduled_action_fireโดย consume ที่trigger-workerและ publish จากcron-scheduler
จุดเชื่อมต่อกับ Service อื่น
- ตาราง
scheduled_action— คอลัมน์สำคัญคือexecute_at,expires_at,status(pending, processing, done, failed, cancelled),retry_count,action_type,action_config,source_config,workflow_idและworkflow_node_path - ตารางอื่น:
trigger_rule(rule ต้นทาง) และline_user(สำหรับประเมินเงื่อนไขซ้ำ) - Redis — ใช้ตรวจ tracking click ที่มาถึงก่อนเวลา
- RabbitMQ — queue
scheduled_action_fireโดยมีปลายทางสุดท้ายคือaction_execute - เชื่อมโดยตรงกับ เครื่องยนต์ Trigger / Workflow Automation และ การลงมือทำ Action ของ Workflow