Skip to main content

การจ่ายงาน Campaign ตามกำหนดเวลา (process_campaign)

ภาพรวม

Campaign ที่ผู้ใช้ตั้งเวลาส่งไว้ใน CMS จะถูกบันทึกเป็นแถวในตาราง campaign ด้วยสถานะ scheduled พร้อมกับค่า start_date

งานนี้คือ scanner ที่ทำงานทุกนาที คอยค้นหา campaign ที่ถึงเวลาส่งแล้ว จองสิทธิ์แบบ atomic จากนั้นโยนงานจริงเข้า queue line_broadcast_rich_message หรือ line_multicast_rich_message ตามค่า cast_type

หัวใจของ feature นี้คือการป้องกันการส่งซ้ำ ไม่ว่าจะมี worker หลาย replica scan พร้อมกัน หรือ message ถูก redeliver ผู้ใช้ปลายทางจะต้องไม่ได้รับข้อความซ้ำเด็ดขาด

Business Flow

  1. Tick — cron */1 * * * * บน profile cron-scheduler (CampaignScheduleScannerService) publish message เปล่า 1 ใบเข้า queue process_campaign
    • ใช้ Redis SET NX บน key CAMPAIGN_SCHEDULE_TICK:<yyyy-MM-ddTHH:mm> (TTL 55 วินาที) เพื่อให้มีเพียง replica เดียวเป็นผู้ publish
    • หาก Redis ล่ม ระบบจะ publish ต่อไปตามปกติ เพราะการจองสิทธิ์ที่ระดับฐานข้อมูลรับประกันความถูกต้องอยู่แล้ว
    • เดิม tick นี้มาจาก infrastructure ภายนอก ปัจจุบันย้ายเข้ามาอยู่ใน repository แล้ว
  2. Reaper — profile main รับ message นี้แล้วทำงานแรกคือกู้คืน campaign ที่ค้างสถานะ sending โดยมี claimed_at เก่ากว่า 15 นาที กลับไปเป็น scheduled เพราะถือว่า worker ที่จองไว้ตายไปแล้ว
  3. Atomic claim — สั่ง UPDATE campaign SET status='sending', claimed_at=NOW() โดยมีเงื่อนไข status='scheduled' AND start_date <= NOW() AND deleted_date IS NULL พร้อม RETURNING ทำให้ได้กลับมาเฉพาะแถวที่ tick รอบนี้จองได้จริง scanner ตัวอื่นจึงหยิบซ้ำไม่ได้
  4. Enqueue — วนทุก campaign ที่จองได้ แล้ว publish payload ที่ประกอบด้วย campaignId, audienceId, organizationId และ lineOaId เข้า queue ตาม cast_type โดยค่า broadcast ไปที่ line_broadcast_rich_message ส่วนค่าอื่นไปที่ line_multicast_rich_message
  5. Rollback เมื่อ publish ล้มเหลว — ระบบจะ revert status กลับเป็น scheduled ทันทีแบบ best-effort เพื่อให้ tick รอบถัดไปหยิบไปทำใหม่ ไม่ต้องรอ reaper นานถึง 15 นาที
  6. Heartbeat ระหว่างส่ง — consumer ที่ส่งจริงจะอัปเดต claimed_at ทุก 5 นาทีระหว่างการส่ง เพื่อไม่ให้ reaper มายึด campaign คืนกลางคัน (ดูการส่ง Campaign แบบ Multicast)
  7. ตรวจซ้ำที่ปลายทาง — delivery consumer ตรวจสถานะอีกชั้นหนึ่ง หาก campaign ไม่ได้อยู่ในสถานะ sending แล้ว เช่นเปลี่ยนเป็น sent ไปแล้ว จะทิ้ง message นั้นเงียบ ๆ ด้วย permanent error โดยไม่เปลี่ยนสถานะเป็น failed

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

  • internal/campaign/campaign.go
    • Service.ProcessCampaign(ctx) — รวมงาน reaper, atomic claim และ enqueue
    • ค่าคงที่ campaignStatusScheduled, campaignStatusSending, campaignStatusSent, campaignStatusCancel, campaignStatusDraft, campaignStatusFailed และ castTypeBroadcast, castTypeMulticast
  • internal/campaign/consumer.goConsumer.HandleProcessCampaign สำหรับ queue process_campaign และ Consumer.Register()
  • internal/cronscheduler/campaign_schedule_scanner.goCampaignScheduleScannerService.Run(ctx) และ campaignTickKeyPrefix
  • cmd/worker/integration.gocampaignForLineMessageApi.GetCampaignAndValidate() ซึ่งเป็นจุดที่ตรวจสถานะ sending และป้องกัน duplicate delivery
  • Queue: consume process_campaign บน profile main และ publish เข้า line_broadcast_rich_message หรือ line_multicast_rich_message

จุดเชื่อมต่อกับ Service อื่น

  • ตาราง campaign — อ่านและเขียนคอลัมน์ status, claimed_at, start_date, cast_type, audience_id, organization_id และ line_oa_id
  • Redis — key CAMPAIGN_SCHEDULE_TICK:* สำหรับ dedup tick ข้าม replica
  • RabbitMQ — consume process_campaign และ publish เข้า queue ที่ทำหน้าที่ส่งจริง
  • cms-api-go — เป็นผู้สร้างและแก้ไข campaign รวมถึงตั้งค่า start_date ในโดเมน campaign-management กรณีที่ผู้ใช้สั่งส่งทันที cms-api-go จะเปลี่ยนสถานะเป็น sending แล้ว enqueue เองโดยไม่รอ scanner
  • Feature นี้ ไม่มี การเรียก LINE API โดยตรง การเรียกจริงเกิดขึ้นในขั้นตอน delivery