mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-10-07 15:57:53 +08:00
feat(admin): 为账号批量删除增加并发限制
This commit is contained in:
@@ -1531,6 +1531,141 @@ func (h *AccountHandler) RevertProxyFallback(c *gin.Context) {
|
||||
response.Success(c, gin.H{"message": "reverted"})
|
||||
}
|
||||
|
||||
// BatchDelete handles deleting multiple accounts with bounded concurrency.
|
||||
// POST /api/v1/admin/accounts/batch-delete
|
||||
func (h *AccountHandler) BatchDelete(c *gin.Context) {
|
||||
var req struct {
|
||||
AccountIDs []int64 `json:"account_ids"`
|
||||
}
|
||||
if err := c.ShouldBindJSON(&req); err != nil {
|
||||
response.BadRequest(c, "Invalid request: "+err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
accountIDs := normalizeInt64IDList(req.AccountIDs)
|
||||
if len(accountIDs) == 0 {
|
||||
response.BadRequest(c, "account_ids is required")
|
||||
return
|
||||
}
|
||||
|
||||
accounts, err := h.adminService.GetAccountsByIDs(c.Request.Context(), accountIDs)
|
||||
if err != nil {
|
||||
response.ErrorFrom(c, err)
|
||||
return
|
||||
}
|
||||
|
||||
type deleteError struct {
|
||||
AccountID int64 `json:"account_id"`
|
||||
Error string `json:"error"`
|
||||
}
|
||||
|
||||
requestedIDs := make(map[int64]struct{}, len(accountIDs))
|
||||
for _, accountID := range accountIDs {
|
||||
requestedIDs[accountID] = struct{}{}
|
||||
}
|
||||
accountsByID := make(map[int64]*service.Account, len(accounts))
|
||||
for _, account := range accounts {
|
||||
if account != nil {
|
||||
accountsByID[account.ID] = account
|
||||
}
|
||||
}
|
||||
|
||||
rootIDs := make([]int64, 0, len(accountIDs))
|
||||
dependentIDs := make(map[int64][]int64)
|
||||
failedIDs := make([]int64, 0)
|
||||
errorsByAccount := make([]deleteError, 0)
|
||||
for _, accountID := range accountIDs {
|
||||
account := accountsByID[accountID]
|
||||
if account == nil {
|
||||
failedIDs = append(failedIDs, accountID)
|
||||
errorsByAccount = append(errorsByAccount, deleteError{
|
||||
AccountID: accountID,
|
||||
Error: "account not found",
|
||||
})
|
||||
continue
|
||||
}
|
||||
|
||||
rootID := accountID
|
||||
visited := map[int64]struct{}{accountID: {}}
|
||||
for {
|
||||
current := accountsByID[rootID]
|
||||
if current == nil || current.ParentAccountID == nil {
|
||||
break
|
||||
}
|
||||
parentID := *current.ParentAccountID
|
||||
if _, selected := requestedIDs[parentID]; !selected {
|
||||
break
|
||||
}
|
||||
if _, exists := accountsByID[parentID]; !exists {
|
||||
break
|
||||
}
|
||||
if _, cyclic := visited[parentID]; cyclic {
|
||||
rootID = accountID
|
||||
break
|
||||
}
|
||||
visited[parentID] = struct{}{}
|
||||
rootID = parentID
|
||||
}
|
||||
|
||||
if rootID != accountID {
|
||||
dependentIDs[rootID] = append(dependentIDs[rootID], accountID)
|
||||
continue
|
||||
}
|
||||
rootIDs = append(rootIDs, accountID)
|
||||
}
|
||||
|
||||
const maxConcurrency = 5
|
||||
g, gctx := errgroup.WithContext(c.Request.Context())
|
||||
g.SetLimit(maxConcurrency)
|
||||
|
||||
var mu sync.Mutex
|
||||
successIDs := make([]int64, 0, len(accountIDs))
|
||||
|
||||
// Every worker returns nil so one account failure does not cancel the remaining deletions.
|
||||
for _, id := range rootIDs {
|
||||
accountID := id
|
||||
g.Go(func() error {
|
||||
err := h.adminService.DeleteAccount(gctx, accountID)
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
affectedIDs := append([]int64{accountID}, dependentIDs[accountID]...)
|
||||
if err != nil {
|
||||
for _, affectedID := range affectedIDs {
|
||||
failedIDs = append(failedIDs, affectedID)
|
||||
errorsByAccount = append(errorsByAccount, deleteError{
|
||||
AccountID: affectedID,
|
||||
Error: err.Error(),
|
||||
})
|
||||
}
|
||||
return nil
|
||||
}
|
||||
successIDs = append(successIDs, affectedIDs...)
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
if err := g.Wait(); err != nil {
|
||||
response.ErrorFrom(c, err)
|
||||
return
|
||||
}
|
||||
|
||||
sort.Slice(successIDs, func(i, j int) bool { return successIDs[i] < successIDs[j] })
|
||||
sort.Slice(failedIDs, func(i, j int) bool { return failedIDs[i] < failedIDs[j] })
|
||||
sort.Slice(errorsByAccount, func(i, j int) bool {
|
||||
return errorsByAccount[i].AccountID < errorsByAccount[j].AccountID
|
||||
})
|
||||
|
||||
response.Success(c, gin.H{
|
||||
"total": len(accountIDs),
|
||||
"success": len(successIDs),
|
||||
"failed": len(failedIDs),
|
||||
"success_ids": successIDs,
|
||||
"failed_ids": failedIDs,
|
||||
"errors": errorsByAccount,
|
||||
})
|
||||
}
|
||||
|
||||
// BatchClearError handles batch clearing account errors
|
||||
// POST /api/v1/admin/accounts/batch-clear-error
|
||||
func (h *AccountHandler) BatchClearError(c *gin.Context) {
|
||||
|
||||
@@ -0,0 +1,173 @@
|
||||
package admin
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/Wei-Shaw/sub2api/internal/service"
|
||||
)
|
||||
|
||||
type batchDeleteAdminService struct {
|
||||
*stubAdminService
|
||||
|
||||
mu sync.Mutex
|
||||
active int
|
||||
maxActive int
|
||||
deletedIDs []int64
|
||||
deleteErrorsByID map[int64]error
|
||||
accountsByID map[int64]*service.Account
|
||||
}
|
||||
|
||||
func (s *batchDeleteAdminService) GetAccountsByIDs(_ context.Context, ids []int64) ([]*service.Account, error) {
|
||||
accounts := make([]*service.Account, 0, len(ids))
|
||||
for _, id := range ids {
|
||||
if s.accountsByID != nil {
|
||||
if account, ok := s.accountsByID[id]; ok {
|
||||
accounts = append(accounts, account)
|
||||
}
|
||||
continue
|
||||
}
|
||||
accounts = append(accounts, &service.Account{ID: id})
|
||||
}
|
||||
return accounts, nil
|
||||
}
|
||||
|
||||
func (s *batchDeleteAdminService) DeleteAccount(ctx context.Context, id int64) error {
|
||||
s.mu.Lock()
|
||||
s.active++
|
||||
if s.active > s.maxActive {
|
||||
s.maxActive = s.active
|
||||
}
|
||||
s.mu.Unlock()
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
s.mu.Lock()
|
||||
s.active--
|
||||
s.mu.Unlock()
|
||||
return ctx.Err()
|
||||
case <-time.After(10 * time.Millisecond):
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.active--
|
||||
s.deletedIDs = append(s.deletedIDs, id)
|
||||
return s.deleteErrorsByID[id]
|
||||
}
|
||||
|
||||
func setupAccountBatchDeleteRouter(adminSvc *batchDeleteAdminService) *gin.Engine {
|
||||
gin.SetMode(gin.TestMode)
|
||||
router := gin.New()
|
||||
handler := NewAccountHandler(adminSvc, nil, nil, nil, nil, nil, nil, nil, nil, nil, nil, nil, nil, nil)
|
||||
router.POST("/api/v1/admin/accounts/batch-delete", handler.BatchDelete)
|
||||
return router
|
||||
}
|
||||
|
||||
func TestAccountHandlerBatchDeleteReturnsStablePerAccountResults(t *testing.T) {
|
||||
adminSvc := &batchDeleteAdminService{
|
||||
stubAdminService: newStubAdminService(),
|
||||
deleteErrorsByID: map[int64]error{
|
||||
3: errors.New("delete failed"),
|
||||
},
|
||||
}
|
||||
router := setupAccountBatchDeleteRouter(adminSvc)
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
req := httptest.NewRequest(
|
||||
http.MethodPost,
|
||||
"/api/v1/admin/accounts/batch-delete",
|
||||
bytes.NewBufferString(`{"account_ids":[5,4,3,2,1,2,0,-1]}`),
|
||||
)
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
router.ServeHTTP(rec, req)
|
||||
|
||||
require.Equal(t, http.StatusOK, rec.Code)
|
||||
|
||||
var payload struct {
|
||||
Data struct {
|
||||
Total int `json:"total"`
|
||||
Success int `json:"success"`
|
||||
Failed int `json:"failed"`
|
||||
SuccessIDs []int64 `json:"success_ids"`
|
||||
FailedIDs []int64 `json:"failed_ids"`
|
||||
Errors []struct {
|
||||
AccountID int64 `json:"account_id"`
|
||||
Error string `json:"error"`
|
||||
} `json:"errors"`
|
||||
} `json:"data"`
|
||||
}
|
||||
require.NoError(t, json.Unmarshal(rec.Body.Bytes(), &payload))
|
||||
require.Equal(t, 5, payload.Data.Total)
|
||||
require.Equal(t, 4, payload.Data.Success)
|
||||
require.Equal(t, 1, payload.Data.Failed)
|
||||
require.Equal(t, []int64{1, 2, 4, 5}, payload.Data.SuccessIDs)
|
||||
require.Equal(t, []int64{3}, payload.Data.FailedIDs)
|
||||
require.Equal(t, int64(3), payload.Data.Errors[0].AccountID)
|
||||
require.Equal(t, "delete failed", payload.Data.Errors[0].Error)
|
||||
require.LessOrEqual(t, adminSvc.maxActive, 5)
|
||||
require.Greater(t, adminSvc.maxActive, 1)
|
||||
}
|
||||
|
||||
func TestAccountHandlerBatchDeleteDoesNotRaceSelectedShadowWithParent(t *testing.T) {
|
||||
parentID := int64(1)
|
||||
adminSvc := &batchDeleteAdminService{
|
||||
stubAdminService: newStubAdminService(),
|
||||
accountsByID: map[int64]*service.Account{
|
||||
1: {ID: 1},
|
||||
2: {ID: 2, ParentAccountID: &parentID},
|
||||
3: {ID: 3},
|
||||
},
|
||||
}
|
||||
router := setupAccountBatchDeleteRouter(adminSvc)
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
req := httptest.NewRequest(
|
||||
http.MethodPost,
|
||||
"/api/v1/admin/accounts/batch-delete",
|
||||
bytes.NewBufferString(`{"account_ids":[1,2,3]}`),
|
||||
)
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
router.ServeHTTP(rec, req)
|
||||
|
||||
require.Equal(t, http.StatusOK, rec.Code)
|
||||
|
||||
var payload struct {
|
||||
Data struct {
|
||||
SuccessIDs []int64 `json:"success_ids"`
|
||||
FailedIDs []int64 `json:"failed_ids"`
|
||||
} `json:"data"`
|
||||
}
|
||||
require.NoError(t, json.Unmarshal(rec.Body.Bytes(), &payload))
|
||||
require.Equal(t, []int64{1, 2, 3}, payload.Data.SuccessIDs)
|
||||
require.Empty(t, payload.Data.FailedIDs)
|
||||
require.ElementsMatch(t, []int64{1, 3}, adminSvc.deletedIDs)
|
||||
}
|
||||
|
||||
func TestAccountHandlerBatchDeleteRejectsEmptyNormalizedIDs(t *testing.T) {
|
||||
adminSvc := &batchDeleteAdminService{
|
||||
stubAdminService: newStubAdminService(),
|
||||
}
|
||||
router := setupAccountBatchDeleteRouter(adminSvc)
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
req := httptest.NewRequest(
|
||||
http.MethodPost,
|
||||
"/api/v1/admin/accounts/batch-delete",
|
||||
bytes.NewBufferString(`{"account_ids":[0,-1]}`),
|
||||
)
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
router.ServeHTTP(rec, req)
|
||||
|
||||
require.Equal(t, http.StatusBadRequest, rec.Code)
|
||||
}
|
||||
@@ -395,6 +395,7 @@ func registerAccountRoutes(admin *gin.RouterGroup, h *handler.Handlers, stepUpAu
|
||||
accounts.POST("/batch-update-credentials", h.Admin.Account.BatchUpdateCredentials)
|
||||
accounts.POST("/batch-refresh-tier", h.Admin.Account.BatchRefreshTier)
|
||||
accounts.POST("/bulk-update", h.Admin.Account.BulkUpdate)
|
||||
accounts.POST("/batch-delete", h.Admin.Account.BatchDelete)
|
||||
accounts.POST("/batch-clear-error", h.Admin.Account.BatchClearError)
|
||||
accounts.POST("/batch-refresh", h.Admin.Account.BatchRefresh)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user