การจ่ายงาน 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
- Tick — cron
*/1 * * * *บน profilecron-scheduler(CampaignScheduleScannerService) publish message เปล่า 1 ใบเข้า queueprocess_campaign- ใช้ Redis
SET NXบน keyCAMPAIGN_SCHEDULE_TICK:<yyyy-MM-ddTHH:mm>(TTL 55 วินาที) เพื่อให้มีเพียง replica เดียวเป็นผู้ publish - หาก Redis ล่ม ระบบจะ publish ต่อไปตามปกติ เพราะการจองสิทธิ์ที่ระดับฐานข้อมูลรับประกันความถูกต้องอยู่แล้ว
- เดิม tick นี้มาจาก infrastructure ภายนอก ปัจจุบันย้ายเข้ามาอยู่ใน repository แล้ว
- ใช้ Redis
- Reaper — profile
mainรับ message นี้แล้วทำงานแรกคือกู้คืน campaign ที่ค้างสถานะsendingโดยมีclaimed_atเก่ากว่า 15 นาที กลับไปเป็นscheduledเพราะถือว่า worker ที่จองไว้ตายไปแล้ว - Atomic claim — สั่ง
UPDATE campaign SET status='sending', claimed_at=NOW()โดยมีเงื่อนไขstatus='scheduled' AND start_date <= NOW() AND deleted_date IS NULLพร้อมRETURNINGทำให้ได้กลับมาเฉพาะแถวที่ tick รอบนี้จองได้จริง scanner ตัวอื่นจึงหยิบซ้ำไม่ได้ - Enqueue — วนทุก campaign ที่จองได้ แล้ว publish payload ที่ประกอบด้วย
campaignId,audienceId,organizationIdและlineOaIdเข้า queue ตามcast_typeโดยค่าbroadcastไปที่line_broadcast_rich_messageส่วนค่าอื่นไปที่line_multicast_rich_message - Rollback เมื่อ publish ล้มเหลว — ระบบจะ revert
statusกลับเป็นscheduledทันทีแบบ best-effort เพื่อให้ tick รอบถัดไปหยิบไปทำใหม่ ไม่ต้องรอ reaper นานถึง 15 นาที - Heartbeat ระหว่างส่ง — consumer ที่ส่งจริงจะอัปเดต
claimed_atทุก 5 นาทีระหว่างการส่ง เพื่อไม่ให้ reaper มายึด campaign คืนกลางคัน (ดูการส่ง Campaign แบบ Multicast) - ตรวจซ้ำที่ปลายทาง — delivery consumer ตรวจสถานะอีกชั้นหนึ่ง หาก campaign ไม่ได้อยู่ในสถานะ
sendingแล้ว เช่นเปลี่ยนเป็นsentไปแล้ว จะทิ้ง message นั้นเงียบ ๆ ด้วย permanent error โดยไม่เปลี่ยนสถานะเป็น failed
ไฟล์และฟังก์ชันหลัก
internal/campaign/campaign.goService.ProcessCampaign(ctx)— รวมงาน reaper, atomic claim และ enqueue- ค่าคงที่
campaignStatusScheduled,campaignStatusSending,campaignStatusSent,campaignStatusCancel,campaignStatusDraft,campaignStatusFailedและcastTypeBroadcast,castTypeMulticast
internal/campaign/consumer.go—Consumer.HandleProcessCampaignสำหรับ queueprocess_campaignและConsumer.Register()internal/cronscheduler/campaign_schedule_scanner.go—CampaignScheduleScannerService.Run(ctx)และcampaignTickKeyPrefixcmd/worker/integration.go—campaignForLineMessageApi.GetCampaignAndValidate()ซึ่งเป็นจุดที่ตรวจสถานะsendingและป้องกัน duplicate delivery- Queue: consume
process_campaignบน profilemainและ 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