package worker import ( "context" "encoding/json" "fmt" "time" "doormile/db" "doormile/utils" nats "github.com/nats-io/nats.go" ) func StartBookingWorker() { defer func() { if r := recover(); r != nil { utils.Error("BookingWorker: panic recovered, restarting", "error", r) time.Sleep(5 * time.Second) go StartBookingWorker() } }() if db.Js == nil { utils.Warn("BookingWorker: NATS JetStream not available, worker not started") return } type consumerDef struct { subject string consumer string } consumers := []consumerDef{ {"api.v1.bookings.create", "bookings_create"}, {"api.v1.bookings.update", "bookings_update"}, {"api.v1.bookings.cancel", "bookings_cancel"}, } type subscription struct { subject string sub *nats.Subscription } var subs []subscription for _, c := range consumers { sub, err := db.Js.PullSubscribe(c.subject, c.consumer, nats.BindStream("BOOKINGS")) if err != nil { utils.Error("BookingWorker: subscribe failed", "subject", c.subject, "consumer", c.consumer, "error", err) continue } subs = append(subs, subscription{c.subject, sub}) utils.Info("BookingWorker: subscribed", "subject", c.subject, "consumer", c.consumer) } if len(subs) == 0 { utils.Error("BookingWorker: no subscriptions established, worker exiting") return } utils.Info("BookingWorker: running, waiting for messages...") for { for _, s := range subs { msgs, err := s.sub.Fetch(10, nats.MaxWait(1*time.Second)) if err != nil { if err == nats.ErrTimeout { continue } utils.Warn("BookingWorker: fetch error", "subject", s.subject, "error", err) time.Sleep(1 * time.Second) continue } for _, msg := range msgs { processMessage(msg) } } } } func processMessage(msg *nats.Msg) { defer func() { if r := recover(); r != nil { utils.Error("BookingWorker: panic in message handler", "error", r) msg.Nak() } }() var payload map[string]interface{} if err := json.Unmarshal(msg.Data, &payload); err != nil { utils.Error("BookingWorker: invalid JSON", "subject", msg.Subject, "error", err) msg.Term() return } utils.Info("BookingWorker: processing message", "subject", msg.Subject, "booking_id", payload["booking_id"]) var ok bool switch msg.Subject { case "api.v1.bookings.create": ok = handleCreate(payload) case "api.v1.bookings.update": ok = handleUpdate(payload) case "api.v1.bookings.cancel": ok = handleCancel(payload) default: utils.Warn("BookingWorker: no handler for subject", "subject", msg.Subject) msg.Ack() return } if ok { msg.Ack() } else { msg.Nak() } } func handleCreate(payload map[string]interface{}) bool { bid, ok := extractBookingID(payload) if !ok { utils.Error("BookingWorker: missing booking_id in create payload") return false } ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() mapping := pickFields(payload, "booking_id", "booking_no", "customer_id", "pickup_address", "pickup_pincode", "delivery_address", "delivery_pincode", "status", "created_at") if err := db.Rdb.HSet(ctx, fmt.Sprintf("bookings:%s", bid), mapping).Err(); err != nil { utils.Error("BookingWorker: HSET failed on create", "booking_id", bid, "error", err) return false } ctx2, cancel2 := context.WithTimeout(context.Background(), 3*time.Second) defer cancel2() db.Rdb.SAdd(ctx2, "bookings:all", bid) if cid := strOf(payload["customer_id"]); cid != "" { ctx3, cancel3 := context.WithTimeout(context.Background(), 3*time.Second) defer cancel3() db.Rdb.SAdd(ctx3, fmt.Sprintf("bookings:customer:%s", cid), bid) } utils.Info("BookingWorker: created booking in Redis", "booking_id", bid) return true } func handleUpdate(payload map[string]interface{}) bool { bid, ok := extractBookingID(payload) if !ok { return false } ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() exists, err := db.Rdb.Exists(ctx, fmt.Sprintf("bookings:%s", bid)).Result() if err != nil || exists == 0 { utils.Warn("BookingWorker: booking not in cache for update", "booking_id", bid) return false } ctx2, cancel2 := context.WithTimeout(context.Background(), 3*time.Second) defer cancel2() mapping := pickFields(payload, "status", "miler_id", "updated_at") if err := db.Rdb.HSet(ctx2, fmt.Sprintf("bookings:%s", bid), mapping).Err(); err != nil { utils.Error("BookingWorker: HSET failed on update", "booking_id", bid, "error", err) return false } utils.Info("BookingWorker: updated booking in Redis", "booking_id", bid) return true } func handleCancel(payload map[string]interface{}) bool { bid, ok := extractBookingID(payload) if !ok { return false } ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() mapping := pickFields(payload, "cancelled_at") mapping["status"] = "Cancelled" if err := db.Rdb.HSet(ctx, fmt.Sprintf("bookings:%s", bid), mapping).Err(); err != nil { utils.Error("BookingWorker: HSET failed on cancel", "booking_id", bid, "error", err) return false } utils.Info("BookingWorker: cancelled booking in Redis", "booking_id", bid) return true } func extractBookingID(payload map[string]interface{}) (string, bool) { v, exists := payload["booking_id"] if !exists || v == nil { return "", false } s := strOf(v) return s, s != "" } func strOf(v interface{}) string { if v == nil { return "" } return fmt.Sprintf("%v", v) } func pickFields(payload map[string]interface{}, keys ...string) map[string]interface{} { m := make(map[string]interface{}, len(keys)) for _, k := range keys { if v, ok := payload[k]; ok && v != nil { m[k] = strOf(v) } } return m }