D
Engineering

Broadcast Scheduler

Broadcast fan-out dengan rate limiter 8 msg/s — SubmitBroadcastScheduler dan SubmitBroadcastMessage spawn goroutine RunningSubmitBroadcastMessage yang kirim ke WA gateway via EngineUri. Context timeout 60 menit.

Broadcast mengirim pesan ke banyak penerima sekaligus. Go API memicu fan-out, Node engine tidak terlibat di sini — pengiriman langsung ke WA gateway :5002 via submitMessage.

Dua endpoint

EndpointHandlerMode
POST /v1/broadcast/schedulerSubmitBroadcastScheduler (broadcast.go:29)Jadwal — baca dari broadcast_scedulers collection
POST /v1/send_broadcastSubmitBroadcastMessage (broadcast.go:237)One-shot — payload langsung

Keduanya spawn goroutine RunningSubmitBroadcastMessage (broadcast.go:266) dan return immediately.

Scheduler flow — SubmitBroadcastScheduler

text
1. Parse PayloadBroadcastScheduler
   └─ require BroadcastSchedulerId (broadcast.go:32-39)

2. Generate batchId = uuid.New()
   └─ create broadcast logger (broadcast.go:49-54)

3. Read broadcast_scedulers collection (broadcast.go:57)
   └─ FindOne {"id": broadcastId} → BroadcastSchedulerStruct
   └─ 404 jika tidak ditemukan

4. Aggregation $lookup (broadcast.go:78-101)
   ├─ join devices (device ↔ id)
   └─ join broadcast_messages (message_id ↔ id)
   └─ $unwind preserveNullAndEmptyArrays: true

5. Validate (broadcast.go:133-148)
   ├─ DeviceInfo.Token non-empty
   └─ Messages non-empty

6. Normalize contact/product/status/tag/segment lists (broadcast.go:150-195)
   └─ ke []string / []TagList

7. Build PayloadBroadcast (broadcast.go:197-220)

8. go RunningSubmitBroadcastMessage(formattedPayload, &scheduler.DeviceInfo) (broadcast.go:229)
   └─ fire and forget

9. Return {"status":"success","message":"Broadcast Sedang Dalam Pengiriman"}

One-shot flow — SubmitBroadcastMessage

text
1. Parse PayloadBroadcast (broadcast.go:238-239)
   └─ no validation

2. GetDevicebyFilter({"token": payload.Token}) (broadcast.go:247-256)

3. go RunningSubmitBroadcastMessage(*payload, orderPaket) (broadcast.go:258)

4. Return immediately

Fan-out worker — RunningSubmitBroadcastMessage

broadcast.go:266 — goroutine utama.

text
1. defer recovery (broadcast.go:281-292)
   └─ jika final status "failed" dan mode "broadcast_sceduler":
      UpdateScheduler({"id": idScheduler}, {"onprogress": false})

2. ctx, cancel := context.WithTimeout(context.Background(), 60*time.Minute) (broadcast.go:294)
   └─ overall timeout 60 menit

3. Rate limiter (broadcast.go:324-325)
   └─ messagesPerSecond := 8
   └─ rateLimiter := utils.NewRateLimiter(messagesPerSecond)

4. Fan-out loop (per penerima, line 600+):
   ├─ rateLimiter.Wait(ctx)              (broadcast.go:679)  — blocking
   ├─ uri = EngineUri + "send-broadcast" (broadcast.go:689)
   ├─ formData = url.Values{number, instance, message}
   └─ submitMessage(ctx, uri, formData, ...) (broadcast.go:699)

5. Media variants → EngineUri + "send-media"
   └─ broadcast.go:794, 808, 822, 944, 958, 972

6. Per-message log → broadcast_logs collection (broadcast.go:705+)

Rate limiter — utils/rate_limiter.go

go
type RateLimiter struct {
    limiter *rate.Limiter
}

func NewRateLimiter(messagesPerSecond int) *RateLimiter {
    return &RateLimiter{
        limiter: rate.NewLimiter(rate.Limit(messagesPerSecond), messagesPerSecond),
    }
}

func (rl *RateLimiter) Wait(ctx context.Context) error {
    return rl.limiter.Wait(ctx)
}

rate.NewLimiter(rate.Limit(8), 8) = token bucket:

  • Rate: 8 tokens/detik
  • Burst capacity: 8 (limiter segar mengizinkan 8 token sekaligus, lalu refill 8/detik)

submitMessage — HTTP sender

broadcast.go:1110:

go
func submitMessage(ctx context.Context, uri string, formData url.Values,
    number string, messageIdx int, totalMessages int,
    batchId string, schedularId string) (*ResponseSubmitMessage, error)
{
    // nested timeout 90 detik di dalam parent 60 menit
    ctxWithTimeout, cancel := context.WithTimeout(ctx, 90*time.Second)  // broadcast.go:1136
    // POST form-encoded
    // decode ResponseSubmitMessage
}

Tenant CRM

handlers/broadcast_crm.go:127:

go
uri := EngineUriCRM + "send-broadcast"   // :5003, bukan :5002
// ...
submitPoint(uri, formData, nil)          // non-context sender lama

Tenant CRM memakai submitPoint (non-context, defined di chat_live.go:2828), bukan submitMessage dengan context. Tidak ada rate limiter.

Diagram alur

text
POST /v1/broadcast/scheduler
    │
    ▼
SubmitBroadcastScheduler
    ├─ read broadcast_scedulers collection
    ├─ $lookup devices + broadcast_messages
    ├─ build PayloadBroadcast
    └─ go RunningSubmitBroadcastMessage(...)  ─┐
                                                 │
POST /v1/send_broadcast                          │
    │                                            │
    ▼                                            │
SubmitBroadcastMessage                           │
    ├─ GetDevicebyFilter                         │
    └─ go RunningSubmitBroadcastMessage(...)  ───┤
                                                 │
                                                 ▼
                                    RunningSubmitBroadcastMessage
                                        ├─ ctx 60 menit
                                        ├─ rateLimiter 8 msg/s
                                        └─ loop per penerima:
                                            ├─ Wait(ctx)
                                            ├─ POST :5002/send-broadcast
                                            └─ log → broadcast_logs

Langkah berikutnya