Skip to main content

RabbitMQ Publisher และ Topology

ภาพรวม

webhook-go เป็น producer-only gateway กล่าวคือทุก endpoint ทำงานเบา ๆ แล้วโยนงานลง RabbitMQ ในโปรเจกต์นี้ไม่มี consumer ไม่มี cron และไม่มี background job ใด ๆ เลย consumer ทั้งหมด อยู่ที่ worker-go

รูปแบบการส่งข้อความที่ใช้คือ bare JSON กล่าวคือ body ของ message เป็น payload object ตรง ๆ ไม่มี envelope wrapper แบบ {type, data, meta} ที่ template ของบริษัทใช้เป็นมาตรฐาน เหตุผลคือ consumer ฝั่ง worker-go ซึ่ง port มาจาก NestJS parse object ดิบอยู่แล้ว จึงต้องรักษา compatibility ไว้ ในโค้ดมี comment เตือนเรื่องนี้ไว้อย่างชัดเจน

Routing ใช้รูปแบบที่ง่ายที่สุด คือ exchange เป็นชนิด direct และ routing key เท่ากับชื่อคิวเสมอ

Business Flow

Topology ที่ประกาศตอน boot

cmd/api/main.go เรียก AMQP.DeclareTopology(...) หลังจาก publisher ออนไลน์แล้ว หากล้มเหลวจะบันทึก log ระดับ warn แล้วทำงานต่อ เนื่องจากคิวอาจถูกสร้างไว้ล่วงหน้าโดยฝั่ง infra

  1. ประกาศ exchange หลัก line_exchange ชนิด direct แบบ durable
  2. ประกาศ DLQ exchange line_exchange_dlq ชนิด direct แบบ durable
  3. สำหรับทุกคิวในรายการ
    • assertQueue แบบ durable พร้อม argument x-max-priority: 10, x-dead-letter-exchange: line_exchange_dlq และ x-dead-letter-routing-key เป็นชื่อคิว
    • bindQueue คิวเข้ากับ line_exchange โดยใช้ routing key เป็นชื่อคิว
    • assertQueue คิว .dlq ที่คู่กันแบบ durable
    • bindQueue คิว .dlq เข้ากับ line_exchange_dlq โดยใช้ routing key เป็นชื่อคิวหลัก

คิวทั้งหมดที่ประกาศ และผู้ใช้งานจริง

คิวenvpublisher ในโปรเจกต์นี้
line_webhookRABBITMQ_QUEUE_LINE_WEBHOOKLINE Webhook Gateway เส้นทาง default และ cache miss
line_forward_webhookRABBITMQ_QUEUE_LINE_FORWARD_WEBHOOKไม่มี — ประกาศไว้เฉย ๆ เพราะการ forward ทำผ่าน HTTP โดยตรง ดู Webhook Forward
tracking_logRABBITMQ_QUEUE_TRACKING_LOGบันทึก Tracking Log
mookept_audienceRABBITMQ_QUEUE_MOOKEPT_AUDIENCE_QUEUEจัดการสมาชิก Audience เฉพาะ PUT
mbox_callbackRABBITMQ_QUEUE_MBOX_CALLBACKMBOX Chatwoot Callback
mbox_handoffRABBITMQ_QUEUE_MBOX_HANDOFFMBOX Agent Handoff
booking_notificationRABBITMQ_QUEUE_BOOKING_NOTIFICATIONBooking Cancel Postback
note

ไฟล์ .env.example ยังมีค่า RABBITMQ_QUEUE_LINE_CHANGE_RICHMENU=line_change_richmenu อยู่ แต่ internal/config/api.go ไม่ได้อ่านค่านี้ ถือเป็นของค้างจาก service เดิม

การ publish ผ่าน Manager.PublishRaw

  1. แปลง payload เป็น bytes ด้วย json.Marshal
  2. เปิด channel ใหม่ต่อการส่งหนึ่งครั้ง เพื่อไม่ให้ channel ที่เสียหายทำให้ publisher ทั้งตัวค้าง
  3. เปิด publisher confirm และลงทะเบียน listener ก่อนที่จะ publish เพื่อปิดช่องว่างของ race condition
  4. ส่งข้อความด้วย DeliveryMode: Persistent, ContentType: application/json และค่า Priority ซึ่งปกติเป็น 0
  5. รอ ack หากไม่ได้รับ ack จะคืน error และหาก context หมดเวลาจะคืน ctx.Err()
  6. ผู้เรียกแต่ละรายจัดการ error ต่างกัน
    • เส้นทางแบบ fire-and-forget เช่น LINE gateway และ tracking จะบันทึก log อย่างเดียว ข้อความที่ส่งไม่สำเร็จจะสูญหาย
    • เส้นทางแบบ synchronous เช่น mbox callback และ PUT /api/audience จะตอบ error กลับไปยัง client

การเชื่อมต่อ

  • AMQP_URLS รองรับการระบุหลาย URL สำหรับการหมุนเวียนแบบ HA จัดการโดย internal/amqp/manager.go
  • การ publish ถูกป้องกันด้วย mutex เนื่องจาก AMQP channel ไม่ concurrency-safe
  • publisher ถูกห่อเป็น deps.Dependency ผ่าน internal/deps/amqppub จึงเข้าสู่ระบบ readiness, supervisor และ graceful shutdown โดยอัตโนมัติ ดู Health Check, Metrics และการเฝ้าระวัง Dependency
  • หาก AMQP_URLS ว่าง deps.AMQP จะเป็น nil และทุกการ publish จะคืน ErrNotConfigured แต่ service ยังคง boot และตอบ health check ได้ตามปกติ

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

ไฟล์ของสำคัญ
internal/amqp/publisher.goManager.PublishRaw(ctx, exchange, queue, payload, priority) — ช่องทางเดียวที่ทุกโมดูลใช้
internal/amqp/topology.gotype Topology และ Manager.DeclareTopology(t)
internal/amqp/manager.goconnection manager แบบหลาย URL, Channel(), Ready(), การ reconnect
internal/amqp/envelope.goenvelope {type,data,meta} ของ template ซึ่งไม่ได้ถูกใช้ในโปรเจกต์นี้
internal/deps/amqppub/amqppub.goห่อ publisher เป็น deps.Dependency (Name/Connect/Ping/Close)
internal/config/api.gostruct RabbitMQ และเมธอด Queues() ที่กำหนดลำดับการประกาศ
internal/config/api_load.goอ่านค่าจาก env พร้อมค่า default ของทุกคิว
cmd/api/main.goเรียก DeclareTopology ตอน boot โดย guard ด้วยเงื่อนไข AMQP != nil && Ready()

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

  • worker-go — เป็น consumer ของทุกคิวข้างต้นและเป็นคู่สัญญาหลัก การเปลี่ยนรูปร่างของ payload ที่นี่ถือเป็น breaking change ต่อฝั่ง worker ทันที
  • RabbitMQ cluster — คิวและ DLX อาจถูกสร้างล่วงหน้าโดยทีม infra หรือผ่าน CRD ซึ่งไม่มีปัญหา เพราะ DeclareTopology เป็น idempotent จึงประกาศซ้ำได้
  • DLQ — ทุกคิวมี mirror .dlq คู่กัน แต่ฝั่ง publisher นี้ไม่เคยใช้งาน DLQ จะทำงานเมื่อ consumer ฝั่ง worker reject ข้อความหรือ retry จนครบจำนวน
  • ค่า environment ทั้งหมดขึ้นต้นด้วย RABBITMQ_ ดูรายการเต็มได้ที่ .env.example