Skip to main content

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 แทน

  1. recoverStaleProcessing กู้แถวที่ค้างอยู่ในสถานะ processing เกิน 5 นาที ซึ่งเกิดจาก worker ตายกลางคัน
  2. claimBatch จองงานแบบ atomic ครั้งละ 500 แถว ด้วย FOR UPDATE SKIP LOCKED ทำให้หลาย replica ทำงานพร้อมกันได้โดยไม่ชนกัน
  3. ประมวลผลทีละ action ผ่าน processAction
    • โหลด line_user ปัจจุบัน
    • หากมี sourceConfig เก็บไว้ จะเรียก EvaluateConditions ซ้ำ หากไม่ผ่านแล้วก็ยกเลิก action นั้น
    • หากผ่าน จะสร้าง synthetic TriggerRule ขึ้นมาแล้วเรียก ExecuteAction ของ trigger engine
    • retry ได้สูงสุด 3 ครั้ง (maxRetries) ก่อน mark เป็น failed
  4. handleExpired จัดการ action ที่เลย expires_at แล้วแต่ยังไม่ถูกปลดล็อก ตามที่ config ไว้ เช่น เดินเส้นทาง branch "ไม่คลิก"
  5. วน claimBatch ต่อไปจนไม่มีงานเหลือ

การปลดล็อกก่อนถึงเวลา

หากผู้ใช้กดลิงก์ระหว่างที่ระบบกำลังรออยู่ tracking consumer จะเรียก checkAndCompleteDelayWait เพื่อให้ action ของ workflow นั้นทำงานทันที

ไฟล์และฟังก์ชันหลัก

  • internal/scheduledaction/scheduledaction.go
    • ScheduledActionService.ScheduleAction(ctx, params) — บันทึกลงตาราง scheduled_action
    • ProcessDueActions(ctx) — cron รายนาทีของ profile trigger-worker
    • claimBatch() (ใช้ FOR UPDATE SKIP LOCKED), processAction(), recoverStaleProcessing(), handleExpired(), updateStatus(), getLineUser()
    • ค่าคงที่ batchSize = 500, maxRetries = 3, staleProcessingMinute = 5
    • interface TriggerExecutor ที่มี EvaluateConditions และ ExecuteAction ใช้ตัดวงจร circular dependency
  • internal/actionexecute/scheduler.goActionSchedulerService.ScheduleAction() ซึ่งเป็น เวอร์ชันเบาที่ใช้ที่ profile main โดย insert อย่างเดียวและไม่มี cron
  • internal/cronscheduler/scheduled_action_scanner.goScheduledActionScannerService.Run() ของ profile cron-scheduler ที่ claim แล้ว publish เข้า scheduled_action_fire
  • internal/trigger/service_events.goProcessScheduledActionFire() และ handleScheduledActionExpired()
  • cmd/worker/main.gorunTriggerWorker() ที่ 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