mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-06 07:36:37 +08:00
refactor(service): adopt feature-based architecture and rename pkg/diskcache
- Moved GORM/Redis CAPTCHA manager from internal/service/cap directly into the cohesive CAPTCHA app folder at internal/apps/cap/. - Moved background system cleanup handler from internal/service/cleanup.go into internal/apps/upload/cleanup.go. - Completely removed the global internal/service directory to keep module logic self-contained. - Renamed the core utility engine pkg/diskcache to pkg/cache/disk to separate underlying utility code from db/config integrations. - Renamed DiskCache struct in pkg/cache/disk to Cache to resolve revive package-name stuttering warning. - Regenerated Swagger API documentation and confirmed all tests compile and pass with 0 linter issues.
This commit is contained in:
@@ -17,7 +17,6 @@ import ("bytes"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/user"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/service"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"github.com/Rain-kl/Wavelet/internal/util"
|
||||
@@ -92,7 +91,7 @@ func TestListTaskTypes(t *testing.T) {
|
||||
foundCleanup := false
|
||||
foundWarmImageCache := false
|
||||
for _, m := range taskMetas {
|
||||
if m.Type == service.TaskTypeSystemCleanup {
|
||||
if m.Type == upload.TaskTypeSystemCleanup {
|
||||
foundCleanup = true
|
||||
}
|
||||
if m.Type == upload.TaskTypeWarmImageCache {
|
||||
@@ -100,7 +99,7 @@ func TestListTaskTypes(t *testing.T) {
|
||||
}
|
||||
}
|
||||
if !foundCleanup {
|
||||
t.Errorf("expected task type %s to be listed", service.TaskTypeSystemCleanup)
|
||||
t.Errorf("expected task type %s to be listed", upload.TaskTypeSystemCleanup)
|
||||
}
|
||||
if !foundWarmImageCache {
|
||||
t.Errorf("expected task type %s to be listed", upload.TaskTypeWarmImageCache)
|
||||
@@ -116,7 +115,7 @@ func TestDispatchTask(t *testing.T) {
|
||||
|
||||
t.Run("dispatch valid task successfully", func(t *testing.T) {
|
||||
payload := DispatchTaskRequest{
|
||||
TaskType: service.TaskTypeSystemCleanup,
|
||||
TaskType: upload.TaskTypeSystemCleanup,
|
||||
}
|
||||
body, _ := json.Marshal(payload)
|
||||
req, _ := http.NewRequest("POST", "/api/v1/admin/tasks/dispatch", bytes.NewBuffer(body))
|
||||
|
||||
@@ -0,0 +1,282 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package cap provides CAPTCHA and proof-of-work (PoW) verification services.
|
||||
package cap
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
pkgcap "github.com/Rain-kl/Wavelet/pkg/cap"
|
||||
)
|
||||
|
||||
const (
|
||||
managerDefaultChallengeCount = 1
|
||||
managerDefaultChallengeSize = 32
|
||||
defaultChallengeDifficulty = 4
|
||||
defaultChallengeTTL = 10 * time.Minute
|
||||
defaultTokenTTL = 20 * time.Minute
|
||||
redeemTokenIDLength = 8 // 兑换 Token ID 字节长度
|
||||
redeemVerTokenLength = 15 // 兑换验证 Token 字节长度
|
||||
tokenPartsCount = 2 // 兑换 Token 由两部分组成
|
||||
valuePartsCount = 2 // 存储值由 scope 和过期时间组成
|
||||
)
|
||||
|
||||
// Config holds settings for the CAPTCHA manager
|
||||
type Config struct {
|
||||
Secret []byte // HMAC signing key
|
||||
ChallengeCount int // Number of PoW puzzles
|
||||
ChallengeSize int // Size of the salt string
|
||||
ChallengeDifficulty int // Length of difficulty target prefix
|
||||
ChallengeTTL time.Duration // Lifespan of the challenge JWT
|
||||
TokenTTL time.Duration // Lifespan of the redeem token
|
||||
}
|
||||
|
||||
// Manager orchestrates challenge generation and solution validation
|
||||
type Manager struct {
|
||||
conf Config
|
||||
store pkgcap.Store
|
||||
}
|
||||
|
||||
// NewManager creates a new CAPTCHA Manager
|
||||
func NewManager(conf Config, store pkgcap.Store) *Manager {
|
||||
if conf.ChallengeCount <= 0 {
|
||||
conf.ChallengeCount = managerDefaultChallengeCount
|
||||
}
|
||||
if conf.ChallengeSize <= 0 {
|
||||
conf.ChallengeSize = managerDefaultChallengeSize
|
||||
}
|
||||
if conf.ChallengeDifficulty <= 0 {
|
||||
conf.ChallengeDifficulty = defaultChallengeDifficulty
|
||||
}
|
||||
if conf.ChallengeTTL <= 0 {
|
||||
conf.ChallengeTTL = defaultChallengeTTL
|
||||
}
|
||||
if conf.TokenTTL <= 0 {
|
||||
conf.TokenTTL = defaultTokenTTL
|
||||
}
|
||||
return &Manager{
|
||||
conf: conf,
|
||||
store: store,
|
||||
}
|
||||
}
|
||||
|
||||
// Generate creates a challenge response
|
||||
func (m *Manager) Generate(ctx context.Context, scope string) (*pkgcap.ChallengeResponse, error) {
|
||||
c := pkgcap.ChallengeConfig{
|
||||
Count: m.getChallengeCount(ctx),
|
||||
Size: m.getChallengeSize(ctx),
|
||||
Difficulty: m.getChallengeDifficulty(ctx),
|
||||
Expires: m.getChallengeTTL(ctx),
|
||||
}
|
||||
return pkgcap.GenerateChallenge(m.conf.Secret, c, scope)
|
||||
}
|
||||
|
||||
// RedeemResponse is returned to the client on redeem
|
||||
type RedeemResponse struct {
|
||||
Success bool `json:"success"`
|
||||
Token string `json:"token,omitempty"`
|
||||
Expires int64 `json:"expires,omitempty"`
|
||||
Error string `json:"error,omitempty"`
|
||||
}
|
||||
|
||||
// Redeem verifies PoW solutions and returns a one-time redeem token
|
||||
func (m *Manager) Redeem(ctx context.Context, token string, solutions []int, scope string) (*RedeemResponse, error) {
|
||||
sigHex := pkgcap.JwtSigHex(token)
|
||||
if sigHex == "" {
|
||||
return &RedeemResponse{Success: false, Error: "invalid_token"}, nil
|
||||
}
|
||||
|
||||
nonceKey := "cap:nonce:" + sigHex
|
||||
|
||||
// Atomically claim the nonce slot BEFORE verifying solutions.
|
||||
payload, err := pkgcap.VerifyChallengeSolutions(token, solutions, m.conf.Secret, scope)
|
||||
if err != nil {
|
||||
return &RedeemResponse{Success: false, Error: err.Error()}, nil //nolint:nilerr // expected behavior: validation error is returned as response, not system error
|
||||
}
|
||||
|
||||
// Calculate remaining lifetime of the challenge JWT for the nonce TTL.
|
||||
now := time.Now().UnixNano() / int64(time.Millisecond)
|
||||
nonceTTL := time.Duration(payload.Expires-now) * time.Millisecond
|
||||
if nonceTTL < time.Second {
|
||||
nonceTTL = time.Second
|
||||
}
|
||||
|
||||
// Atomic claim: if another goroutine already redeemed this JWT the SetNX
|
||||
// will return false and we reject the request without issuing a token.
|
||||
set, err := m.store.SetNX(ctx, nonceKey, "1", nonceTTL)
|
||||
if err != nil {
|
||||
return &RedeemResponse{Success: false, Error: "nonce_store_error"}, err
|
||||
}
|
||||
if !set {
|
||||
return &RedeemResponse{Success: false, Error: "already_redeemed"}, nil
|
||||
}
|
||||
|
||||
// Generate a redeem token formatted as "id:verToken"
|
||||
id := pkgcap.RandomHex(redeemTokenIDLength)
|
||||
verToken := pkgcap.RandomHex(redeemVerTokenLength)
|
||||
verHashBytes := sha256.Sum256([]byte(verToken))
|
||||
verHashHex := hex.EncodeToString(verHashBytes[:])
|
||||
|
||||
tokenKey := "cap:token:" + id + ":" + verHashHex
|
||||
tokenTTL := m.getTokenTTL(ctx)
|
||||
tokenExpires := time.Now().Add(tokenTTL)
|
||||
|
||||
// Value stored is "expiresNano|scope"
|
||||
storeVal := strconv.FormatInt(tokenExpires.UnixNano(), 10) + "|" + scope
|
||||
|
||||
if err := m.store.Set(ctx, tokenKey, storeVal, tokenTTL); err != nil {
|
||||
return &RedeemResponse{Success: false, Error: "token_store_error"}, err
|
||||
}
|
||||
|
||||
return &RedeemResponse{
|
||||
Success: true,
|
||||
Token: id + ":" + verToken,
|
||||
Expires: tokenExpires.UnixNano() / int64(time.Millisecond),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// VerifyToken validates and consumes the redeem token (single-use).
|
||||
func (m *Manager) VerifyToken(ctx context.Context, token string, expectedScope string) (bool, error) {
|
||||
if token == "" {
|
||||
return false, nil
|
||||
}
|
||||
parts := strings.Split(token, ":")
|
||||
if len(parts) != tokenPartsCount {
|
||||
return false, nil
|
||||
}
|
||||
id := parts[0]
|
||||
verToken := parts[1]
|
||||
|
||||
verHashBytes := sha256.Sum256([]byte(verToken))
|
||||
verHashHex := hex.EncodeToString(verHashBytes[:])
|
||||
|
||||
tokenKey := "cap:token:" + id + ":" + verHashHex
|
||||
|
||||
// Atomically retrieve-and-delete
|
||||
val, exists, err := sGetAndDelete(ctx, m.store, tokenKey)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
if !exists {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
valParts := strings.Split(val, "|")
|
||||
if len(valParts) != valuePartsCount {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
expNano, err := strconv.ParseInt(valParts[0], 10, 64)
|
||||
if err != nil {
|
||||
return false, nil //nolint:nilerr // expected behavior: invalid format is treated as validation failure, not system error
|
||||
}
|
||||
tokenScope := valParts[1]
|
||||
|
||||
if expectedScope != "" && tokenScope != expectedScope {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
if time.Now().UnixNano() > expNano {
|
||||
return false, nil // Expired
|
||||
}
|
||||
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// sGetAndDelete safely calls store.GetAndDelete, treating a nil store as a miss.
|
||||
func sGetAndDelete(ctx context.Context, store pkgcap.Store, key string) (string, bool, error) {
|
||||
if store == nil {
|
||||
return "", false, nil
|
||||
}
|
||||
return store.GetAndDelete(ctx, key)
|
||||
}
|
||||
|
||||
func (m *Manager) getChallengeCount(ctx context.Context) int {
|
||||
val, err := model.GetIntByKey(ctx, model.ConfigKeyCapChallengeCount)
|
||||
if err != nil || val <= 0 {
|
||||
return m.conf.ChallengeCount
|
||||
}
|
||||
return val
|
||||
}
|
||||
|
||||
func (m *Manager) getChallengeSize(ctx context.Context) int {
|
||||
val, err := model.GetIntByKey(ctx, model.ConfigKeyCapChallengeSize)
|
||||
if err != nil || val <= 0 {
|
||||
return m.conf.ChallengeSize
|
||||
}
|
||||
return val
|
||||
}
|
||||
|
||||
func (m *Manager) getChallengeDifficulty(ctx context.Context) int {
|
||||
val, err := model.GetIntByKey(ctx, model.ConfigKeyCapChallengeDifficulty)
|
||||
if err != nil || val <= 0 {
|
||||
return m.conf.ChallengeDifficulty
|
||||
}
|
||||
return val
|
||||
}
|
||||
|
||||
func (m *Manager) getChallengeTTL(ctx context.Context) time.Duration {
|
||||
val, err := model.GetIntByKey(ctx, model.ConfigKeyCapChallengeTTL)
|
||||
if err != nil || val <= 0 {
|
||||
return m.conf.ChallengeTTL
|
||||
}
|
||||
return time.Duration(val) * time.Second
|
||||
}
|
||||
|
||||
func (m *Manager) getTokenTTL(ctx context.Context) time.Duration {
|
||||
val, err := model.GetIntByKey(ctx, model.ConfigKeyCapTokenTTL)
|
||||
if err != nil || val <= 0 {
|
||||
return m.conf.TokenTTL
|
||||
}
|
||||
return time.Duration(val) * time.Second
|
||||
}
|
||||
|
||||
var (
|
||||
defaultManager *Manager
|
||||
once sync.Once
|
||||
)
|
||||
|
||||
// GetDefaultManager yields the global singleton CAPTCHA manager
|
||||
func GetDefaultManager() *Manager {
|
||||
once.Do(func() {
|
||||
var secret []byte
|
||||
if config.Config != nil && config.Config.App.SessionSecret != "" {
|
||||
secret = []byte(config.Config.App.SessionSecret)
|
||||
} else {
|
||||
secret = []byte("default-captcha-secret-key-at-least-16-bytes")
|
||||
}
|
||||
|
||||
challengeCount := managerDefaultChallengeCount
|
||||
challengeSize := managerDefaultChallengeSize
|
||||
challengeDifficulty := defaultChallengeDifficulty
|
||||
challengeTTL := defaultChallengeTTL
|
||||
tokenTTL := defaultTokenTTL
|
||||
|
||||
var store pkgcap.Store
|
||||
if config.Config != nil && config.Config.Redis.Enabled && db.Redis != nil {
|
||||
store = pkgcap.NewRedisStore(db.Redis)
|
||||
} else {
|
||||
store = pkgcap.NewMemoryStore(1 * time.Minute)
|
||||
}
|
||||
|
||||
defaultManager = NewManager(Config{
|
||||
Secret: secret,
|
||||
ChallengeCount: challengeCount,
|
||||
ChallengeSize: challengeSize,
|
||||
ChallengeDifficulty: challengeDifficulty,
|
||||
ChallengeTTL: challengeTTL,
|
||||
TokenTTL: tokenTTL,
|
||||
}, store)
|
||||
})
|
||||
return defaultManager
|
||||
}
|
||||
@@ -0,0 +1,172 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package cap
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
pkgcap "github.com/Rain-kl/Wavelet/pkg/cap"
|
||||
)
|
||||
|
||||
func TestCapFullFlow(t *testing.T) {
|
||||
secret := []byte("a-very-long-secret-key-at-least-16-bytes")
|
||||
store := pkgcap.NewMemoryStore(1 * time.Minute)
|
||||
|
||||
manager := NewManager(Config{
|
||||
Secret: secret,
|
||||
ChallengeCount: 3, // small count for fast test
|
||||
ChallengeSize: 32,
|
||||
ChallengeDifficulty: 3, // small difficulty for fast test
|
||||
ChallengeTTL: 5 * time.Second,
|
||||
TokenTTL: 10 * time.Second,
|
||||
}, store)
|
||||
|
||||
scope := "test-scope"
|
||||
ctx := context.Background()
|
||||
resp, err := manager.Generate(ctx, scope)
|
||||
if err != nil {
|
||||
t.Fatalf("Generate failed: %v", err)
|
||||
}
|
||||
|
||||
if resp.Challenge.C != 3 {
|
||||
t.Errorf("Expected count 3, got %d", resp.Challenge.C)
|
||||
}
|
||||
|
||||
// Solve the challenge (acting as client)
|
||||
solutions := pkgcap.Solve(resp.Token, resp.Challenge.C, resp.Challenge.S, resp.Challenge.D)
|
||||
|
||||
// Redeem
|
||||
redeemResp, err := manager.Redeem(ctx, resp.Token, solutions, scope)
|
||||
if err != nil {
|
||||
t.Fatalf("Redeem failed: %v", err)
|
||||
}
|
||||
if !redeemResp.Success {
|
||||
t.Fatalf("Redeem returned success=false: %s", redeemResp.Error)
|
||||
}
|
||||
if redeemResp.Token == "" {
|
||||
t.Fatalf("Expected token, got empty")
|
||||
}
|
||||
|
||||
// Verify the token
|
||||
valid, err := manager.VerifyToken(ctx, redeemResp.Token, scope)
|
||||
if err != nil {
|
||||
t.Fatalf("VerifyToken failed: %v", err)
|
||||
}
|
||||
if !valid {
|
||||
t.Fatalf("Expected redeem token to be valid")
|
||||
}
|
||||
|
||||
// Verify token is one-time use
|
||||
validAgain, err := manager.VerifyToken(ctx, redeemResp.Token, scope)
|
||||
if err != nil {
|
||||
t.Fatalf("VerifyToken second call failed: %v", err)
|
||||
}
|
||||
if validAgain {
|
||||
t.Fatalf("Expected redeem token to be single-use (invalidated after verification)")
|
||||
}
|
||||
}
|
||||
|
||||
// TestRedeemConcurrentRace verifies that when N goroutines simultaneously call
|
||||
// Redeem with the same challenge JWT, exactly one succeeds and the rest are
|
||||
// rejected with "already_redeemed". This guards against the TOCTOU fix.
|
||||
func TestRedeemConcurrentRace(t *testing.T) {
|
||||
const goroutines = 50
|
||||
|
||||
secret := []byte("race-test-secret-key-at-least-16-bytes")
|
||||
store := pkgcap.NewMemoryStore(1 * time.Minute)
|
||||
manager := NewManager(Config{
|
||||
Secret: secret,
|
||||
ChallengeCount: 1,
|
||||
ChallengeSize: 32,
|
||||
ChallengeDifficulty: 3,
|
||||
ChallengeTTL: 30 * time.Second,
|
||||
TokenTTL: 30 * time.Second,
|
||||
}, store)
|
||||
|
||||
ctx := context.Background()
|
||||
resp, err := manager.Generate(ctx, "login")
|
||||
if err != nil {
|
||||
t.Fatalf("Generate failed: %v", err)
|
||||
}
|
||||
solutions := pkgcap.Solve(resp.Token, resp.Challenge.C, resp.Challenge.S, resp.Challenge.D)
|
||||
|
||||
var (
|
||||
wg sync.WaitGroup
|
||||
success atomic.Int32
|
||||
barrier = make(chan struct{}) // synchronise goroutine start
|
||||
)
|
||||
|
||||
for i := 0; i < goroutines; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-barrier // wait for the gun
|
||||
r, _ := manager.Redeem(ctx, resp.Token, solutions, "login")
|
||||
if r != nil && r.Success {
|
||||
success.Add(1)
|
||||
}
|
||||
}()
|
||||
}
|
||||
close(barrier) // fire all goroutines at once
|
||||
wg.Wait()
|
||||
|
||||
if n := success.Load(); n != 1 {
|
||||
t.Fatalf("Expected exactly 1 successful Redeem, got %d", n)
|
||||
}
|
||||
}
|
||||
|
||||
// TestVerifyTokenConcurrentRace verifies that when N goroutines simultaneously
|
||||
// call VerifyToken with the same cap token, exactly one succeeds and the rest
|
||||
// fail. This guards against the GetAndDelete fix.
|
||||
func TestVerifyTokenConcurrentRace(t *testing.T) {
|
||||
const goroutines = 50
|
||||
|
||||
secret := []byte("race-test-secret-key-at-least-16-bytes")
|
||||
store := pkgcap.NewMemoryStore(1 * time.Minute)
|
||||
manager := NewManager(Config{
|
||||
Secret: secret,
|
||||
ChallengeCount: 1,
|
||||
ChallengeSize: 32,
|
||||
ChallengeDifficulty: 3,
|
||||
ChallengeTTL: 30 * time.Second,
|
||||
TokenTTL: 30 * time.Second,
|
||||
}, store)
|
||||
|
||||
ctx := context.Background()
|
||||
resp, _ := manager.Generate(ctx, "login")
|
||||
solutions := pkgcap.Solve(resp.Token, resp.Challenge.C, resp.Challenge.S, resp.Challenge.D)
|
||||
redeemResp, err := manager.Redeem(ctx, resp.Token, solutions, "login")
|
||||
if err != nil || !redeemResp.Success {
|
||||
t.Fatalf("Redeem failed: %v %+v", err, redeemResp)
|
||||
}
|
||||
capToken := redeemResp.Token
|
||||
|
||||
var (
|
||||
wg sync.WaitGroup
|
||||
success atomic.Int32
|
||||
barrier = make(chan struct{})
|
||||
)
|
||||
|
||||
for i := 0; i < goroutines; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-barrier
|
||||
ok, _ := manager.VerifyToken(ctx, capToken, "login")
|
||||
if ok {
|
||||
success.Add(1)
|
||||
}
|
||||
}()
|
||||
}
|
||||
close(barrier)
|
||||
wg.Wait()
|
||||
|
||||
if n := success.Load(); n != 1 {
|
||||
t.Fatalf("Expected exactly 1 successful VerifyToken, got %d", n)
|
||||
}
|
||||
}
|
||||
@@ -3,16 +3,17 @@
|
||||
|
||||
package cap
|
||||
|
||||
import ("net/http"
|
||||
import (
|
||||
"net/http"
|
||||
|
||||
caputil "github.com/Rain-kl/Wavelet/internal/service/cap"
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response")
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||
)
|
||||
|
||||
// VerifyMiddleware returns a Gin middleware that checks and consumes the X-Cap-Token header.
|
||||
// enabledFunc is an optional callback allowing dynamic check of whether captcha protection is turned on.
|
||||
func VerifyMiddleware(mgr *caputil.Manager, scope string, enabledFunc func() bool) gin.HandlerFunc {
|
||||
func VerifyMiddleware(mgr *Manager, scope string, enabledFunc func() bool) gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
if enabledFunc != nil && !enabledFunc() {
|
||||
c.Next()
|
||||
|
||||
@@ -6,8 +6,7 @@ package cap
|
||||
import (
|
||||
"net/http"
|
||||
|
||||
capService "github.com/Rain-kl/Wavelet/internal/service/cap"
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
type challengeRequest struct {
|
||||
@@ -28,7 +27,7 @@ type redeemRequest struct {
|
||||
// @Produce json
|
||||
// @Param request body challengeRequest false "可选范围限制参数"
|
||||
// @Success 200 {object} cap.ChallengeResponse "成功返回 PoW 难题"
|
||||
// @Failure 500 {object} capService.RedeemResponse "内部服务错误"
|
||||
// @Failure 500 {object} RedeemResponse "内部服务错误"
|
||||
// @Router /api/cap/challenge [post]
|
||||
func Challenge(c *gin.Context) {
|
||||
var req challengeRequest
|
||||
@@ -38,10 +37,10 @@ func Challenge(c *gin.Context) {
|
||||
req.Scope = "login"
|
||||
}
|
||||
|
||||
mgr := capService.GetDefaultManager()
|
||||
mgr := GetDefaultManager()
|
||||
resp, err := mgr.Generate(c.Request.Context(), req.Scope)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, capService.RedeemResponse{
|
||||
c.JSON(http.StatusInternalServerError, RedeemResponse{
|
||||
Success: false,
|
||||
Error: err.Error(),
|
||||
})
|
||||
@@ -58,14 +57,14 @@ func Challenge(c *gin.Context) {
|
||||
// @Accept json
|
||||
// @Produce json
|
||||
// @Param request body redeemRequest true "难题 Token 与解答 solutions 数组"
|
||||
// @Success 200 {object} capService.RedeemResponse "核销成功,返回 X-Cap-Token"
|
||||
// @Failure 400 {object} capService.RedeemResponse "参数错误或核销失败"
|
||||
// @Failure 500 {object} capService.RedeemResponse "内部服务错误"
|
||||
// @Success 200 {object} RedeemResponse "核销成功,返回 X-Cap-Token"
|
||||
// @Failure 400 {object} RedeemResponse "参数错误或核销失败"
|
||||
// @Failure 500 {object} RedeemResponse "内部服务错误"
|
||||
// @Router /api/cap/redeem [post]
|
||||
func Redeem(c *gin.Context) {
|
||||
var req redeemRequest
|
||||
if err := c.ShouldBindJSON(&req); err != nil {
|
||||
c.JSON(http.StatusBadRequest, capService.RedeemResponse{
|
||||
c.JSON(http.StatusBadRequest, RedeemResponse{
|
||||
Success: false,
|
||||
Error: "无效的参数",
|
||||
})
|
||||
@@ -76,10 +75,10 @@ func Redeem(c *gin.Context) {
|
||||
req.Scope = "login"
|
||||
}
|
||||
|
||||
mgr := capService.GetDefaultManager()
|
||||
mgr := GetDefaultManager()
|
||||
resp, err := mgr.Redeem(c.Request.Context(), req.Token, req.Solutions, req.Scope)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, capService.RedeemResponse{
|
||||
c.JSON(http.StatusInternalServerError, RedeemResponse{
|
||||
Success: false,
|
||||
Error: err.Error(),
|
||||
})
|
||||
|
||||
@@ -13,8 +13,7 @@ import ("bytes"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
capUtil "github.com/Rain-kl/Wavelet/internal/service/cap"
|
||||
pkgcap "github.com/Rain-kl/Wavelet/pkg/cap"
|
||||
pkgcap "github.com/Rain-kl/Wavelet/pkg/cap"
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response")
|
||||
@@ -34,7 +33,7 @@ func TestCapEndpointsAndMiddleware(t *testing.T) {
|
||||
}
|
||||
|
||||
// Login endpoint with CAPTCHA middleware
|
||||
r.POST("/api/v1/user/login", VerifyMiddleware(capUtil.GetDefaultManager(), "login", func() bool {
|
||||
r.POST("/api/v1/user/login", VerifyMiddleware(GetDefaultManager(), "login", func() bool {
|
||||
enabled, err := model.GetBoolByKey(context.Background(), model.ConfigKeyCapLoginEnabled)
|
||||
if err != nil {
|
||||
return false
|
||||
@@ -106,7 +105,7 @@ func TestCapEndpointsAndMiddleware(t *testing.T) {
|
||||
t.Fatalf("expected 200 OK for redeem, got %d. Body: %s", w.Code, w.Body.String())
|
||||
}
|
||||
|
||||
var redeemResp capUtil.RedeemResponse
|
||||
var redeemResp RedeemResponse
|
||||
if err := json.Unmarshal(w.Body.Bytes(), &redeemResp); err != nil {
|
||||
t.Fatalf("failed to unmarshal redeem response: %v", err)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,150 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package upload implements upload tasks and file cleanup services.
|
||||
package upload
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// 异步任务名称与管理类型定义
|
||||
const (
|
||||
// SystemCleanupTask 系统定期垃圾清理任务标识
|
||||
SystemCleanupTask = "system:cleanup"
|
||||
// TaskTypeSystemCleanup 系统定期垃圾清理管理类型
|
||||
TaskTypeSystemCleanup = "system_cleanup"
|
||||
|
||||
// 错误描述常量
|
||||
errStorageReadOnly = "存储迁移维护中,当前仅允许读取文件"
|
||||
errQueryUnusedUploadsFailed = "查询未使用的上传文件失败: %w"
|
||||
)
|
||||
|
||||
// SystemCleanupMeta represents the task metadata.
|
||||
var SystemCleanupMeta = task.TaskMeta{
|
||||
Type: TaskTypeSystemCleanup,
|
||||
AsynqTask: SystemCleanupTask,
|
||||
Name: "系统垃圾清理",
|
||||
Description: "定期清理超过1小时的未使用上传文件和超过7天的历史推送记录",
|
||||
SupportsTime: false,
|
||||
MaxRetry: task.DefaultMaxRetry,
|
||||
Queue: task.QueueDefault,
|
||||
Retryable: true,
|
||||
}
|
||||
|
||||
// SystemCleanupHandler 系统定期垃圾清理异步任务处理器
|
||||
type SystemCleanupHandler struct{}
|
||||
|
||||
// Execute 执行系统清理(包含文件清理和历史消息推送日志清理)
|
||||
func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.TaskResult, error) {
|
||||
if storageReadOnly(ctx) {
|
||||
return nil, errors.New(errStorageReadOnly)
|
||||
}
|
||||
const batchSize = 100 // 每批处理100个文件
|
||||
var lastID uint64
|
||||
var totalProcessed int
|
||||
var totalDeleted int
|
||||
|
||||
// 计算1小时前的时间
|
||||
oneHourAgo := time.Now().Add(-1 * time.Hour)
|
||||
|
||||
task.AppendLog(ctx, "开始扫描未使用上传文件,阈值: %s", oneHourAgo.Format(time.RFC3339))
|
||||
|
||||
for {
|
||||
// 使用游标分页查询未使用且超过1小时的上传记录
|
||||
var unusedUploads []model.Upload
|
||||
if err := db.DB(ctx).
|
||||
Where("id > ? AND status = ? AND created_at < ?", lastID, model.UploadStatusPending, oneHourAgo).
|
||||
Order("id ASC").
|
||||
Limit(batchSize).
|
||||
Find(&unusedUploads).Error; err != nil {
|
||||
task.AppendLog(ctx, "查询未使用的上传文件失败: %v", err)
|
||||
return nil, fmt.Errorf(errQueryUnusedUploadsFailed, err)
|
||||
}
|
||||
|
||||
// 没有更多数据,退出循环
|
||||
if len(unusedUploads) == 0 {
|
||||
break
|
||||
}
|
||||
|
||||
task.AppendLog(ctx, "本批次找到 %d 个需要清理的上传文件", len(unusedUploads))
|
||||
|
||||
// 处理每个未使用的上传文件
|
||||
for _, u := range unusedUploads {
|
||||
totalProcessed++
|
||||
|
||||
if err := db.DB(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
// 更新上传记录状态
|
||||
if err := tx.Model(&model.Upload{}).
|
||||
Where("id = ? AND status = ?", u.ID, model.UploadStatusPending).
|
||||
Update("status", model.UploadStatusDeleted).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
driver := storage.Driver(u.StorageDriver)
|
||||
if driver == "" {
|
||||
driver = storage.DriverLocal
|
||||
}
|
||||
backend, err := storage.ForDriver(ctx, driver)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := backend.Delete(ctx, u.FilePath); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}); err != nil {
|
||||
task.AppendLog(ctx, "清理上传文件失败 [ID:%d]: %v", u.ID, err)
|
||||
lastID = u.ID
|
||||
continue
|
||||
}
|
||||
|
||||
totalDeleted++
|
||||
lastID = u.ID
|
||||
}
|
||||
}
|
||||
|
||||
// 2. 清理超过7天的历史推送日志
|
||||
task.AppendLog(ctx, "开始清理历史推送审计日志,只保留最近7天数据...")
|
||||
cutoff := time.Now().AddDate(0, 0, -7)
|
||||
var pushHistoryCount int64
|
||||
if err := db.DB(ctx).Model(&model.PushHistory{}).Where("created_at < ?", cutoff).Count(&pushHistoryCount).Error; err != nil {
|
||||
task.AppendLog(ctx, "统计待清理的历史推送记录失败: %v", err)
|
||||
} else if pushHistoryCount > 0 {
|
||||
if err := db.DB(ctx).Where("created_at < ?", cutoff).Delete(&model.PushHistory{}).Error; err != nil {
|
||||
task.AppendLog(ctx, "删除历史推送记录失败: %v", err)
|
||||
} else {
|
||||
task.AppendLog(ctx, "成功删除 %d 条历史推送记录 (截止时间: %s)", pushHistoryCount, cutoff.Format("2006-01-02 15:04:05"))
|
||||
}
|
||||
} else {
|
||||
task.AppendLog(ctx, "没有需要清理的历史推送记录 (截止时间: %s)", cutoff.Format("2006-01-02 15:04:05"))
|
||||
}
|
||||
|
||||
msg := fmt.Sprintf("系统清理完成。成功清理未使用的上传文件 %d/%d 个;清理历史推送审计日志 %d 条。", totalDeleted, totalProcessed, pushHistoryCount)
|
||||
task.AppendLog(ctx, "%s", msg)
|
||||
return &task.TaskResult{Message: msg}, nil
|
||||
}
|
||||
|
||||
func storageReadOnly(ctx context.Context) bool {
|
||||
var execution model.TaskExecution
|
||||
err := db.DB(ctx).Where("task_type = ?", "storage:migrate").Order("id DESC").First(&execution).Error
|
||||
if err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return false
|
||||
}
|
||||
logger.ErrorF(ctx, "读取存储维护状态失败: %v", err)
|
||||
return true
|
||||
}
|
||||
return execution.Status != model.TaskExecutionStatusSucceeded
|
||||
}
|
||||
@@ -20,7 +20,6 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/diskcache"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/service"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
@@ -111,7 +110,7 @@ func TestSystemCleanupHandler_Execute(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
|
||||
// 执行 handler
|
||||
handler := &service.SystemCleanupHandler{}
|
||||
handler := &SystemCleanupHandler{}
|
||||
result, err := handler.Execute(ctx, nil)
|
||||
|
||||
// 验证结果
|
||||
@@ -162,7 +161,7 @@ func TestSystemCleanupHandler_ExecuteNoFiles(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
|
||||
// 没有任何上传记录
|
||||
handler := &service.SystemCleanupHandler{}
|
||||
handler := &SystemCleanupHandler{}
|
||||
result, err := handler.Execute(ctx, nil)
|
||||
|
||||
require.NoError(t, err)
|
||||
@@ -172,7 +171,7 @@ func TestSystemCleanupHandler_ExecuteNoFiles(t *testing.T) {
|
||||
|
||||
func TestSystemCleanupHandler_ImplementsTaskHandler(t *testing.T) {
|
||||
// 编译期验证 SystemCleanupHandler 实现了 TaskHandler 接口
|
||||
var _ task.TaskHandler = (*service.SystemCleanupHandler)(nil)
|
||||
var _ task.TaskHandler = (*SystemCleanupHandler)(nil)
|
||||
}
|
||||
|
||||
func TestWarmImageCacheHandlerValidatePayload(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user