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