226 lines
5.6 KiB
Go
226 lines
5.6 KiB
Go
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
|
|
}
|