diff --git a/backend/cmd/server/VERSION b/backend/cmd/server/VERSION index afecc7a0a2..417cd44107 100644 --- a/backend/cmd/server/VERSION +++ b/backend/cmd/server/VERSION @@ -1 +1 @@ -0.1.108.138 +0.1.108.139 diff --git a/backend/internal/payment/load_balancer.go b/backend/internal/payment/load_balancer.go index a36644b89d..afe607e0cf 100644 --- a/backend/internal/payment/load_balancer.go +++ b/backend/internal/payment/load_balancer.go @@ -50,22 +50,56 @@ func NewDefaultLoadBalancer(db *dbent.Client, encryptionKey []byte) *DefaultLoad return &DefaultLoadBalancer{db: db, encryptionKey: encryptionKey} } -// SelectInstance picks an enabled instance for the given provider key and payment type. -func (lb *DefaultLoadBalancer) SelectInstance(ctx context.Context, providerKey string, paymentType PaymentType, strategy Strategy, orderAmount float64) (*InstanceSelection, error) { - candidates, err := lb.getCandidates(ctx, providerKey, paymentType, orderAmount) - if err != nil { - return nil, err - } - - selected, err := lb.pickByStrategy(ctx, candidates, strategy) - if err != nil { - return nil, err - } - - return lb.buildSelection(selected) +// instanceCandidate pairs an instance with its pre-fetched daily usage. +type instanceCandidate struct { + inst *dbent.PaymentProviderInstance + dailyUsed float64 // includes PENDING orders } -func (lb *DefaultLoadBalancer) getCandidates(ctx context.Context, providerKey string, paymentType PaymentType, orderAmount float64) ([]*dbent.PaymentProviderInstance, error) { +// SelectInstance picks an enabled instance for the given provider key and payment type. +// +// Flow: +// 1. Query all enabled instances for providerKey, filter by supported paymentType +// 2. Batch-query daily usage (PENDING + PAID + COMPLETED + RECHARGING) for all candidates +// 3. Filter out instances where: single-min/max violated OR daily remaining < orderAmount +// 4. Pick from survivors using the configured strategy (round-robin / least-amount) +// 5. If all filtered out, fall back to full list (let the provider itself reject) +func (lb *DefaultLoadBalancer) SelectInstance( + ctx context.Context, + providerKey string, + paymentType PaymentType, + strategy Strategy, + orderAmount float64, +) (*InstanceSelection, error) { + // Step 1: query enabled instances matching payment type. + instances, err := lb.queryEnabledInstances(ctx, providerKey, paymentType) + if err != nil { + return nil, err + } + + // Step 2: batch-fetch daily usage for all candidates. + candidates := lb.attachDailyUsage(ctx, instances) + + // Step 3: filter by limits. + available := filterByLimits(candidates, paymentType, orderAmount) + if len(available) == 0 { + slog.Warn("all instances exceeded limits, using full candidate list", + "provider", providerKey, "payment_type", paymentType, + "order_amount", orderAmount, "count", len(candidates)) + available = candidates + } + + // Step 4: pick by strategy. + selected := lb.pickByStrategy(available, strategy) + return lb.buildSelection(selected.inst) +} + +// queryEnabledInstances returns enabled instances for providerKey that support paymentType. +func (lb *DefaultLoadBalancer) queryEnabledInstances( + ctx context.Context, + providerKey string, + paymentType PaymentType, +) ([]*dbent.PaymentProviderInstance, error) { instances, err := lb.db.PaymentProviderInstance.Query(). Where( paymentproviderinstance.ProviderKey(providerKey), @@ -77,72 +111,97 @@ func (lb *DefaultLoadBalancer) getCandidates(ctx context.Context, providerKey st return nil, fmt.Errorf("query provider instances: %w", err) } - // Filter by supported types. - // When paymentType equals providerKey (e.g. "stripe"), all instances of that - // provider are candidates — the sub-type filtering is handled internally. - var candidates []*dbent.PaymentProviderInstance + var matched []*dbent.PaymentProviderInstance for _, inst := range instances { if paymentType == providerKey || InstanceSupportsType(inst.SupportedTypes, paymentType) { - candidates = append(candidates, inst) + matched = append(matched, inst) } } - if len(candidates) == 0 { + if len(matched) == 0 { return nil, fmt.Errorf("no enabled instance for provider %s type %s", providerKey, paymentType) } + return matched, nil +} - // Exclude instances that cannot accommodate this order (daily limit or single amount range). - filtered := lb.filterByLimits(ctx, candidates, paymentType, orderAmount) - if len(filtered) > 0 { - return filtered, nil +// attachDailyUsage queries daily usage for each instance in a single pass. +// Usage includes PENDING orders to avoid over-committing capacity. +func (lb *DefaultLoadBalancer) attachDailyUsage( + ctx context.Context, + instances []*dbent.PaymentProviderInstance, +) []instanceCandidate { + todayStart := startOfDay(time.Now()) + + // Collect instance IDs. + ids := make([]string, len(instances)) + for i, inst := range instances { + ids[i] = fmt.Sprintf("%d", inst.ID) } - // All instances exhausted — fall back to full candidate list so the order - // can still be attempted (the provider itself may reject it). - slog.Warn("all instances exceeded limits, using full candidate list", - "provider", providerKey, "payment_type", paymentType, - "order_amount", orderAmount, "count", len(candidates)) - return candidates, nil + + // Batch query: sum pay_amount grouped by provider_instance_id. + type row struct { + InstanceID string `json:"provider_instance_id"` + Sum float64 `json:"sum"` + } + var rows []row + err := lb.db.PaymentOrder.Query(). + Where( + paymentorder.ProviderInstanceIDIn(ids...), + paymentorder.StatusIn( + OrderStatusPending, OrderStatusPaid, + OrderStatusCompleted, OrderStatusRecharging, + ), + paymentorder.CreatedAtGTE(todayStart), + ). + GroupBy(paymentorder.FieldProviderInstanceID). + Aggregate(dbent.Sum(paymentorder.FieldPayAmount)). + Scan(ctx, &rows) + if err != nil { + slog.Warn("batch daily usage query failed, treating all as zero", "error", err) + } + + usageMap := make(map[string]float64, len(rows)) + for _, r := range rows { + usageMap[r.InstanceID] = r.Sum + } + + candidates := make([]instanceCandidate, len(instances)) + for i, inst := range instances { + candidates[i] = instanceCandidate{ + inst: inst, + dailyUsed: usageMap[fmt.Sprintf("%d", inst.ID)], + } + } + return candidates } // filterByLimits removes instances that cannot accommodate the order: -// - daily remaining capacity < orderAmount -// - orderAmount outside single-transaction min/max range -func (lb *DefaultLoadBalancer) filterByLimits(ctx context.Context, candidates []*dbent.PaymentProviderInstance, paymentType PaymentType, orderAmount float64) []*dbent.PaymentProviderInstance { - var filtered []*dbent.PaymentProviderInstance - for _, inst := range candidates { - cl := getInstanceChannelLimits(inst, paymentType) +// - orderAmount outside single-transaction [min, max] +// - daily remaining capacity (limit - used) < orderAmount +func filterByLimits(candidates []instanceCandidate, paymentType PaymentType, orderAmount float64) []instanceCandidate { + var result []instanceCandidate + for _, c := range candidates { + cl := getInstanceChannelLimits(c.inst, paymentType) - // Check single-transaction range. if cl.SingleMin > 0 && orderAmount < cl.SingleMin { slog.Info("order below instance single min, skipping", - "instance_id", inst.ID, "order", orderAmount, "min", cl.SingleMin) + "instance_id", c.inst.ID, "order", orderAmount, "min", cl.SingleMin) continue } if cl.SingleMax > 0 && orderAmount > cl.SingleMax { slog.Info("order above instance single max, skipping", - "instance_id", inst.ID, "order", orderAmount, "max", cl.SingleMax) + "instance_id", c.inst.ID, "order", orderAmount, "max", cl.SingleMax) + continue + } + if cl.DailyLimit > 0 && c.dailyUsed+orderAmount > cl.DailyLimit { + slog.Info("instance daily remaining insufficient, skipping", + "instance_id", c.inst.ID, "used", c.dailyUsed, + "order", orderAmount, "limit", cl.DailyLimit) continue } - // Check daily remaining capacity. - if cl.DailyLimit > 0 { - used, err := lb.GetInstanceDailyAmount(ctx, fmt.Sprintf("%d", inst.ID)) - if err != nil { - slog.Warn("failed to query daily amount, keeping instance", - "instance_id", inst.ID, "error", err) - filtered = append(filtered, inst) - continue - } - if used+orderAmount > cl.DailyLimit { - slog.Info("instance daily remaining insufficient, skipping", - "instance_id", inst.ID, "used", used, - "order", orderAmount, "limit", cl.DailyLimit) - continue - } - } - - filtered = append(filtered, inst) + result = append(result, c) } - return filtered + return result } // getInstanceChannelLimits returns the channel limits for a specific payment type. @@ -165,30 +224,26 @@ func getInstanceChannelLimits(inst *dbent.PaymentProviderInstance, paymentType P return ChannelLimits{} } -func (lb *DefaultLoadBalancer) pickByStrategy(ctx context.Context, candidates []*dbent.PaymentProviderInstance, strategy Strategy) (*dbent.PaymentProviderInstance, error) { +// pickByStrategy selects one instance from the available candidates. +func (lb *DefaultLoadBalancer) pickByStrategy(candidates []instanceCandidate, strategy Strategy) instanceCandidate { if strategy == StrategyLeastAmount && len(candidates) > 1 { - return lb.pickLeastAmount(ctx, candidates) + return pickLeastAmount(candidates) } - // Default: round-robin + // Default: round-robin. idx := lb.counter.Add(1) % uint64(len(candidates)) - return candidates[idx], nil + return candidates[idx] } -func (lb *DefaultLoadBalancer) pickLeastAmount(ctx context.Context, candidates []*dbent.PaymentProviderInstance) (*dbent.PaymentProviderInstance, error) { - var selected *dbent.PaymentProviderInstance - minAmount := -1.0 - for _, c := range candidates { - amount, err := lb.GetInstanceDailyAmount(ctx, fmt.Sprintf("%d", c.ID)) - if err != nil { - slog.Warn("failed to get instance daily amount, falling back", "instance_id", c.ID, "error", err) - amount = 0 - } - if minAmount < 0 || amount < minAmount { - minAmount = amount - selected = c +// pickLeastAmount selects the instance with the lowest daily usage. +// No extra DB queries — usage was pre-fetched in attachDailyUsage. +func pickLeastAmount(candidates []instanceCandidate) instanceCandidate { + best := candidates[0] + for _, c := range candidates[1:] { + if c.dailyUsed < best.dailyUsed { + best = c } } - return selected, nil + return best } func (lb *DefaultLoadBalancer) buildSelection(selected *dbent.PaymentProviderInstance) (*InstanceSelection, error) { @@ -197,7 +252,6 @@ func (lb *DefaultLoadBalancer) buildSelection(selected *dbent.PaymentProviderIns return nil, fmt.Errorf("decrypt instance %d config: %w", selected.ID, err) } - // Inject payment_mode into config so providers can read it uniformly if selected.PaymentMode != "" { config["paymentMode"] = selected.PaymentMode } @@ -224,8 +278,7 @@ func (lb *DefaultLoadBalancer) decryptConfig(encrypted string) (map[string]strin // GetInstanceDailyAmount returns the total completed order amount for an instance today. func (lb *DefaultLoadBalancer) GetInstanceDailyAmount(ctx context.Context, instanceID string) (float64, error) { - now := time.Now() - todayStart := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, now.Location()) + todayStart := startOfDay(time.Now()) var result []struct { Sum float64 `json:"sum"` @@ -247,6 +300,10 @@ func (lb *DefaultLoadBalancer) GetInstanceDailyAmount(ctx context.Context, insta return 0, nil } +func startOfDay(t time.Time) time.Time { + return time.Date(t.Year(), t.Month(), t.Day(), 0, 0, 0, 0, t.Location()) +} + // InstanceSupportsType checks if the given supported types string includes the target type. // An empty supportedTypes string means all types are supported. func InstanceSupportsType(supportedTypes string, target PaymentType) bool {