From 03e86bfdcf8074640d5b18414bf4e761dcfe8e5b Mon Sep 17 00:00:00 2001 From: Phloraxx Date: Tue, 1 Sep 2026 05:15:10 +0000 Subject: [PATCH 1/2] feat(v4): deliver signed webhooks durably --- internal/v4/webhooks/service.go | 336 +++++++++++++++++++++++++++ internal/v4/webhooks/service_test.go | 255 ++++++++++++++++++++ 2 files changed, 591 insertions(+) create mode 100644 internal/v4/webhooks/service.go create mode 100644 internal/v4/webhooks/service_test.go diff --git a/internal/v4/webhooks/service.go b/internal/v4/webhooks/service.go new file mode 100644 index 0000000..5d2e65a --- /dev/null +++ b/internal/v4/webhooks/service.go @@ -0,0 +1,336 @@ +package webhooks + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "strconv" + "strings" + "sync" + "time" + + "github.com/Phloraxx/payment-api/internal/v4/storage" +) + +const ( + defaultMaxAttempts = 8 + defaultBatchSize = 50 + defaultLease = 2 * time.Minute +) + +var ErrInvalidConfig = errors.New("invalid webhook configuration") + +type Config struct { + Endpoint string + Secret string + AllowInsecureHTTP bool +} +type Service struct { + DB *storage.DB + Config Config + HTTPClient *http.Client + Now func() time.Time + MaxAttempts int + BatchSize int + Lease time.Duration + + wake chan struct{} + mu sync.Mutex +} + +type Delivery struct { + ID string + Body string + Attempts int +} + +func NewService(db *storage.DB, cfg Config) *Service { + return &Service{ + DB: db, Config: cfg, HTTPClient: newHTTPClient(), Now: time.Now, + MaxAttempts: defaultMaxAttempts, BatchSize: defaultBatchSize, + Lease: defaultLease, wake: make(chan struct{}, 1), + } +} + +func (s *Service) Enabled() bool { + return s != nil && strings.TrimSpace(s.Config.Endpoint) != "" && strings.TrimSpace(s.Config.Secret) != "" +} +func (s *Service) Wake() { + if !s.Enabled() { + return + } + select { + case s.wake <- struct{}{}: + default: + } +} + +func (s *Service) Run(ctx context.Context) { + if !s.Enabled() { + return + } + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + s.Wake() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + case <-s.wake: + } + _, _ = s.SendPending(ctx) + } +} + +func (s *Service) SendPending(ctx context.Context) (int, error) { + if !s.Enabled() { + return 0, nil + } + if err := s.ready(); err != nil { + return 0, err + } + s.mu.Lock() + defer s.mu.Unlock() + now := s.now() + limit := s.BatchSize + if limit <= 0 || limit > 500 { + limit = defaultBatchSize + } + rows, err := s.DB.SQL.QueryContext(ctx, `SELECT id FROM webhook_deliveries + WHERE status IN ('pending','retry') AND COALESCE(next_attempt_at,0) <= ? + ORDER BY COALESCE(next_attempt_at,0), created_at, rowid LIMIT ?`, now.UnixMilli(), limit) + if err != nil { + return 0, fmt.Errorf("list due webhooks: %w", err) + } + var ids []string + for rows.Next() { + var id string + if err := rows.Scan(&id); err != nil { + rows.Close() + return 0, fmt.Errorf("scan due webhook: %w", err) + } + ids = append(ids, id) + } + if err := rows.Close(); err != nil { + return 0, err + } + if err := rows.Err(); err != nil { + return 0, fmt.Errorf("iterate due webhooks: %w", err) + } + + processed := 0 + for _, id := range ids { + if err := ctx.Err(); err != nil { + return processed, err + } + delivery, err := s.claim(ctx, id, now) + if err != nil { + return processed, err + } + if delivery == nil { + continue + } + s.deliver(ctx, *delivery) + processed++ + } + return processed, nil +} +func (s *Service) claim(ctx context.Context, id string, now time.Time) (*Delivery, error) { + lease := s.Lease + if lease <= 0 { + lease = defaultLease + } + var out *Delivery + err := s.DB.WithImmediateTx(ctx, func(tx *storage.ImmediateTx) error { + result, err := tx.ExecContext(ctx, `UPDATE webhook_deliveries + SET attempts=attempts+1,next_attempt_at=? + WHERE id=? AND status IN ('pending','retry') AND COALESCE(next_attempt_at,0) <= ?`, + now.Add(lease).UnixMilli(), id, now.UnixMilli()) + if err != nil { + return fmt.Errorf("claim webhook %s: %w", id, err) + } + rows, err := result.RowsAffected() + if err != nil { + return err + } + if rows != 1 { + return nil + } + var delivery Delivery + if err := tx.QueryRowContext(ctx, `SELECT id,payload_json,attempts FROM webhook_deliveries WHERE id=?`, id). + Scan(&delivery.ID, &delivery.Body, &delivery.Attempts); err != nil { + return fmt.Errorf("load claimed webhook %s: %w", id, err) + } + out = &delivery + return nil + }) + return out, err +} + +func (s *Service) deliver(ctx context.Context, delivery Delivery) { + now := s.now() + timestamp := strconv.FormatInt(now.Unix(), 10) + signature := Sign(s.Config.Secret, timestamp, []byte(delivery.Body)) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, s.Config.Endpoint, strings.NewReader(delivery.Body)) + if err == nil { + req.Header.Set("Content-Type", "application/json") + req.Header.Set("User-Agent", "PayGate/4") + req.Header.Set("PayGate-Event-Id", delivery.ID) + req.Header.Set("PayGate-Timestamp", timestamp) + req.Header.Set("PayGate-Signature", "v1="+signature) + } + statusCode := 0 + retryable := true + if err == nil { + var response *http.Response + response, err = s.client().Do(req) + if response != nil { + statusCode = response.StatusCode + _, _ = io.Copy(io.Discard, io.LimitReader(response.Body, 64<<10)) + _ = response.Body.Close() + if statusCode < 200 || statusCode >= 300 { + err = fmt.Errorf("webhook returned HTTP %d", statusCode) + retryable = retryableHTTPStatus(statusCode) + } + } + } + _ = s.finish(context.Background(), delivery, statusCode, retryable, err) +} + +func (s *Service) finish(ctx context.Context, delivery Delivery, statusCode int, retryable bool, deliveryErr error) error { + now := s.now() + maxAttempts := s.MaxAttempts + if maxAttempts <= 0 { + maxAttempts = defaultMaxAttempts + } + return s.DB.WithImmediateTx(ctx, func(tx *storage.ImmediateTx) error { + if deliveryErr == nil { + _, err := tx.ExecContext(ctx, `UPDATE webhook_deliveries SET status='delivered',next_attempt_at=NULL,last_http_status=?,last_error=NULL,delivered_at=? WHERE id=?`, + nullableStatus(statusCode), now.UnixMilli(), delivery.ID) + return err + } + + message := truncate(deliveryErr.Error(), 4000) + if !retryable || delivery.Attempts >= maxAttempts { + _, err := tx.ExecContext(ctx, `UPDATE webhook_deliveries SET status='exhausted',next_attempt_at=NULL,last_http_status=?,last_error=?,delivered_at=NULL WHERE id=?`, + nullableStatus(statusCode), message, delivery.ID) + return err + } + _, err := tx.ExecContext(ctx, `UPDATE webhook_deliveries SET status='retry',next_attempt_at=?,last_http_status=?,last_error=?,delivered_at=NULL WHERE id=?`, + now.Add(retryDelay(delivery.Attempts)).UnixMilli(), nullableStatus(statusCode), message, delivery.ID) + return err + }) +} +func (s *Service) RetryOne(ctx context.Context, id string) error { + if err := s.ready(); err != nil { + return err + } + id = strings.TrimSpace(id) + if id == "" { + return errors.New("webhook id is required") + } + result, err := s.DB.SQL.ExecContext(ctx, `UPDATE webhook_deliveries + SET status='pending',attempts=0,next_attempt_at=?,last_http_status=NULL,last_error=NULL,delivered_at=NULL + WHERE id=? AND status IN ('retry','exhausted')`, s.now().UnixMilli(), id) + if err != nil { + return fmt.Errorf("retry webhook: %w", err) + } + if rows, _ := result.RowsAffected(); rows != 1 { + return errors.New("webhook is not retryable") + } + s.Wake() + return nil +} + +func (s *Service) ready() error { + if s == nil || s.DB == nil || s.DB.SQL == nil { + return errors.New("webhook storage is required") + } + endpoint := strings.TrimSpace(s.Config.Endpoint) + secret := strings.TrimSpace(s.Config.Secret) + if endpoint == "" || len(secret) < 32 { + return fmt.Errorf("%w: endpoint and at least 32-byte secret are required", ErrInvalidConfig) + } + u, err := url.Parse(endpoint) + if err != nil || u.Host == "" || u.User != nil || u.RawQuery != "" || u.Fragment != "" { + return fmt.Errorf("%w: endpoint must be an absolute URL without credentials/query/fragment", ErrInvalidConfig) + } + if u.Scheme != "https" && !(s.Config.AllowInsecureHTTP && u.Scheme == "http") { + return fmt.Errorf("%w: HTTPS endpoint is required", ErrInvalidConfig) + } + return nil +} + +func Sign(secret, timestamp string, body []byte) string { + mac := hmac.New(sha256.New, []byte(secret)) + _, _ = mac.Write([]byte(timestamp)) + _, _ = mac.Write([]byte(".")) + _, _ = mac.Write(body) + return hex.EncodeToString(mac.Sum(nil)) +} +func retryableHTTPStatus(status int) bool { + return status == http.StatusRequestTimeout || status == http.StatusTooEarly || status == http.StatusTooManyRequests || status >= 500 +} + +func retryDelay(attempt int) time.Duration { + delays := []time.Duration{ + time.Minute, + 5 * time.Minute, + 30 * time.Minute, + 2 * time.Hour, + 6 * time.Hour, + 12 * time.Hour, + 24 * time.Hour, + } + index := attempt - 1 + if index < 0 { + index = 0 + } + if index >= len(delays) { + return delays[len(delays)-1] + } + return delays[index] +} + +func (s *Service) now() time.Time { + if s.Now == nil { + return time.Now().UTC() + } + return s.Now().UTC() +} + +func (s *Service) client() *http.Client { + if s.HTTPClient == nil { + return newHTTPClient() + } + return s.HTTPClient +} +func newHTTPClient() *http.Client { + return &http.Client{ + Timeout: 10 * time.Second, + CheckRedirect: func(_ *http.Request, _ []*http.Request) error { + return http.ErrUseLastResponse + }, + } +} + +func nullableStatus(status int) any { + if status == 0 { + return nil + } + return status +} + +func truncate(value string, max int) string { + if len(value) <= max { + return value + } + return value[:max] +} diff --git a/internal/v4/webhooks/service_test.go b/internal/v4/webhooks/service_test.go new file mode 100644 index 0000000..67c748b --- /dev/null +++ b/internal/v4/webhooks/service_test.go @@ -0,0 +1,255 @@ +package webhooks + +import ( + "context" + "crypto/hmac" + "encoding/json" + "errors" + "io" + "net/http" + "net/http/httptest" + "path/filepath" + "sync/atomic" + "testing" + "time" + + "github.com/Phloraxx/payment-api/internal/v4/payments" + "github.com/Phloraxx/payment-api/internal/v4/profiles" + "github.com/Phloraxx/payment-api/internal/v4/storage" +) + +const testSecret = "0123456789abcdef0123456789abcdef" + +type webhookFixture struct { + db *storage.DB + payment payments.Payment + now *time.Time +} + +func newWebhookFixture(t *testing.T) webhookFixture { + t.Helper() + db, err := storage.Open(context.Background(), filepath.Join(t.TempDir(), "paygate.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = db.Close() }) + now := time.Date(2026, 9, 1, 9, 0, 0, 0, time.UTC) + profileService := profiles.NewService(db) + ctx := context.Background() + if _, err := profileService.Upsert(ctx, profiles.UpsertInput{ + ID: "paytm", Label: "Paytm", UPIID: "paygate@paytm", PayeeName: "PayGate", + Parser: "paytm_notification", Enabled: true, + }); err != nil { + t.Fatal(err) + } + if _, err := profileService.Activate(ctx, "paytm"); err != nil { + t.Fatal(err) + } + paymentService := payments.NewService(db) + paymentService.Now = func() time.Time { return now } + created, err := paymentService.Create(ctx, payments.CreateInput{ + RequestedAmountPaise: 10000, Name: "Sourav P Bijoy", ExternalID: "evt_test", + IdempotencyScope: "test", IdempotencyKey: "webhook-fixture", + }) + if err != nil { + t.Fatal(err) + } + return webhookFixture{db: db, payment: created.Payment, now: &now} +} + +func newTestService(f webhookFixture, endpoint string) *Service { + s := NewService(f.db, Config{Endpoint: endpoint, Secret: testSecret, AllowInsecureHTTP: true}) + s.Now = func() time.Time { return *f.now } + return s +} +func TestSendPendingSignsAndDelivers(t *testing.T) { + f := newWebhookFixture(t) + var seen int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + atomic.AddInt32(&seen, 1) + body, err := io.ReadAll(r.Body) + if err != nil { + t.Fatal(err) + } + eventID := r.Header.Get("PayGate-Event-Id") + timestamp := r.Header.Get("PayGate-Timestamp") + signature := r.Header.Get("PayGate-Signature") + if eventID == "" || timestamp == "" || signature == "" { + t.Fatalf("missing PayGate headers: %v", r.Header) + } + want := "v1=" + Sign(testSecret, timestamp, body) + if !hmac.Equal([]byte(signature), []byte(want)) { + t.Fatalf("signature=%q want=%q", signature, want) + } + var payload map[string]any + if err := json.Unmarshal(body, &payload); err != nil { + t.Fatal(err) + } + if payload["id"] != eventID || payload["type"] != "payment.created" { + t.Fatalf("payload=%v event=%s", payload, eventID) + } + w.WriteHeader(http.StatusNoContent) + })) + defer server.Close() + + s := newTestService(f, server.URL) + processed, err := s.SendPending(context.Background()) + if err != nil || processed != 1 || atomic.LoadInt32(&seen) != 1 { + t.Fatalf("processed=%d seen=%d err=%v", processed, seen, err) + } + var status string + var attempts int + var httpStatus int + var delivered int64 + if err := f.db.SQL.QueryRow(`SELECT status,attempts,last_http_status,delivered_at FROM webhook_deliveries WHERE payment_id=?`, f.payment.ID). + Scan(&status, &attempts, &httpStatus, &delivered); err != nil { + t.Fatal(err) + } + if status != "delivered" || attempts != 1 || httpStatus != http.StatusNoContent || delivered != f.now.UnixMilli() { + t.Fatalf("delivery state status=%s attempts=%d http=%d delivered=%d", status, attempts, httpStatus, delivered) + } +} + +func TestRetryable500BacksOffThenSucceeds(t *testing.T) { + f := newWebhookFixture(t) + var calls int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + call := atomic.AddInt32(&calls, 1) + if call == 1 { + http.Error(w, "temporary", http.StatusInternalServerError) + return + } + w.WriteHeader(http.StatusOK) + })) + defer server.Close() + + s := newTestService(f, server.URL) + if processed, err := s.SendPending(context.Background()); err != nil || processed != 1 { + t.Fatalf("first pass processed=%d err=%v", processed, err) + } + var status string + var attempts int + var next int64 + if err := f.db.SQL.QueryRow(`SELECT status,attempts,next_attempt_at FROM webhook_deliveries WHERE payment_id=?`, f.payment.ID). + Scan(&status, &attempts, &next); err != nil { + t.Fatal(err) + } + if status != "retry" || attempts != 1 || next != f.now.Add(time.Minute).UnixMilli() { + t.Fatalf("retry state status=%s attempts=%d next=%d", status, attempts, next) + } + *f.now = f.now.Add(time.Minute) + if processed, err := s.SendPending(context.Background()); err != nil || processed != 1 { + t.Fatalf("second pass processed=%d err=%v", processed, err) + } + if err := f.db.SQL.QueryRow(`SELECT status,attempts FROM webhook_deliveries WHERE payment_id=?`, f.payment.ID). + Scan(&status, &attempts); err != nil { + t.Fatal(err) + } + if status != "delivered" || attempts != 2 || atomic.LoadInt32(&calls) != 2 { + t.Fatalf("final retry state status=%s attempts=%d calls=%d", status, attempts, calls) + } +} + +func TestPermanent404ExhaustsImmediately(t *testing.T) { + f := newWebhookFixture(t) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + http.NotFound(w, nil) + })) + defer server.Close() + + s := newTestService(f, server.URL) + if processed, err := s.SendPending(context.Background()); err != nil || processed != 1 { + t.Fatalf("processed=%d err=%v", processed, err) + } + var status string + var attempts int + var next *int64 + var httpStatus int + if err := f.db.SQL.QueryRow(`SELECT status,attempts,next_attempt_at,last_http_status FROM webhook_deliveries WHERE payment_id=?`, f.payment.ID). + Scan(&status, &attempts, &next, &httpStatus); err != nil { + t.Fatal(err) + } + if status != "exhausted" || attempts != 1 || next != nil || httpStatus != http.StatusNotFound { + t.Fatalf("404 state status=%s attempts=%d next=%v http=%d", status, attempts, next, httpStatus) + } +} +func TestRetryOneResetsOnlyExplicitFailedEvent(t *testing.T) { + f := newWebhookFixture(t) + var statusCode int32 = http.StatusNotFound + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(int(atomic.LoadInt32(&statusCode))) + })) + defer server.Close() + s := newTestService(f, server.URL) + if _, err := s.SendPending(context.Background()); err != nil { + t.Fatal(err) + } + var eventID string + if err := f.db.SQL.QueryRow(`SELECT id FROM webhook_deliveries WHERE payment_id=?`, f.payment.ID).Scan(&eventID); err != nil { + t.Fatal(err) + } + atomic.StoreInt32(&statusCode, http.StatusOK) + if err := s.RetryOne(context.Background(), eventID); err != nil { + t.Fatal(err) + } + if _, err := s.SendPending(context.Background()); err != nil { + t.Fatal(err) + } + var status string + var attempts int + if err := f.db.SQL.QueryRow(`SELECT status,attempts FROM webhook_deliveries WHERE id=?`, eventID).Scan(&status, &attempts); err != nil { + t.Fatal(err) + } + if status != "delivered" || attempts != 1 { + t.Fatalf("manual retry status=%s attempts=%d", status, attempts) + } + if err := s.RetryOne(context.Background(), eventID); err == nil { + t.Fatal("delivered webhook should not be manually requeued") + } +} + +func TestClaimLeaseRecoversAfterInterruptedDelivery(t *testing.T) { + f := newWebhookFixture(t) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusOK) })) + defer server.Close() + s := newTestService(f, server.URL) + var eventID string + if err := f.db.SQL.QueryRow(`SELECT id FROM webhook_deliveries WHERE payment_id=?`, f.payment.ID).Scan(&eventID); err != nil { + t.Fatal(err) + } + claimed, err := s.claim(context.Background(), eventID, *f.now) + if err != nil || claimed == nil || claimed.Attempts != 1 { + t.Fatalf("claim=%+v err=%v", claimed, err) + } + if processed, err := s.SendPending(context.Background()); err != nil || processed != 0 { + t.Fatalf("leased event processed=%d err=%v", processed, err) + } + *f.now = f.now.Add(defaultLease + time.Second) + if processed, err := s.SendPending(context.Background()); err != nil || processed != 1 { + t.Fatalf("recovered event processed=%d err=%v", processed, err) + } + var status string + var attempts int + if err := f.db.SQL.QueryRow(`SELECT status,attempts FROM webhook_deliveries WHERE id=?`, eventID).Scan(&status, &attempts); err != nil { + t.Fatal(err) + } + if status != "delivered" || attempts != 2 { + t.Fatalf("recovered state status=%s attempts=%d", status, attempts) + } +} + +func TestConfigurationValidation(t *testing.T) { + f := newWebhookFixture(t) + cases := []Config{ + {Endpoint: "http://example.com/hook", Secret: testSecret}, + {Endpoint: "https://user:pass@example.com/hook", Secret: testSecret}, + {Endpoint: "https://example.com/hook?token=x", Secret: testSecret}, + {Endpoint: "https://example.com/hook", Secret: "short"}, + } + for _, cfg := range cases { + s := NewService(f.db, cfg) + if _, err := s.SendPending(context.Background()); !errors.Is(err, ErrInvalidConfig) { + t.Fatalf("config %+v err=%v", cfg, err) + } + } +} From db2a86af6dd30843656766a091005a672b6e4de7 Mon Sep 17 00:00:00 2001 From: Phloraxx Date: Tue, 1 Sep 2026 05:25:39 +0000 Subject: [PATCH 2/2] feat(v4): add operator overview activity and settings --- internal/v4/operator/service.go | 346 +++++++++++++++++++++++++++ internal/v4/operator/service_test.go | 167 +++++++++++++ internal/v4/operator/settings.go | 175 ++++++++++++++ internal/v4/webhooks/service.go | 66 +++-- 4 files changed, 737 insertions(+), 17 deletions(-) create mode 100644 internal/v4/operator/service.go create mode 100644 internal/v4/operator/service_test.go create mode 100644 internal/v4/operator/settings.go diff --git a/internal/v4/operator/service.go b/internal/v4/operator/service.go new file mode 100644 index 0000000..6f6c5ee --- /dev/null +++ b/internal/v4/operator/service.go @@ -0,0 +1,346 @@ +package operator + +import ( + "context" + "database/sql" + "errors" + "fmt" + "sort" + "time" + + "github.com/Phloraxx/payment-api/internal/v4/storage" +) + +type Service struct { + DB *storage.DB + Now func() time.Time + Location *time.Location +} + +type DailyVolume struct { + Date string `json:"date"` + AmountPaise int64 `json:"amount_paise"` + Payments int `json:"payments"` +} + +type Overview struct { + CollectedTodayPaise int64 + PaymentsToday int + PaidToday int + Pending int + ExpiredToday int + StatusCounts map[string]int + Volume []DailyVolume + ActiveProfile *ProfileSummary + Relay RelaySummary + Webhooks WebhookSummary +} +type ProfileSummary struct { + ID string + Label string +} + +type RelaySummary struct { + Connected bool + Name string + LastSeenAt *time.Time + AppVersion string +} + +type WebhookSummary struct { + Pending int + Exhausted int + LastDeliveredAt *time.Time +} + +type ActivityEntry struct { + At time.Time + Kind string + Status string + Source string + Title string + PaymentID string + AmountPaise *int64 + Detail string +} + +func NewService(db *storage.DB) *Service { + return &Service{DB: db, Now: time.Now, Location: time.FixedZone("IST", 5*60*60+30*60)} +} + +func (s *Service) ready() error { + if s == nil || s.DB == nil || s.DB.SQL == nil { + return errors.New("operator storage is required") + } + return nil +} +func (s *Service) Overview(ctx context.Context) (Overview, error) { + if err := s.ready(); err != nil { + return Overview{}, err + } + now := s.now() + start, end := localDayBounds(now, s.location()) + out := Overview{StatusCounts: map[string]int{}} + + if err := s.DB.SQL.QueryRowContext(ctx, `SELECT COALESCE(SUM(payable_amount_paise),0),COUNT(*) FROM payments + WHERE status='paid' AND paid_at>=? AND paid_at=? AND created_at=? AND grace_until=? ORDER BY paid_at`, start.UnixMilli()) + if err != nil { + return nil, fmt.Errorf("read volume trend: %w", err) + } + defer rows.Close() + byDate := make(map[string]*DailyVolume, days) + for i := 0; i < days; i++ { + date := startDay.AddDate(0, 0, i).Format("2006-01-02") + byDate[date] = &DailyVolume{Date: date} + } + for rows.Next() { + var paidAt, amount int64 + if err := rows.Scan(&paidAt, &amount); err != nil { + return nil, fmt.Errorf("scan volume trend: %w", err) + } + date := time.UnixMilli(paidAt).In(loc).Format("2006-01-02") + if bucket := byDate[date]; bucket != nil { + bucket.AmountPaise += amount + bucket.Payments++ + } + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate volume trend: %w", err) + } + out := make([]DailyVolume, 0, days) + for i := 0; i < days; i++ { + date := startDay.AddDate(0, 0, i).Format("2006-01-02") + out = append(out, *byDate[date]) + } + return out, nil +} +func (s *Service) activeProfile(ctx context.Context) (*ProfileSummary, error) { + var profile ProfileSummary + err := s.DB.SQL.QueryRowContext(ctx, `SELECT id,label FROM collection_profiles WHERE active=1 AND enabled=1`).Scan(&profile.ID, &profile.Label) + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("read active profile: %w", err) + } + return &profile, nil +} + +func (s *Service) loadRelay(ctx context.Context, out *RelaySummary) error { + var name string + var lastSeen sql.NullInt64 + var appVersion sql.NullString + err := s.DB.SQL.QueryRowContext(ctx, `SELECT COALESCE(name,''),last_seen_at,app_version FROM relay_devices WHERE enabled=1 LIMIT 1`). + Scan(&name, &lastSeen, &appVersion) + if errors.Is(err, sql.ErrNoRows) { + return nil + } + if err != nil { + return fmt.Errorf("read relay summary: %w", err) + } + out.Connected = true + out.Name = name + out.AppVersion = appVersion.String + if lastSeen.Valid { + value := time.UnixMilli(lastSeen.Int64).UTC() + out.LastSeenAt = &value + } + return nil +} + +func (s *Service) loadWebhookSummary(ctx context.Context, out *WebhookSummary) error { + if err := s.DB.SQL.QueryRowContext(ctx, `SELECT COUNT(*) FROM webhook_deliveries WHERE status IN ('pending','retry')`).Scan(&out.Pending); err != nil { + return fmt.Errorf("read pending webhook count: %w", err) + } + if err := s.DB.SQL.QueryRowContext(ctx, `SELECT COUNT(*) FROM webhook_deliveries WHERE status='exhausted'`).Scan(&out.Exhausted); err != nil { + return fmt.Errorf("read exhausted webhook count: %w", err) + } + var delivered sql.NullInt64 + if err := s.DB.SQL.QueryRowContext(ctx, `SELECT MAX(delivered_at) FROM webhook_deliveries WHERE status='delivered'`).Scan(&delivered); err != nil { + return fmt.Errorf("read last webhook delivery: %w", err) + } + if delivered.Valid { + value := time.UnixMilli(delivered.Int64).UTC() + out.LastDeliveredAt = &value + } + return nil +} +func (s *Service) Activity(ctx context.Context, limit int) ([]ActivityEntry, error) { + if err := s.ready(); err != nil { + return nil, err + } + if limit <= 0 { + limit = 100 + } + if limit > 500 { + limit = 500 + } + entries := make([]ActivityEntry, 0, limit*2) + if err := s.appendPaymentHistory(ctx, &entries, limit); err != nil { + return nil, err + } + if err := s.appendObservations(ctx, &entries, limit); err != nil { + return nil, err + } + if err := s.appendWebhooks(ctx, &entries, limit); err != nil { + return nil, err + } + sort.SliceStable(entries, func(i, j int) bool { + if entries[i].At.Equal(entries[j].At) { + return entries[i].Kind < entries[j].Kind + } + return entries[i].At.After(entries[j].At) + }) + if len(entries) > limit { + entries = entries[:limit] + } + return entries, nil +} + +func (s *Service) appendPaymentHistory(ctx context.Context, out *[]ActivityEntry, limit int) error { + rows, err := s.DB.SQL.QueryContext(ctx, `SELECT created_at,type,actor,summary,payment_id FROM payment_history ORDER BY created_at DESC,rowid DESC LIMIT ?`, limit) + if err != nil { + return fmt.Errorf("read payment activity: %w", err) + } + defer rows.Close() + for rows.Next() { + var at int64 + var entry ActivityEntry + if err := rows.Scan(&at, &entry.Status, &entry.Source, &entry.Title, &entry.PaymentID); err != nil { + return fmt.Errorf("scan payment activity: %w", err) + } + entry.At = time.UnixMilli(at).UTC() + entry.Kind = "payment" + *out = append(*out, entry) + } + return rows.Err() +} +func (s *Service) appendObservations(ctx context.Context, out *[]ActivityEntry, limit int) error { + rows, err := s.DB.SQL.QueryContext(ctx, `SELECT received_at,match_result,source,COALESCE(matched_payment_id,''),amount_paise,COALESCE(payer_name,'') + FROM payment_observations ORDER BY received_at DESC,rowid DESC LIMIT ?`, limit) + if err != nil { + return fmt.Errorf("read observation activity: %w", err) + } + defer rows.Close() + for rows.Next() { + var at, amount int64 + var entry ActivityEntry + var payer string + if err := rows.Scan(&at, &entry.Status, &entry.Source, &entry.PaymentID, &amount, &payer); err != nil { + return fmt.Errorf("scan observation activity: %w", err) + } + entry.At = time.UnixMilli(at).UTC() + entry.Kind = "payment_detected" + entry.Title = "Incoming payment detected" + entry.AmountPaise = &amount + entry.Detail = payer + *out = append(*out, entry) + } + return rows.Err() +} + +func (s *Service) appendWebhooks(ctx context.Context, out *[]ActivityEntry, limit int) error { + rows, err := s.DB.SQL.QueryContext(ctx, `SELECT COALESCE(delivered_at,created_at),status,event_type,payment_id,last_http_status,COALESCE(last_error,'') + FROM webhook_deliveries ORDER BY COALESCE(delivered_at,created_at) DESC,rowid DESC LIMIT ?`, limit) + if err != nil { + return fmt.Errorf("read webhook activity: %w", err) + } + defer rows.Close() + for rows.Next() { + var at int64 + var entry ActivityEntry + var httpStatus sql.NullInt64 + if err := rows.Scan(&at, &entry.Status, &entry.Source, &entry.PaymentID, &httpStatus, &entry.Detail); err != nil { + return fmt.Errorf("scan webhook activity: %w", err) + } + entry.At = time.UnixMilli(at).UTC() + entry.Kind = "webhook" + entry.Title = entry.Source + if httpStatus.Valid { + if entry.Detail != "" { + entry.Detail = fmt.Sprintf("HTTP %d ยท %s", httpStatus.Int64, entry.Detail) + } else { + entry.Detail = fmt.Sprintf("HTTP %d", httpStatus.Int64) + } + } + *out = append(*out, entry) + } + return rows.Err() +} + +func (s *Service) now() time.Time { + if s.Now == nil { + return time.Now().UTC() + } + return s.Now().UTC() +} + +func (s *Service) location() *time.Location { + if s.Location == nil { + return time.FixedZone("IST", 5*60*60+30*60) + } + return s.Location +} + +func localDayBounds(now time.Time, loc *time.Location) (time.Time, time.Time) { + local := now.In(loc) + startLocal := time.Date(local.Year(), local.Month(), local.Day(), 0, 0, 0, 0, loc) + return startLocal.UTC(), startLocal.AddDate(0, 0, 1).UTC() +} diff --git a/internal/v4/operator/service_test.go b/internal/v4/operator/service_test.go new file mode 100644 index 0000000..d1be214 --- /dev/null +++ b/internal/v4/operator/service_test.go @@ -0,0 +1,167 @@ +package operator + +import ( + "context" + "path/filepath" + "testing" + "time" + + "github.com/Phloraxx/payment-api/internal/v4/adminpayments" + "github.com/Phloraxx/payment-api/internal/v4/payments" + "github.com/Phloraxx/payment-api/internal/v4/profiles" + "github.com/Phloraxx/payment-api/internal/v4/storage" + "github.com/Phloraxx/payment-api/internal/v4/webhooks" +) + +type operatorFixture struct { + db *storage.DB + payments *payments.Service + admin *adminpayments.Service + operator *Service + now *time.Time +} + +func newOperatorFixture(t *testing.T) operatorFixture { + t.Helper() + db, err := storage.Open(context.Background(), filepath.Join(t.TempDir(), "paygate.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = db.Close() }) + now := time.Date(2026, 9, 1, 5, 0, 0, 0, time.UTC) + profileService := profiles.NewService(db) + ctx := context.Background() + if _, err := profileService.Upsert(ctx, profiles.UpsertInput{ + ID: "paytm", Label: "Paytm", UPIID: "paygate@paytm", PayeeName: "PayGate", + Parser: "paytm_notification", Enabled: true, + }); err != nil { + t.Fatal(err) + } + if _, err := profileService.Activate(ctx, "paytm"); err != nil { + t.Fatal(err) + } + paymentService := payments.NewService(db) + paymentService.Now = func() time.Time { return now } + adminService := adminpayments.NewService(db) + adminService.Now = func() time.Time { return now } + operatorService := NewService(db) + operatorService.Now = func() time.Time { return now } + return operatorFixture{db: db, payments: paymentService, admin: adminService, operator: operatorService, now: &now} +} + +func (f operatorFixture) create(t *testing.T, amount int64, idem string) payments.Payment { + t.Helper() + result, err := f.payments.Create(context.Background(), payments.CreateInput{ + RequestedAmountPaise: amount * 100, Name: "Sourav P Bijoy", ExternalID: "evt_test", + IdempotencyScope: "operator-test", IdempotencyKey: idem, + }) + if err != nil { + t.Fatal(err) + } + return result.Payment +} +func TestOverviewUsesIndiaLocalDayAndShowsOperationalSummary(t *testing.T) { + f := newOperatorFixture(t) + paidPayment := f.create(t, 100, "paid") + status := "paid" + if _, err := f.admin.Edit(context.Background(), paidPayment.ID, adminpayments.EditInput{Status: &status}); err != nil { + t.Fatal(err) + } + _ = f.create(t, 200, "pending") + if _, err := f.db.SQL.Exec(`INSERT INTO relay_devices(id,name,public_key_pem,enabled,enrolled_at,last_seen_at,app_version) + VALUES('device-1','Edge 60 Stylus','pem',1,?,?,?)`, f.now.Add(-time.Hour).UnixMilli(), f.now.Add(-time.Minute).UnixMilli(), "0.5.0"); err != nil { + t.Fatal(err) + } + + overview, err := f.operator.Overview(context.Background()) + if err != nil { + t.Fatal(err) + } + if overview.PaymentsToday != 2 || overview.PaidToday != 1 || overview.Pending != 1 { + t.Fatalf("overview counts = %+v", overview) + } + if overview.CollectedTodayPaise != paidPayment.PayableAmountPaise { + t.Fatalf("collected=%d want=%d", overview.CollectedTodayPaise, paidPayment.PayableAmountPaise) + } + if overview.ActiveProfile == nil || overview.ActiveProfile.ID != "paytm" { + t.Fatalf("active profile=%+v", overview.ActiveProfile) + } + if !overview.Relay.Connected || overview.Relay.Name != "Edge 60 Stylus" || overview.Relay.LastSeenAt == nil { + t.Fatalf("relay=%+v", overview.Relay) + } + if len(overview.Volume) != 7 || overview.Volume[6].Date != "2026-09-01" || overview.Volume[6].Payments != 1 { + t.Fatalf("volume=%+v", overview.Volume) + } +} +func TestActivityCombinesPaymentObservationAndWebhookEvents(t *testing.T) { + f := newOperatorFixture(t) + payment := f.create(t, 100, "activity") + if _, err := f.db.SQL.Exec(`INSERT INTO relay_devices(id,name,public_key_pem,enabled,enrolled_at) VALUES('device-1','Phone','pem',1,?)`, f.now.UnixMilli()); err != nil { + t.Fatal(err) + } + if _, err := f.db.SQL.Exec(`INSERT INTO relay_events(id,device_id,source_event_id,package_name,posted_at,received_at,amount_hint_paise,status) + VALUES('relay-1','device-1','source-1','com.paytm.business',?,?,?,?)`, f.now.UnixMilli(), f.now.UnixMilli(), payment.PayableAmountPaise, "matched"); err != nil { + t.Fatal(err) + } + if _, err := f.db.SQL.Exec(`INSERT INTO payment_observations(id,relay_event_id,source,collection_profile_id,amount_paise,payer_name,occurred_at,occurred_at_source,received_at,matched_payment_id,match_result) + VALUES('obs-1','relay-1','paytm_notification','paytm',?,'Bijoy P',?,'notification_posted_at',?,?,'matched')`, + payment.PayableAmountPaise, f.now.UnixMilli(), f.now.UnixMilli(), payment.ID); err != nil { + t.Fatal(err) + } + entries, err := f.operator.Activity(context.Background(), 20) + if err != nil { + t.Fatal(err) + } + var sawPayment, sawObservation, sawWebhook bool + for _, entry := range entries { + switch entry.Kind { + case "payment": + sawPayment = true + case "payment_detected": + sawObservation = true + case "webhook": + sawWebhook = true + } + } + if !sawPayment || !sawObservation || !sawWebhook { + t.Fatalf("activity kinds missing: %+v", entries) + } +} +func TestWebhookSettingsGenerateHideRotateAndApplyLive(t *testing.T) { + f := newOperatorFixture(t) + worker := webhooks.NewService(f.db, webhooks.Config{}) + settings := NewSettingsService(f.db, worker) + + current, err := settings.Webhook(context.Background()) + if err != nil || current.Enabled || current.SecretConfigured { + t.Fatalf("initial settings=%+v err=%v", current, err) + } + configured, secret, err := settings.ConfigureWebhook(context.Background(), "https://example.com/paygate-hook", false) + if err != nil { + t.Fatal(err) + } + if !configured.Enabled || !configured.SecretConfigured || secret == "" { + t.Fatalf("configured=%+v secret=%q", configured, secret) + } + if got := worker.ConfigSnapshot(); got.Endpoint != configured.Endpoint || got.Secret != secret { + t.Fatalf("worker config=%+v", got) + } + ordinary, err := settings.Webhook(context.Background()) + if err != nil || !ordinary.SecretConfigured { + t.Fatalf("ordinary settings=%+v err=%v", ordinary, err) + } + + rotated, newSecret, err := settings.ConfigureWebhook(context.Background(), configured.Endpoint, true) + if err != nil { + t.Fatal(err) + } + if !rotated.Enabled || newSecret == "" || newSecret == secret { + t.Fatalf("rotation secret old=%q new=%q settings=%+v", secret, newSecret, rotated) + } + if _, _, err := settings.ConfigureWebhook(context.Background(), "", false); err != nil { + t.Fatal(err) + } + if worker.Enabled() { + t.Fatal("worker remained enabled after webhook was disabled") + } +} diff --git a/internal/v4/operator/settings.go b/internal/v4/operator/settings.go new file mode 100644 index 0000000..909ba10 --- /dev/null +++ b/internal/v4/operator/settings.go @@ -0,0 +1,175 @@ +package operator + +import ( + "context" + "crypto/rand" + "database/sql" + "encoding/base64" + "errors" + "fmt" + "io" + "strings" + "time" + + "github.com/Phloraxx/payment-api/internal/v4/storage" + "github.com/Phloraxx/payment-api/internal/v4/webhooks" +) + +const ( + webhookEndpointKey = "webhook_endpoint" + webhookSecretKey = "webhook_secret" +) + +type SettingsService struct { + DB *storage.DB + Worker *webhooks.Service + Random io.Reader + Now func() time.Time +} + +type WebhookSettings struct { + Enabled bool + Endpoint string + SecretConfigured bool +} + +func NewSettingsService(db *storage.DB, worker *webhooks.Service) *SettingsService { + return &SettingsService{DB: db, Worker: worker, Random: rand.Reader, Now: time.Now} +} +func (s *SettingsService) Webhook(ctx context.Context) (WebhookSettings, error) { + cfg, err := s.loadWebhookConfig(ctx) + if err != nil { + return WebhookSettings{}, err + } + return WebhookSettings{ + Enabled: strings.TrimSpace(cfg.Endpoint) != "" && strings.TrimSpace(cfg.Secret) != "", + Endpoint: cfg.Endpoint, + SecretConfigured: strings.TrimSpace(cfg.Secret) != "", + }, nil +} + +func (s *SettingsService) ConfigureWebhook(ctx context.Context, endpoint string, rotateSecret bool) (WebhookSettings, string, error) { + if s == nil || s.DB == nil || s.DB.SQL == nil { + return WebhookSettings{}, "", errors.New("settings storage is required") + } + endpoint = strings.TrimSpace(endpoint) + if endpoint == "" { + if err := s.storeWebhookConfig(ctx, webhooks.Config{}); err != nil { + return WebhookSettings{}, "", err + } + if s.Worker != nil { + _ = s.Worker.UpdateConfig(webhooks.Config{}) + } + return WebhookSettings{}, "", nil + } + existing, err := s.loadWebhookConfig(ctx) + if err != nil { + return WebhookSettings{}, "", err + } + secret := existing.Secret + newSecret := "" + if secret == "" || rotateSecret { + secret, err = s.generateWebhookSecret() + if err != nil { + return WebhookSettings{}, "", err + } + newSecret = secret + } + cfg := webhooks.Config{Endpoint: endpoint, Secret: secret} + if err := webhooks.ValidateConfig(cfg); err != nil { + return WebhookSettings{}, "", err + } + if err := s.storeWebhookConfig(ctx, cfg); err != nil { + return WebhookSettings{}, "", err + } + if s.Worker != nil { + if err := s.Worker.UpdateConfig(cfg); err != nil { + return WebhookSettings{}, "", err + } + } + return WebhookSettings{Enabled: true, Endpoint: endpoint, SecretConfigured: true}, newSecret, nil +} +func (s *SettingsService) ApplyPersistedWebhook(ctx context.Context) error { + if s == nil || s.Worker == nil { + return nil + } + cfg, err := s.loadWebhookConfig(ctx) + if err != nil { + return err + } + if strings.TrimSpace(cfg.Endpoint) == "" && strings.TrimSpace(cfg.Secret) == "" { + return s.Worker.UpdateConfig(webhooks.Config{}) + } + return s.Worker.UpdateConfig(cfg) +} + +func (s *SettingsService) loadWebhookConfig(ctx context.Context) (webhooks.Config, error) { + if s == nil || s.DB == nil || s.DB.SQL == nil { + return webhooks.Config{}, errors.New("settings storage is required") + } + endpoint, err := readSetting(ctx, s.DB.SQL, webhookEndpointKey) + if err != nil { + return webhooks.Config{}, err + } + secret, err := readSetting(ctx, s.DB.SQL, webhookSecretKey) + if err != nil { + return webhooks.Config{}, err + } + return webhooks.Config{Endpoint: endpoint, Secret: secret}, nil +} +func (s *SettingsService) storeWebhookConfig(ctx context.Context, cfg webhooks.Config) error { + now := time.Now().UTC() + if s.Now != nil { + now = s.Now().UTC() + } + return s.DB.WithImmediateTx(ctx, func(tx *storage.ImmediateTx) error { + if err := writeSetting(ctx, tx, webhookEndpointKey, strings.TrimSpace(cfg.Endpoint), now); err != nil { + return err + } + if err := writeSetting(ctx, tx, webhookSecretKey, strings.TrimSpace(cfg.Secret), now); err != nil { + return err + } + return nil + }) +} + +func (s *SettingsService) generateWebhookSecret() (string, error) { + r := s.Random + if r == nil { + r = rand.Reader + } + raw := make([]byte, 32) + if _, err := io.ReadFull(r, raw); err != nil { + return "", fmt.Errorf("generate webhook secret: %w", err) + } + return "whsec_" + base64.RawURLEncoding.EncodeToString(raw), nil +} + +type settingQuerier interface { + QueryRowContext(context.Context, string, ...any) *sql.Row +} + +type settingExecer interface { + ExecContext(context.Context, string, ...any) (sql.Result, error) +} + +func readSetting(ctx context.Context, db settingQuerier, key string) (string, error) { + var value string + err := db.QueryRowContext(ctx, `SELECT value FROM settings WHERE key=?`, key).Scan(&value) + if errors.Is(err, sql.ErrNoRows) { + return "", nil + } + if err != nil { + return "", fmt.Errorf("read setting %s: %w", key, err) + } + return value, nil +} + +func writeSetting(ctx context.Context, db settingExecer, key, value string, at time.Time) error { + _, err := db.ExecContext(ctx, `INSERT INTO settings(key,value,updated_at) VALUES(?,?,?) + ON CONFLICT(key) DO UPDATE SET value=excluded.value,updated_at=excluded.updated_at`, key, value, at.UnixMilli()) + if err != nil { + return fmt.Errorf("write setting %s: %w", key, err) + } + return nil +} diff --git a/internal/v4/webhooks/service.go b/internal/v4/webhooks/service.go index 5d2e65a..577b949 100644 --- a/internal/v4/webhooks/service.go +++ b/internal/v4/webhooks/service.go @@ -33,7 +33,8 @@ type Config struct { } type Service struct { DB *storage.DB - Config Config + config Config + configMu sync.RWMutex HTTPClient *http.Client Now func() time.Time MaxAttempts int @@ -52,14 +53,40 @@ type Delivery struct { func NewService(db *storage.DB, cfg Config) *Service { return &Service{ - DB: db, Config: cfg, HTTPClient: newHTTPClient(), Now: time.Now, + DB: db, config: cfg, HTTPClient: newHTTPClient(), Now: time.Now, MaxAttempts: defaultMaxAttempts, BatchSize: defaultBatchSize, Lease: defaultLease, wake: make(chan struct{}, 1), } } func (s *Service) Enabled() bool { - return s != nil && strings.TrimSpace(s.Config.Endpoint) != "" && strings.TrimSpace(s.Config.Secret) != "" + if s == nil { + return false + } + cfg := s.ConfigSnapshot() + return strings.TrimSpace(cfg.Endpoint) != "" && strings.TrimSpace(cfg.Secret) != "" +} + +func (s *Service) ConfigSnapshot() Config { + if s == nil { + return Config{} + } + s.configMu.RLock() + defer s.configMu.RUnlock() + return s.config +} + +func (s *Service) UpdateConfig(cfg Config) error { + if strings.TrimSpace(cfg.Endpoint) != "" || strings.TrimSpace(cfg.Secret) != "" { + if err := ValidateConfig(cfg); err != nil { + return err + } + } + s.configMu.Lock() + s.config = cfg + s.configMu.Unlock() + s.Wake() + return nil } func (s *Service) Wake() { if !s.Enabled() { @@ -72,9 +99,6 @@ func (s *Service) Wake() { } func (s *Service) Run(ctx context.Context) { - if !s.Enabled() { - return - } ticker := time.NewTicker(30 * time.Second) defer ticker.Stop() s.Wake() @@ -90,10 +114,14 @@ func (s *Service) Run(ctx context.Context) { } func (s *Service) SendPending(ctx context.Context) (int, error) { - if !s.Enabled() { + if err := s.readyStorage(); err != nil { + return 0, err + } + cfg := s.ConfigSnapshot() + if strings.TrimSpace(cfg.Endpoint) == "" && strings.TrimSpace(cfg.Secret) == "" { return 0, nil } - if err := s.ready(); err != nil { + if err := ValidateConfig(cfg); err != nil { return 0, err } s.mu.Lock() @@ -137,7 +165,7 @@ func (s *Service) SendPending(ctx context.Context) (int, error) { if delivery == nil { continue } - s.deliver(ctx, *delivery) + s.deliver(ctx, cfg, *delivery) processed++ } return processed, nil @@ -174,11 +202,11 @@ func (s *Service) claim(ctx context.Context, id string, now time.Time) (*Deliver return out, err } -func (s *Service) deliver(ctx context.Context, delivery Delivery) { +func (s *Service) deliver(ctx context.Context, cfg Config, delivery Delivery) { now := s.now() timestamp := strconv.FormatInt(now.Unix(), 10) - signature := Sign(s.Config.Secret, timestamp, []byte(delivery.Body)) - req, err := http.NewRequestWithContext(ctx, http.MethodPost, s.Config.Endpoint, strings.NewReader(delivery.Body)) + signature := Sign(cfg.Secret, timestamp, []byte(delivery.Body)) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, cfg.Endpoint, strings.NewReader(delivery.Body)) if err == nil { req.Header.Set("Content-Type", "application/json") req.Header.Set("User-Agent", "PayGate/4") @@ -229,7 +257,7 @@ func (s *Service) finish(ctx context.Context, delivery Delivery, statusCode int, }) } func (s *Service) RetryOne(ctx context.Context, id string) error { - if err := s.ready(); err != nil { + if err := s.readyStorage(); err != nil { return err } id = strings.TrimSpace(id) @@ -249,12 +277,16 @@ func (s *Service) RetryOne(ctx context.Context, id string) error { return nil } -func (s *Service) ready() error { +func (s *Service) readyStorage() error { if s == nil || s.DB == nil || s.DB.SQL == nil { return errors.New("webhook storage is required") } - endpoint := strings.TrimSpace(s.Config.Endpoint) - secret := strings.TrimSpace(s.Config.Secret) + return nil +} + +func ValidateConfig(cfg Config) error { + endpoint := strings.TrimSpace(cfg.Endpoint) + secret := strings.TrimSpace(cfg.Secret) if endpoint == "" || len(secret) < 32 { return fmt.Errorf("%w: endpoint and at least 32-byte secret are required", ErrInvalidConfig) } @@ -262,7 +294,7 @@ func (s *Service) ready() error { if err != nil || u.Host == "" || u.User != nil || u.RawQuery != "" || u.Fragment != "" { return fmt.Errorf("%w: endpoint must be an absolute URL without credentials/query/fragment", ErrInvalidConfig) } - if u.Scheme != "https" && !(s.Config.AllowInsecureHTTP && u.Scheme == "http") { + if u.Scheme != "https" && !(cfg.AllowInsecureHTTP && u.Scheme == "http") { return fmt.Errorf("%w: HTTPS endpoint is required", ErrInvalidConfig) } return nil