package gateway import ( "context" "database/sql" "net/http" "net/http/httptest" "path/filepath" "testing" "time" "github.com/example/notify-gateway/internal/config" "github.com/example/notify-gateway/internal/model" "github.com/example/notify-gateway/internal/outbox" ) func TestQueueRetriesOnlyFailedDestinationAndFreezesDryRun(t *testing.T) { ctx := context.Background() dir := t.TempDir() store, err := config.Open(filepath.Join(dir, "config.json")) if err != nil { t.Fatal(err) } db, err := outbox.Open(filepath.Join(dir, "outbox.db")) if err != nil { t.Fatal(err) } defer db.Close() good, bad := 0, 0 fail := true srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if r.Header.Get("Idempotency-Key") == "" { t.Error("missing delivery id") } if r.URL.Path == "/good" { good++ w.WriteHeader(204) } else { bad++ if fail { w.WriteHeader(400) } else { w.WriteHeader(200) } } })) defer srv.Close() c := store.Get() c.Outbounds = []config.OutboundConfig{{ID: "good", Provider: "webhook", URL: srv.URL + "/good", Live: true}, {ID: "bad", Provider: "webhook", URL: srv.URL + "/bad", Live: true}, {ID: "dry", Provider: "webhook", URL: srv.URL + "/good"}} c.Mappings = nil for _, id := range []string{"good", "bad", "dry"} { c.Mappings = append(c.Mappings, config.Mapping{ID: id, Enabled: true, Target: "webhook", OutboundID: id, TextTemplate: "{{.Message}}"}) } if err := store.Replace(c); err != nil { t.Fatal(err) } q := &Queue{Store: db, Dispatcher: New(store, nil)} msg := model.InboundMessage{Source: "webhook", Channel: "ops", Message: "test", ReceivedAt: time.Now()} first, err := q.Accept(ctx, msg, "request-1") if err != nil { t.Fatal(err) } c.Outbounds[2].Live = true store.Replace(c) for i := 0; i < 3; i++ { if ok, err := q.ProcessOne(ctx); err != nil || !ok { t.Fatalf("process %v %v", ok, err) } } if good != 1 || bad != 1 { t.Fatalf("calls good=%d bad=%d", good, bad) } msg.ReceivedAt = time.Now().Add(time.Hour) duplicate, err := q.Accept(ctx, msg, "request-1") if err != nil || !duplicate.Duplicate || first.ID != duplicate.ID { t.Fatalf("receipt %+v %v", duplicate, err) } rows, _ := db.List(ctx, 50, 0) for _, j := range rows { if j.MappingID == "bad" { if j.State != "dead" { t.Fatal(j.State) } if err := db.Retry(ctx, j.ID); err != nil { t.Fatal(err) } } if j.MappingID == "dry" && j.State != "dry_run" { t.Fatal("dry-run turned live") } } fail = false if ok, err := q.ProcessOne(ctx); !ok || err != nil { t.Fatal(err) } if good != 1 || bad != 2 { t.Fatalf("successful destination repeated: good=%d bad=%d", good, bad) } } func TestQueueBackoffAndAtomicInvalidMapping(t *testing.T) { ctx := context.Background() dir := t.TempDir() store, _ := config.Open(filepath.Join(dir, "config.json")) db, err := outbox.Open(filepath.Join(dir, "outbox.db")) if err != nil { t.Fatal(err) } defer db.Close() calls := 0 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { calls++ w.Header().Set("Retry-After", "120") w.WriteHeader(429) })) defer srv.Close() c := store.Get() c.Outbounds = []config.OutboundConfig{{ID: "rate", Provider: "webhook", URL: srv.URL, Live: true}} c.Mappings = []config.Mapping{{ID: "rate", Enabled: true, Target: "webhook", OutboundID: "rate", TextTemplate: "{{.Message}}"}} if err := store.Replace(c); err != nil { t.Fatal(err) } q := &Queue{Store: db, Dispatcher: New(store, nil)} if _, err := q.Accept(ctx, model.InboundMessage{Message: "one"}, "one"); err != nil { t.Fatal(err) } if _, err := q.ProcessOne(ctx); err != nil { t.Fatal(err) } rows, _ := db.List(ctx, 50, 0) if rows[0].State != "pending" || rows[0].Next < time.Now().Add(119*time.Second).Unix() { t.Fatalf("backoff: %+v", rows) } if worked, err := q.ProcessOne(ctx); err != nil || worked || calls != 1 { t.Fatalf("early retry %v %v %d", worked, err, calls) } c.Mappings = append(c.Mappings, config.Mapping{ID: "invalid", Enabled: true, Target: "discord", OutboundID: "discord", TextTemplate: ""}) c.Outbounds = append(c.Outbounds, config.OutboundConfig{ID: "discord", Provider: "discord", URL: "https://discord.com/api/webhooks/123/token"}) store.Replace(c) if _, err := q.Accept(ctx, model.InboundMessage{Message: "new"}, "new"); err == nil { t.Fatal("expected invalid empty Discord payload") } rows, _ = db.List(ctx, 50, 0) if len(rows) != 1 { t.Fatal("partial message was queued") } // Advance the due time in the isolated database to exercise the retry cap // without waiting for the backoff clock in this integration test. clockDB, err := sql.Open("sqlite", filepath.Join(dir, "outbox.db")) if err != nil { t.Fatal(err) } defer clockDB.Close() for i := 1; i < 8; i++ { if _, err := clockDB.Exec("UPDATE deliveries SET next=0 WHERE state='pending'"); err != nil { t.Fatal(err) } if _, err := q.ProcessOne(ctx); err != nil { t.Fatal(err) } } rows, _ = db.List(ctx, 50, 0) if rows[0].State != "dead" || rows[0].Attempts != 8 || calls != 8 { t.Fatalf("retry limit: %+v calls=%d", rows, calls) } }