เครื่องยนต์ Trigger / Workflow Automation
ภาพรวม
Trigger engine คือหัวใจของระบบ automation ทั้งหมด แนวคิดคือ "ถ้ามีเหตุการณ์ X เกิดขึ้นกับลูกค้าคนหนึ่ง
และเขาเข้าเงื่อนไข Y ให้ทำ Z" โดยกฎเหล่านี้เก็บอยู่ในตาราง trigger_rule และผูกเข้ากับ workflow ได้
Profile trigger-worker รัน consumer ทั้งหมด 8 queue ตามชนิดของเหตุการณ์ต้นทาง แต่ทุก queue
จะมารวมกันที่ pipeline เดียวกัน คือ กรอง rule → ตรวจ dedup และ cooldown → ประเมินเงื่อนไข →
สั่ง action หรือกำหนดเวลาไว้ทำภายหลัง
Business Flow
Source type ทั้ง 8 แบบ
| Queue | Source type | ต้นทางของเหตุการณ์ |
|---|---|---|
attribute_change | attribute_change แบบ batch | pg-listener ที่รับจาก pg_notify('attribute_changed') และเส้นทาง import |
attribute_change_realtime | attribute_change รายคน | เส้นทาง realtime |
audience_membership_trigger | audience_membership | audience consumer เมื่อมีการเพิ่มหรือถอนสมาชิก |
friend_track_event_trigger | friend_track_event | pg-listener ที่รับจาก friend_track_event_inserted |
form_submitted_trigger | form_submitted | pg-listener เมื่อ client-api บันทึกฟอร์ม |
campaign_click_trigger | campaign_click | pg-listener ที่รับจาก trigger บนตาราง tracking_log |
scheduled_trigger_fire | scheduled | cron scanner ที่ profile cron-scheduler |
scheduled_action_fire | — | cron scanner สำหรับ action ที่ถึงเวลาทำ |
Pipeline ต่อหนึ่ง message
- Source-type gate —
HasActiveSourceType(lineOaId, organizationId, sourceType)ตรวจจาก cache (TTL 60 วินาที) ว่า OA นี้มี rule ชนิดนี้เปิดใช้อยู่หรือไม่ หากไม่มีจะ ack ทันทีโดยไม่แตะ DB เพื่อกันไม่ให้ event ท่วมระบบโดยเปล่าประโยชน์ - โหลด rule —
getTriggerRules(cache 60 วินาที) แล้วกรองเฉพาะ rule ที่ตรง source type และตรง target เช่นsourceConfig.audienceIdต้องตรงกับ audience ที่เปลี่ยน isKeyRelevant— สำหรับกรณี attribute change หาก key ที่เปลี่ยนไม่เกี่ยวข้องกับเงื่อนไข ของ rule ก็ข้ามไปเลย- Dedup และ cooldown —
shouldSkipByDedupใช้ Redis กันไม่ให้ rule เดียวกันยิงซ้ำให้ user เดิม ตามค่าfrequency(เช่น ครั้งเดียวตลอดกาล) และcooldownSeconds - ประเมินเงื่อนไข —
evaluateConditionsและevaluateCondition- รองรับ operator เปรียบเทียบครบชุด รวมถึงการเทียบวันที่แบบนับวัน โดย
roundDiffDaysเลียนแบบพฤติกรรมMath.roundของ JS เพื่อให้ผลลัพธ์ตรงกับระบบเดิม - รองรับ split test ผ่าน
evaluateSplitTestซึ่งแบ่งกลุ่มผู้ใช้แบบ deterministic จาก userId
- รองรับ operator เปรียบเทียบครบชุด รวมถึงการเทียบวันที่แบบนับวัน โดย
- ลงมือทำ action —
executeAction(rule, user, triggerDepth)- หาก
actionConfig.delayมีค่า จะ ไม่ทำทันที แต่เรียกScheduleActionเพื่อบันทึกลงตารางscheduled_action(ดู Action แบบหน่วงเวลา / ตั้งเวลา) - หากไม่มี delay จะ publish
ActionExecutePayloadเข้า queueaction_executeหรือเข้าweb_request_executeเมื่อactionType = web_request - ค่าคงที่
maxTriggerDepth = 5ทำหน้าที่ป้องกัน action ที่ไปกระตุ้น trigger ต่อจนกลายเป็นลูป ไม่รู้จบ
- หาก
- บันทึก log — เขียนลง
trigger_logทั้ง rule id, user, triggerOn, status และ action type/config - handler ส่วนใหญ่ return nil เสมอ เพราะ service จัดการ error และ log เอง มีเพียง payload ที่เสียหาย เท่านั้นที่จะเข้า DLQ
scheduled_trigger_fire
cron จะค้นหา rule ที่มี source_type='scheduled' และถึงเวลาทำงาน โดยคำนวณตาม timezone
Asia/Bangkok เป็นรายกฎ แล้ว publish เข้า queue นี้ จากนั้น consumer จะขยายออกเป็นรายผู้ใช้
แบบ batch ละ 500 คน (scheduledFireBatchSize) แล้วเดินเข้า pipeline เดียวกัน
ไฟล์และฟังก์ชันหลัก
internal/trigger/consumer.go—Consumer.Register()ผูกทั้ง 8 queue พร้อม handlerHandleAttributeChange,HandleRealtimeAttributeChange,HandleAudienceMembershipTrigger,HandleFriendTrackEventTrigger,HandleFormSubmittedTrigger,HandleCampaignClickTrigger,HandleScheduledTriggerFire,HandleScheduledActionFireinternal/trigger/service.goTriggerService.EvaluateTriggers(),EvaluateTriggersForBatch(),processAttributeRule()HasActiveSourceType(),getTriggerRules(),shouldSkipByDedup(),isKeyRelevant()evaluateCondition(),evaluateConditions(),evaluateSplitTest(),getUserValue()executeAction(),getLineUser(),mapLineUser()- ค่าคงที่
cacheTTLSecs = 60,maxTriggerDepth = 5,scheduledFireBatchSize = 500
internal/trigger/service_events.goEvaluateAudienceTriggers(),EvaluateFriendTrackTriggers(),EvaluateFormSubmittedTriggers(),EvaluateCampaignClickTriggers()ProcessScheduledTriggerFire(),ProcessScheduledActionFire(),handleScheduledActionExpired()
internal/trigger/repository.go—TriggerRuleRepositoryและTriggerLogRepository.LogTrigger()internal/payloads/payloads.go—ActionExecutePayloadซึ่งเป็นสัญญาข้ามโดเมนcmd/worker/main.go—runTriggerWorker()ที่ทำ cross-wire ระหว่างTriggerServiceกับScheduledActionService- Queue: consume ทั้ง 8 queue ข้างต้น (profile
trigger-worker) และ publish ไปยังaction_executeกับweb_request_execute
จุดเชื่อมต่อกับ Service อื่น
- รับ job จาก: pg-listener (5 channel), audience processing, line user CSV import
และ cron scanner ที่ profile
cron-scheduler - ตารางที่เกี่ยวข้อง:
trigger_rule(ตัวกฎ),trigger_log(บันทึกการยิง),line_user(ข้อมูลผู้ใช้ที่ใช้ประเมินเงื่อนไข),scheduled_action(action ที่หน่วงเวลา) และworkflow(จัดกลุ่ม rule) - Redis: cache ของ rule และ source-type gate (TTL 60 วินาที) รวมถึง key สำหรับ dedup และ cooldown ต่อคู่ (rule, user)
- RabbitMQ: publish ต่อไปยัง การลงมือทำ Action ของ Workflow และ การเรียก API ภายนอกจาก Workflow
- ENV toggle:
ENABLE_TRIGGER_EVALUATIONและENABLE_TRIGGER_SOURCE_TYPE_CACHE - ฝั่ง CMS: cms-api-go โดเมน workflow และ trigger-rule เป็นผู้สร้างและแก้ไขกฎเหล่านี้