Skip to main content

เครื่องยนต์ 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 แบบ

QueueSource typeต้นทางของเหตุการณ์
attribute_changeattribute_change แบบ batchpg-listener ที่รับจาก pg_notify('attribute_changed') และเส้นทาง import
attribute_change_realtimeattribute_change รายคนเส้นทาง realtime
audience_membership_triggeraudience_membershipaudience consumer เมื่อมีการเพิ่มหรือถอนสมาชิก
friend_track_event_triggerfriend_track_eventpg-listener ที่รับจาก friend_track_event_inserted
form_submitted_triggerform_submittedpg-listener เมื่อ client-api บันทึกฟอร์ม
campaign_click_triggercampaign_clickpg-listener ที่รับจาก trigger บนตาราง tracking_log
scheduled_trigger_firescheduledcron scanner ที่ profile cron-scheduler
scheduled_action_firecron scanner สำหรับ action ที่ถึงเวลาทำ

Pipeline ต่อหนึ่ง message

  1. Source-type gateHasActiveSourceType(lineOaId, organizationId, sourceType) ตรวจจาก cache (TTL 60 วินาที) ว่า OA นี้มี rule ชนิดนี้เปิดใช้อยู่หรือไม่ หากไม่มีจะ ack ทันทีโดยไม่แตะ DB เพื่อกันไม่ให้ event ท่วมระบบโดยเปล่าประโยชน์
  2. โหลด rulegetTriggerRules (cache 60 วินาที) แล้วกรองเฉพาะ rule ที่ตรง source type และตรง target เช่น sourceConfig.audienceId ต้องตรงกับ audience ที่เปลี่ยน
  3. isKeyRelevant — สำหรับกรณี attribute change หาก key ที่เปลี่ยนไม่เกี่ยวข้องกับเงื่อนไข ของ rule ก็ข้ามไปเลย
  4. Dedup และ cooldownshouldSkipByDedup ใช้ Redis กันไม่ให้ rule เดียวกันยิงซ้ำให้ user เดิม ตามค่า frequency (เช่น ครั้งเดียวตลอดกาล) และ cooldownSeconds
  5. ประเมินเงื่อนไขevaluateConditions และ evaluateCondition
    • รองรับ operator เปรียบเทียบครบชุด รวมถึงการเทียบวันที่แบบนับวัน โดย roundDiffDays เลียนแบบพฤติกรรม Math.round ของ JS เพื่อให้ผลลัพธ์ตรงกับระบบเดิม
    • รองรับ split test ผ่าน evaluateSplitTest ซึ่งแบ่งกลุ่มผู้ใช้แบบ deterministic จาก userId
  6. ลงมือทำ actionexecuteAction(rule, user, triggerDepth)
    • หาก actionConfig.delay มีค่า จะ ไม่ทำทันที แต่เรียก ScheduleAction เพื่อบันทึกลงตาราง scheduled_action (ดู Action แบบหน่วงเวลา / ตั้งเวลา)
    • หากไม่มี delay จะ publish ActionExecutePayload เข้า queue action_execute หรือเข้า web_request_execute เมื่อ actionType = web_request
    • ค่าคงที่ maxTriggerDepth = 5 ทำหน้าที่ป้องกัน action ที่ไปกระตุ้น trigger ต่อจนกลายเป็นลูป ไม่รู้จบ
  7. บันทึก log — เขียนลง trigger_log ทั้ง rule id, user, triggerOn, status และ action type/config
  8. 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.goConsumer.Register() ผูกทั้ง 8 queue พร้อม handler HandleAttributeChange, HandleRealtimeAttributeChange, HandleAudienceMembershipTrigger, HandleFriendTrackEventTrigger, HandleFormSubmittedTrigger, HandleCampaignClickTrigger, HandleScheduledTriggerFire, HandleScheduledActionFire
  • internal/trigger/service.go
    • TriggerService.EvaluateTriggers(), EvaluateTriggersForBatch(), processAttributeRule()
    • HasActiveSourceType(), getTriggerRules(), shouldSkipByDedup(), isKeyRelevant()
    • evaluateCondition(), evaluateConditions(), evaluateSplitTest(), getUserValue()
    • executeAction(), getLineUser(), mapLineUser()
    • ค่าคงที่ cacheTTLSecs = 60, maxTriggerDepth = 5, scheduledFireBatchSize = 500
  • internal/trigger/service_events.go
    • EvaluateAudienceTriggers(), EvaluateFriendTrackTriggers(), EvaluateFormSubmittedTriggers(), EvaluateCampaignClickTriggers()
    • ProcessScheduledTriggerFire(), ProcessScheduledActionFire(), handleScheduledActionExpired()
  • internal/trigger/repository.goTriggerRuleRepository และ TriggerLogRepository.LogTrigger()
  • internal/payloads/payloads.goActionExecutePayload ซึ่งเป็นสัญญาข้ามโดเมน
  • cmd/worker/main.gorunTriggerWorker() ที่ทำ 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 เป็นผู้สร้างและแก้ไขกฎเหล่านี้