mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-09-24 16:05:44 +08:00
refactor(payment): rewrite load balancer with complete limit checking
- Batch-query daily usage in one SQL (GroupBy) instead of N queries - Include PENDING orders in daily usage to prevent over-committing - Check remaining capacity (used + orderAmount > limit) not just used >= limit - Check SingleMin/SingleMax per-channel limits - Pre-attach usage to candidates, reuse across filter and strategy pick - Fallback to full list with warn log when all instances exhausted chore: bump version to 0.1.108.139
This commit is contained in:
@@ -1 +1 @@
|
||||
0.1.108.138
|
||||
0.1.108.139
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user