diff --git a/backend/go.mod b/backend/go.mod index 66b6cc25b5..12ff846027 100644 --- a/backend/go.mod +++ b/backend/go.mod @@ -54,6 +54,7 @@ require ( github.com/Azure/go-ansiterm v0.0.0-20210617225240-d185dfc1b5a1 // indirect github.com/Microsoft/go-winio v0.6.2 // indirect github.com/agext/levenshtein v1.2.3 // indirect + github.com/alicebob/miniredis/v2 v2.37.0 // indirect github.com/apparentlymart/go-textseg/v15 v15.0.0 // indirect github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.5 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.18 // indirect @@ -158,6 +159,7 @@ require ( github.com/tklauser/numcpus v0.6.1 // indirect github.com/twitchyliquid64/golang-asm v0.15.1 // indirect github.com/ugorji/go/codec v1.2.11 // indirect + github.com/yuin/gopher-lua v1.1.1 // indirect github.com/yusufpapurcu/wmi v1.2.4 // indirect github.com/zclconf/go-cty v1.14.4 // indirect github.com/zclconf/go-cty-yaml v1.1.0 // indirect diff --git a/backend/go.sum b/backend/go.sum index 9312af63e5..0502b30144 100644 --- a/backend/go.sum +++ b/backend/go.sum @@ -16,6 +16,8 @@ github.com/agext/levenshtein v1.2.3 h1:YB2fHEn0UJagG8T1rrWknE3ZQzWM06O8AMAatNn7l github.com/agext/levenshtein v1.2.3/go.mod h1:JEDfjyjHDjOF/1e4FlBE/PkbqA9OfWu2ki2W0IB5558= github.com/agiledragon/gomonkey v2.0.2+incompatible h1:eXKi9/piiC3cjJD1658mEE2o3NjkJ5vDLgYjCQu0Xlw= github.com/agiledragon/gomonkey v2.0.2+incompatible/go.mod h1:2NGfXu1a80LLr2cmWXGBDaHEjb1idR6+FVlX5T3D9hw= +github.com/alicebob/miniredis/v2 v2.37.0 h1:RheObYW32G1aiJIj81XVt78ZHJpHonHLHW7OLIshq68= +github.com/alicebob/miniredis/v2 v2.37.0/go.mod h1:TcL7YfarKPGDAthEtl5NBeHZfeUQj6OXMm/+iu5cLMM= github.com/alitto/pond/v2 v2.6.2 h1:Sphe40g0ILeM1pA2c2K+Th0DGU+pt0A/Kprr+WB24Pw= github.com/alitto/pond/v2 v2.6.2/go.mod h1:xkjYEgQ05RSpWdfSd1nM3OVv7TBhLdy7rMp3+2Nq+yE= github.com/andybalholm/brotli v1.2.0 h1:ukwgCxwYrmACq68yiUqwIWnGY0cTPox/M94sVwToPjQ= @@ -372,6 +374,8 @@ github.com/wechatpay-apiv3/wechatpay-go v0.2.21 h1:uIyMpzvcaHA33W/QPtHstccw+X52H github.com/wechatpay-apiv3/wechatpay-go v0.2.21/go.mod h1:A254AUBVB6R+EqQFo3yTgeh7HtyqRRtN2w9hQSOrd4Q= github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU= github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E= +github.com/yuin/gopher-lua v1.1.1 h1:kYKnWBjvbNP4XLT3+bPEwAXJx262OhaHDWDVOPjL46M= +github.com/yuin/gopher-lua v1.1.1/go.mod h1:GBR0iDaNXjAgGg9zfCvksxSRnQx76gclCIb7kdAd1Pw= github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0= github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0= github.com/zclconf/go-cty v1.14.4 h1:uXXczd9QDGsgu0i/QFR/hzI5NYCHLf6NQw/atrbnhq8= diff --git a/backend/internal/repository/signature_pool_cache_integration_test.go b/backend/internal/repository/signature_pool_cache_integration_test.go new file mode 100644 index 0000000000..5e2afb179f --- /dev/null +++ b/backend/internal/repository/signature_pool_cache_integration_test.go @@ -0,0 +1,205 @@ +//go:build integration + +package repository + +import ( + "fmt" + "testing" + "time" + + "github.com/Wei-Shaw/sub2api/internal/service/signature" + "github.com/stretchr/testify/require" + "github.com/stretchr/testify/suite" +) + +type SignaturePoolCacheSuite struct { + IntegrationRedisSuite + pool signature.SignaturePool +} + +func TestSignaturePoolCacheSuite(t *testing.T) { + suite.Run(t, new(SignaturePoolCacheSuite)) +} + +func (s *SignaturePoolCacheSuite) SetupTest() { + s.IntegrationRedisSuite.SetupTest() + s.pool = NewSignaturePoolCache(s.rdb, signature.DefaultSignatureTTL) +} + +// uniqueBucket returns a test-isolated bucket to avoid cross-test pollution. +func (s *SignaturePoolCacheSuite) uniqueBucket(suffix string) string { + return fmt.Sprintf("test:%s:%d", suffix, time.Now().UnixNano()) +} + +// --- Basic flow --- + +func (s *SignaturePoolCacheSuite) TestAddAndTopN_BasicFlow() { + bucket := s.uniqueBucket("basic") + now := time.Now() + + s.RequireNoError(s.pool.Add(s.ctx, bucket, "sig-A", now, 10)) + s.RequireNoError(s.pool.Add(s.ctx, bucket, "sig-B", now.Add(time.Second), 10)) + s.RequireNoError(s.pool.Add(s.ctx, bucket, "sig-C", now.Add(2*time.Second), 10)) + + sigs, err := s.pool.TopN(s.ctx, bucket, 3) + s.RequireNoError(err) + // Newest first: C, B, A + require.Equal(s.T(), []string{"sig-C", "sig-B", "sig-A"}, sigs) +} + +func (s *SignaturePoolCacheSuite) TestTopN_ReturnsFewerWhenPoolSmall() { + bucket := s.uniqueBucket("fewer") + now := time.Now() + + s.RequireNoError(s.pool.Add(s.ctx, bucket, "only-one", now, 10)) + + sigs, err := s.pool.TopN(s.ctx, bucket, 5) + s.RequireNoError(err) + require.Equal(s.T(), []string{"only-one"}, sigs) +} + +func (s *SignaturePoolCacheSuite) TestTopN_EmptyBucketReturnsNil() { + sigs, err := s.pool.TopN(s.ctx, s.uniqueBucket("empty"), 10) + s.RequireNoError(err) + require.Empty(s.T(), sigs) +} + +// --- Capacity trimming --- + +func (s *SignaturePoolCacheSuite) TestAdd_TrimsToCapacity() { + bucket := s.uniqueBucket("cap") + now := time.Now() + cap := 3 + + for i := 0; i < 6; i++ { + sig := fmt.Sprintf("sig-%d", i) + s.RequireNoError(s.pool.Add(s.ctx, bucket, sig, now.Add(time.Duration(i)*time.Second), cap)) + } + + // Only the newest 3 should survive. + sigs, err := s.pool.TopN(s.ctx, bucket, 10) + s.RequireNoError(err) + require.Equal(s.T(), []string{"sig-5", "sig-4", "sig-3"}, sigs) + + sz, err := s.pool.Size(s.ctx, bucket) + s.RequireNoError(err) + require.Equal(s.T(), int64(3), sz) +} + +// --- Soft TTL lazy expiry --- + +func (s *SignaturePoolCacheSuite) TestAdd_EvictsExpiredEntries() { + bucket := s.uniqueBucket("ttl") + ttl := signature.DefaultSignatureTTL // 1h + + // Add an entry timestamped 2 hours in the past — already expired. + oldTime := time.Now().Add(-2 * ttl) + s.RequireNoError(s.pool.Add(s.ctx, bucket, "old-sig", oldTime, 100)) + + // Verify it exists before a new Add triggers cleanup. + sz, err := s.pool.Size(s.ctx, bucket) + s.RequireNoError(err) + require.Equal(s.T(), int64(1), sz, "old entry should exist before lazy cleanup") + + // Add a fresh entry — lazy cleanup should remove the expired one. + s.RequireNoError(s.pool.Add(s.ctx, bucket, "new-sig", time.Now(), 100)) + + sigs, err := s.pool.TopN(s.ctx, bucket, 10) + s.RequireNoError(err) + require.Equal(s.T(), []string{"new-sig"}, sigs, "expired entry should be evicted by lazy cleanup") + + sz, err = s.pool.Size(s.ctx, bucket) + s.RequireNoError(err) + require.Equal(s.T(), int64(1), sz) +} + +func (s *SignaturePoolCacheSuite) TestTopN_DoesNotFilterByTTL() { + // Per design: "避免没有签名可用" — expired entries survive until the next + // Add evicts them. + bucket := s.uniqueBucket("no-filter") + ttl := signature.DefaultSignatureTTL + + oldTime := time.Now().Add(-2 * ttl) + s.RequireNoError(s.pool.Add(s.ctx, bucket, "stale-but-valid", oldTime, 100)) + + sigs, err := s.pool.TopN(s.ctx, bucket, 10) + s.RequireNoError(err) + require.Equal(s.T(), []string{"stale-but-valid"}, sigs, + "TopN must NOT filter by TTL — lazy expiry only happens on Add") +} + +// --- De-duplication (ZADD score update) --- + +func (s *SignaturePoolCacheSuite) TestAdd_SameSignatureUpdatesScore() { + bucket := s.uniqueBucket("dedup") + now := time.Now() + + s.RequireNoError(s.pool.Add(s.ctx, bucket, "sig-A", now, 10)) + s.RequireNoError(s.pool.Add(s.ctx, bucket, "sig-B", now.Add(time.Second), 10)) + // Re-add sig-A with a newer timestamp — it should move to the front. + s.RequireNoError(s.pool.Add(s.ctx, bucket, "sig-A", now.Add(2*time.Second), 10)) + + sigs, err := s.pool.TopN(s.ctx, bucket, 10) + s.RequireNoError(err) + // sig-A is now newest (score updated). + require.Equal(s.T(), []string{"sig-A", "sig-B"}, sigs) + + sz, err := s.pool.Size(s.ctx, bucket) + s.RequireNoError(err) + require.Equal(s.T(), int64(2), sz, "duplicate should not increase pool size") +} + +// --- Redis key-level TTL --- + +func (s *SignaturePoolCacheSuite) TestAdd_SetsKeyTTL() { + bucket := s.uniqueBucket("keyttl") + key := signaturePoolKey(bucket) + + s.RequireNoError(s.pool.Add(s.ctx, bucket, "sig", time.Now(), 10)) + + ttl, err := s.rdb.TTL(s.ctx, key).Result() + s.RequireNoError(err) + // key TTL = softTTL * signaturePoolKeyTTLFactor = 1h * 24 = 24h. + // Allow a generous window for test latency. + s.AssertTTLWithin(ttl, 23*time.Hour, 25*time.Hour) +} + +// --- Size --- + +func (s *SignaturePoolCacheSuite) TestSize_ReflectsCurrentEntries() { + bucket := s.uniqueBucket("size") + now := time.Now() + + sz, err := s.pool.Size(s.ctx, bucket) + s.RequireNoError(err) + require.Equal(s.T(), int64(0), sz, "empty bucket") + + for i := 0; i < 5; i++ { + s.RequireNoError(s.pool.Add(s.ctx, bucket, fmt.Sprintf("s%d", i), now.Add(time.Duration(i)*time.Second), 100)) + } + sz, err = s.pool.Size(s.ctx, bucket) + s.RequireNoError(err) + require.Equal(s.T(), int64(5), sz) +} + +// --- Edge cases --- + +func (s *SignaturePoolCacheSuite) TestAdd_EmptyBucketOrSignatureIsNoop() { + require.NoError(s.T(), s.pool.Add(s.ctx, "", "sig", time.Now(), 10)) + require.NoError(s.T(), s.pool.Add(s.ctx, "b", "", time.Now(), 10)) +} + +func (s *SignaturePoolCacheSuite) TestTopN_EmptyBucketStringReturnsNil() { + sigs, err := s.pool.TopN(s.ctx, "", 10) + s.RequireNoError(err) + require.Empty(s.T(), sigs) +} + +func (s *SignaturePoolCacheSuite) TestTopN_ZeroNReturnsNil() { + bucket := s.uniqueBucket("zeron") + s.RequireNoError(s.pool.Add(s.ctx, bucket, "sig", time.Now(), 10)) + + sigs, err := s.pool.TopN(s.ctx, bucket, 0) + s.RequireNoError(err) + require.Empty(s.T(), sigs) +} diff --git a/backend/internal/repository/signature_pool_cache_test.go b/backend/internal/repository/signature_pool_cache_test.go new file mode 100644 index 0000000000..e826b625f4 --- /dev/null +++ b/backend/internal/repository/signature_pool_cache_test.go @@ -0,0 +1,147 @@ +//go:build unit + +package repository + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/Wei-Shaw/sub2api/internal/service/signature" + "github.com/alicebob/miniredis/v2" + "github.com/redis/go-redis/v9" + "github.com/stretchr/testify/require" +) + +func setupPoolTest(t *testing.T, softTTL time.Duration) (context.Context, signature.SignaturePool) { + t.Helper() + mr := miniredis.RunT(t) + rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) + t.Cleanup(func() { _ = rdb.Close() }) + return context.Background(), NewSignaturePoolCache(rdb, softTTL) +} + +func TestSignaturePool_AddAndTopN(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + now := time.Now() + + require.NoError(t, pool.Add(ctx, "b", "A", now, 10)) + require.NoError(t, pool.Add(ctx, "b", "B", now.Add(time.Second), 10)) + require.NoError(t, pool.Add(ctx, "b", "C", now.Add(2*time.Second), 10)) + + sigs, err := pool.TopN(ctx, "b", 3) + require.NoError(t, err) + require.Equal(t, []string{"C", "B", "A"}, sigs) +} + +func TestSignaturePool_TopN_FewerThanN(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + require.NoError(t, pool.Add(ctx, "b", "only", time.Now(), 10)) + + sigs, err := pool.TopN(ctx, "b", 5) + require.NoError(t, err) + require.Equal(t, []string{"only"}, sigs) +} + +func TestSignaturePool_TopN_EmptyBucket(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + sigs, err := pool.TopN(ctx, "empty", 10) + require.NoError(t, err) + require.Empty(t, sigs) +} + +func TestSignaturePool_CapacityTrim(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + now := time.Now() + + for i := 0; i < 6; i++ { + require.NoError(t, pool.Add(ctx, "b", fmt.Sprintf("s%d", i), now.Add(time.Duration(i)*time.Second), 3)) + } + + sigs, err := pool.TopN(ctx, "b", 10) + require.NoError(t, err) + require.Equal(t, []string{"s5", "s4", "s3"}, sigs) + + sz, err := pool.Size(ctx, "b") + require.NoError(t, err) + require.Equal(t, int64(3), sz) +} + +func TestSignaturePool_LazyExpiry(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + + oldTime := time.Now().Add(-2 * time.Hour) + require.NoError(t, pool.Add(ctx, "b", "old", oldTime, 100)) + + sz, _ := pool.Size(ctx, "b") + require.Equal(t, int64(1), sz, "old entry exists before cleanup") + + require.NoError(t, pool.Add(ctx, "b", "new", time.Now(), 100)) + + sigs, err := pool.TopN(ctx, "b", 10) + require.NoError(t, err) + require.Equal(t, []string{"new"}, sigs) + + sz, _ = pool.Size(ctx, "b") + require.Equal(t, int64(1), sz) +} + +func TestSignaturePool_TopN_NoTTLFilter(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + + oldTime := time.Now().Add(-2 * time.Hour) + require.NoError(t, pool.Add(ctx, "b", "stale", oldTime, 100)) + + sigs, err := pool.TopN(ctx, "b", 10) + require.NoError(t, err) + require.Equal(t, []string{"stale"}, sigs, "TopN must not filter by TTL") +} + +func TestSignaturePool_DuplicateUpdatesScore(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + now := time.Now() + + require.NoError(t, pool.Add(ctx, "b", "A", now, 10)) + require.NoError(t, pool.Add(ctx, "b", "B", now.Add(time.Second), 10)) + require.NoError(t, pool.Add(ctx, "b", "A", now.Add(2*time.Second), 10)) + + sigs, err := pool.TopN(ctx, "b", 10) + require.NoError(t, err) + require.Equal(t, []string{"A", "B"}, sigs) + + sz, _ := pool.Size(ctx, "b") + require.Equal(t, int64(2), sz) +} + +func TestSignaturePool_EmptyInputNoop(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + require.NoError(t, pool.Add(ctx, "", "sig", time.Now(), 10)) + require.NoError(t, pool.Add(ctx, "b", "", time.Now(), 10)) + + sigs, _ := pool.TopN(ctx, "", 10) + require.Empty(t, sigs) +} + +func TestSignaturePool_ZeroN(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + require.NoError(t, pool.Add(ctx, "b", "sig", time.Now(), 10)) + + sigs, err := pool.TopN(ctx, "b", 0) + require.NoError(t, err) + require.Empty(t, sigs) +} + +func TestSignaturePool_Size(t *testing.T) { + ctx, pool := setupPoolTest(t, time.Hour) + now := time.Now() + + sz, _ := pool.Size(ctx, "b") + require.Equal(t, int64(0), sz) + + for i := 0; i < 5; i++ { + require.NoError(t, pool.Add(ctx, "b", fmt.Sprintf("s%d", i), now.Add(time.Duration(i)*time.Second), 100)) + } + sz, _ = pool.Size(ctx, "b") + require.Equal(t, int64(5), sz) +}