diff --git a/internal/subscriber/handle.go b/internal/subscriber/handle.go index df71133..ccbe65d 100644 --- a/internal/subscriber/handle.go +++ b/internal/subscriber/handle.go @@ -28,6 +28,8 @@ type ProcessFunc func(ctx context.Context, req notify.Request) (notify.Result, e type HandleInput struct { Body []byte Headers map[string]any + Name string + Queue string SourceName string MaxRetry int Deduper Deduper @@ -81,29 +83,37 @@ func HandleMessage(ctx context.Context, in HandleInput, conv *tradesignal.Conver } } + raw := string(in.Body) + if hash == "" { + hash = MessageHash(in.Body) + } + slog.Info("mq message", "name", in.Name, "queue", in.Queue, "hash", hash, "raw", raw) + event, data, err := conv.Convert(in.Body) if err != nil { - slog.Warn("invalid signal, ack", "error", err) + slog.Warn("invalid signal, ack", "hash", hash, "raw", raw, "error", err) return DispositionAck } src, err := lookup(ctx, in.SourceName) if err != nil || src == nil || src.Status != 1 { - slog.Warn("source unavailable, ack", "source", in.SourceName, "error", err) + slog.Warn("source unavailable, ack", "source", in.SourceName, "hash", hash, "raw", raw, "error", err) return DispositionAck } res, err := process(ctx, notify.Request{Source: src, Event: event, Data: data}) if err == nil { if !res.Matched { - slog.Info("no matching rule", "source", src.Name, "event", event) + slog.Info("no matching rule", "source", src.Name, "event", event, "hash", hash, "raw", raw) } else if res.Filtered { - slog.Info("rule filtered", "source", src.Name, "event", event, "reason", res.Reason) + slog.Info("rule filtered", "source", src.Name, "event", event, "reason", res.Reason, "hash", hash, "raw", raw) + } else { + slog.Info("mq message accepted", "source", src.Name, "event", event, "channels", res.Channels, "hash", hash, "raw", raw) } return DispositionAck } if errors.Is(err, notify.ErrUnprocessable) { - slog.Warn("unprocessable notify, ack", "source", src.Name, "event", event, "error", err) + slog.Warn("unprocessable notify, ack", "source", src.Name, "event", event, "hash", hash, "raw", raw, "error", err) return DispositionAck } diff --git a/internal/subscriber/subscriber.go b/internal/subscriber/subscriber.go index 11cce3b..34bb797 100644 --- a/internal/subscriber/subscriber.go +++ b/internal/subscriber/subscriber.go @@ -121,6 +121,8 @@ func (s *Subscriber) handleDelivery(ch *amqp.Channel, d amqp.Delivery) { disp := HandleMessage(context.Background(), HandleInput{ Body: d.Body, Headers: map[string]any(d.Headers), + Name: s.cfg.Name, + Queue: s.cfg.Queue, SourceName: s.cfg.Source, MaxRetry: s.cfg.MaxRetry, Deduper: s.deduper,