Skip to main content

โครงสร้าง Worker และ Profile การรัน (SVC)

ภาพรวม

line-management-worker-go เป็น background worker service ที่เขียนด้วยภาษา Go โดย port มาจาก NestJS worker เดิมแบบ 1:1

ทั้ง repository ถูก build ออกมาเป็น binary เพียงตัวเดียว (cmd/worker) แต่สามารถทำงานได้ 7 บทบาท (profile) โดยเลือกผ่าน environment variable SVC ตอนรัน แต่ละ profile จะ register consumer และ cron job คนละชุด รวมถึงเชื่อมต่อ dependency คนละชุดกัน ทำให้สามารถ deploy แยก pod ได้ตามภาระงานจริง เช่น scale trigger-worker แยกจาก main ได้อิสระ

ทุก profile ใช้ bootstrap ร่วมกันตัวเดียวคือ internal/app.New ซึ่งทำหน้าที่โหลด config จาก environment variable → เชื่อมต่อ dependency (Postgres / Redis / RabbitMQ ตามที่ profile นั้นต้องการ) → เปิด admin server สำหรับ metrics และ health check → เริ่ม outage supervisor → รอสัญญาณเพื่อ graceful shutdown

Business Flow

  1. กำหนด profileentrypoint.sh หรือ start.sh กำหนดค่า SVC (ค่าเริ่มต้นคือ main) หากใส่ค่าที่ไม่รู้จัก process จะ exit ทันที
  2. จัดสรร admin port — หากไม่ได้ตั้ง ADMIN_ADDR ระบบจะกำหนด admin port ให้อัตโนมัติแยกตาม profile (:9104 ถึง :9110) เพื่อไม่ให้ทั้ง 7 profile แย่งพอร์ต :9100 กันเมื่อรันบนเครื่องเดียวกัน
  3. Bootstrapapp.New โหลด config แล้วตรวจสอบด้วย Validate() แบบ fail-fast จากนั้นเชื่อมต่อ dependency ตามตาราง depsFor(profile) ภายใน timeout 30 วินาที และสร้าง RedisService แบบมี namespace
  4. Dispatch เข้า runner — runner ของแต่ละ profile จะสร้าง mq.Client, assert topology, ประกอบ service ทุกโดเมน แล้ว register คู่ queue กับ handler ลงใน mq.Registrar ตัวเดียวกัน
  5. รอสัญญาณ — main goroutine block รอ ctx.Done() จาก SIGTERM หรือ SIGINT แล้วเรียก a.Shutdown() โดยลำดับคือ fail readiness probe ก่อน → หยุดรับงานใหม่ → drain งานที่ค้างอยู่ → ปิด dependency แบบย้อนลำดับ
  6. เฝ้าระวัง outagesupervisor เฝ้าแต่ละ dependency แยกกัน หาก dependency ล่มจะ retry ทุก 5 วินาทีในนาทีแรก จากนั้นทุก 10 วินาที ยิง alert เข้า chat webhook เมื่อครบ 1 นาที และเมื่อครบ 3 นาทีจะ alert ซ้ำแล้ว exit เพื่อให้ Kubernetes restart pod ให้

ตารางสรุป profile

SVCหน้าที่DependencyAdmin port (ค่าเริ่มต้น)
main23 consumer หลักของงาน delivery พร้อม DLQ monitorPG + Redis + Rabbit:9104
trigger-worker8 consumer ของ trigger engine และ cron สแกน due-actions ทุกนาทีPG + Redis + Rabbit:9105
web-request-workerconsumer ของ queue web_request_executePG + Redis + Rabbit:9106
bigquery-syncdynamic cron ที่อ่าน config จาก DB เพื่อ sync BigQuery เข้า PGPG อย่างเดียว:9107
knowledge-workerconsumer ของ queue knowledge_indexPG + Redis + Rabbit:9108
pg-listenerbridge PG LISTEN/NOTIFY เข้า RabbitMQ พร้อม outbox sweepPG + Rabbit:9109
cron-scheduler10 cron scanner และ 4 consumer (mbox, booking)PG + Redis + Rabbit:9110

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

  • cmd/worker/main.gomain(), validProfiles และ runner ทั้ง 7 ตัว ได้แก่ runMain, runPgListener, runTriggerWorker, runCronScheduler, runBigquerySync, runWebRequestWorker, runKnowledgeWorker
  • cmd/worker/integration.go — adapter ข้ามโดเมนทั้งหมด แต่ละ package จะประกาศ interface แคบ ๆ ของ sibling service ที่ต้องเรียกใช้ แล้ว binary นี้เป็นผู้สร้าง implementation จริงและ inject ให้
  • internal/app/app.goNew(), depsFor(profile), Shutdown()
  • internal/config/config.goConfig, Queues, Load(), Validate(), HealthPort()
  • internal/supervisor/supervisor.go — state machine เฝ้า outage รายตัว โดย evaluate เขียนเป็น pure function ที่มี unit test ครอบคลุม
  • internal/health/health.go — endpoint /livez, /readyz, /healthz
  • internal/obs/admin.go และ internal/obs/metrics.go — admin port ที่ให้บริการ /metrics, /version และ pprof (เปิด/ปิดด้วย PPROF_ENABLED)
  • internal/logx/logx.go — structured JSON log ที่ปรับ level ได้ด้วยสัญญาณ SIGHUP
  • internal/signalx/signalx.go — multiplex สัญญาณ SIGTERM/SIGINT สำหรับ drain และ SIGHUP สำหรับ reload

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

  • PostgreSQL (TYPEORM_*) — shared pgx pool ผ่าน internal/platform/pg ทุกโดเมนใช้ raw SQL ที่คัดลอกมาจาก TypeORM query เดิมแบบ verbatim
  • Redis / Dragonfly (REDIS_URL, REDIS_NAMESPACE) — ใช้เป็น cache, distributed lock, dedup และเก็บ session state
  • RabbitMQ (RABBITMQ_URL) — broker หลักที่รับ job จาก webhook-go, cms-api-go, client-api-go และจาก profile pg-listener ผ่าน pg_notify
  • S3-compatible storage (STORAGE_*) — เก็บไฟล์ CSV ของ audience, import archive และ error log
  • LINE Messaging API (LINE_ENDPOINT) — สร้าง client แยกต่อ OA โดยใช้ line_oa.channel_access_token
  • CMS API (CMS_API_BASE_URL, INTERNAL_API_KEY) — เรียกกลับสำหรับงาน audience refresh
  • Meilisearch และ Embedding API (MEILISEARCH_*, EMBEDDING_API_URL) — ใช้ในงาน knowledge index และ AI response
  • SMTP (MAIL_SMTP_*) — ส่งอีเมลสำหรับ reset password
  • Alert webhook (ALERT_WEBHOOK_URL, DLQ_ALERT_WEBHOOK_URL) — แจ้งเตือนเข้า Google Chat หรือ Slack