โครงสร้าง 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
- กำหนด profile —
entrypoint.shหรือstart.shกำหนดค่าSVC(ค่าเริ่มต้นคือmain) หากใส่ค่าที่ไม่รู้จัก process จะ exit ทันที - จัดสรร admin port — หากไม่ได้ตั้ง
ADMIN_ADDRระบบจะกำหนด admin port ให้อัตโนมัติแยกตาม profile (:9104ถึง:9110) เพื่อไม่ให้ทั้ง 7 profile แย่งพอร์ต:9100กันเมื่อรันบนเครื่องเดียวกัน - Bootstrap —
app.Newโหลด config แล้วตรวจสอบด้วยValidate()แบบ fail-fast จากนั้นเชื่อมต่อ dependency ตามตารางdepsFor(profile)ภายใน timeout 30 วินาที และสร้างRedisServiceแบบมี namespace - Dispatch เข้า runner — runner ของแต่ละ profile จะสร้าง
mq.Client, assert topology, ประกอบ service ทุกโดเมน แล้ว register คู่ queue กับ handler ลงในmq.Registrarตัวเดียวกัน - รอสัญญาณ — main goroutine block รอ
ctx.Done()จาก SIGTERM หรือ SIGINT แล้วเรียกa.Shutdown()โดยลำดับคือ fail readiness probe ก่อน → หยุดรับงานใหม่ → drain งานที่ค้างอยู่ → ปิด dependency แบบย้อนลำดับ - เฝ้าระวัง outage —
supervisorเฝ้าแต่ละ dependency แยกกัน หาก dependency ล่มจะ retry ทุก 5 วินาทีในนาทีแรก จากนั้นทุก 10 วินาที ยิง alert เข้า chat webhook เมื่อครบ 1 นาที และเมื่อครบ 3 นาทีจะ alert ซ้ำแล้ว exit เพื่อให้ Kubernetes restart pod ให้
ตารางสรุป profile
| SVC | หน้าที่ | Dependency | Admin port (ค่าเริ่มต้น) |
|---|---|---|---|
main | 23 consumer หลักของงาน delivery พร้อม DLQ monitor | PG + Redis + Rabbit | :9104 |
trigger-worker | 8 consumer ของ trigger engine และ cron สแกน due-actions ทุกนาที | PG + Redis + Rabbit | :9105 |
web-request-worker | consumer ของ queue web_request_execute | PG + Redis + Rabbit | :9106 |
bigquery-sync | dynamic cron ที่อ่าน config จาก DB เพื่อ sync BigQuery เข้า PG | PG อย่างเดียว | :9107 |
knowledge-worker | consumer ของ queue knowledge_index | PG + Redis + Rabbit | :9108 |
pg-listener | bridge PG LISTEN/NOTIFY เข้า RabbitMQ พร้อม outbox sweep | PG + Rabbit | :9109 |
cron-scheduler | 10 cron scanner และ 4 consumer (mbox, booking) | PG + Redis + Rabbit | :9110 |
ไฟล์และฟังก์ชันหลัก
cmd/worker/main.go—main(),validProfilesและ runner ทั้ง 7 ตัว ได้แก่runMain,runPgListener,runTriggerWorker,runCronScheduler,runBigquerySync,runWebRequestWorker,runKnowledgeWorkercmd/worker/integration.go— adapter ข้ามโดเมนทั้งหมด แต่ละ package จะประกาศ interface แคบ ๆ ของ sibling service ที่ต้องเรียกใช้ แล้ว binary นี้เป็นผู้สร้าง implementation จริงและ inject ให้internal/app/app.go—New(),depsFor(profile),Shutdown()internal/config/config.go—Config,Queues,Load(),Validate(),HealthPort()internal/supervisor/supervisor.go— state machine เฝ้า outage รายตัว โดยevaluateเขียนเป็น pure function ที่มี unit test ครอบคลุมinternal/health/health.go— endpoint/livez,/readyz,/healthzinternal/obs/admin.goและinternal/obs/metrics.go— admin port ที่ให้บริการ/metrics,/versionและ pprof (เปิด/ปิดด้วยPPROF_ENABLED)internal/logx/logx.go— structured JSON log ที่ปรับ level ได้ด้วยสัญญาณ SIGHUPinternal/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 และจาก profilepg-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