mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 05:46:36 +08:00
全仓 go test -race 扫描(93 包)→ 全绿。修复 6 类数据竞争:frpc/frps 测试的锁外读与并发 Wait;oauth/repository 4 个 Pub/Sub 监听器 goroutine 读可变包变量(局部捕获 + done 通道等待);oauth 测试换 db.Redis 前停监听器;【真实生产 bug】tls 响应快照与异步续签 goroutine 并发写 cert 竞争(先快照再起 goroutine);upload/cache 监听器 goroutine 内读 db.Redis(调用方捕获)。
Result: {"status":"keep","total_issues":8,"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_exhaustive":0,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_usestdlibvars":0,"golint_wastedassign":0,"golint_total":8,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_usetesting":0,"golint_test_total":0,"golint_vetx_total":0,"eslint_problems":0,"eslint_errors":0,"eslint_warnings":0,"tsc_errors":0,"vitest_failed":0,"vitest_total":116,"measure_s":74}
This commit is contained in:
@@ -61,8 +61,12 @@ func assertStatusEventually(t *testing.T, m *Manager, relayID string, expectedSt
|
||||
for time.Now().Before(deadline) {
|
||||
m.mu.RLock()
|
||||
proc, ok := m.processes[relayID]
|
||||
var status string
|
||||
if ok {
|
||||
status = proc.Status
|
||||
}
|
||||
m.mu.RUnlock()
|
||||
if ok && proc.Status == expectedStatus {
|
||||
if ok && status == expectedStatus {
|
||||
return
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
|
||||
@@ -32,10 +32,12 @@ var (
|
||||
tokenListenerOnce sync.Once
|
||||
tokenListenerCtx context.Context
|
||||
tokenListenerCancel context.CancelFunc
|
||||
tokenListenerDone chan struct{}
|
||||
|
||||
userListenerOnce sync.Once
|
||||
userListenerCtx context.Context
|
||||
userListenerCancel context.CancelFunc
|
||||
userListenerDone chan struct{}
|
||||
)
|
||||
|
||||
func tokenCacheKey(tokenHash string) string {
|
||||
@@ -55,15 +57,19 @@ func ensureTokenCacheListener() {
|
||||
|
||||
func startTokenCacheInvalidationListener() {
|
||||
tokenListenerCtx, tokenListenerCancel = context.WithCancel(context.Background())
|
||||
tokenListenerDone = make(chan struct{})
|
||||
|
||||
go func() {
|
||||
pubsub := db.Redis.Subscribe(tokenListenerCtx, oauthTokenInvalidationChannel)
|
||||
listenerCtx := tokenListenerCtx
|
||||
defer close(tokenListenerDone)
|
||||
|
||||
pubsub := db.Redis.Subscribe(listenerCtx, oauthTokenInvalidationChannel)
|
||||
defer func() {
|
||||
_ = pubsub.Close()
|
||||
}()
|
||||
|
||||
go func() {
|
||||
<-tokenListenerCtx.Done()
|
||||
<-listenerCtx.Done()
|
||||
_ = pubsub.Close()
|
||||
}()
|
||||
|
||||
@@ -94,15 +100,19 @@ func ensureUserCacheListener() {
|
||||
|
||||
func startUserCacheInvalidationListener() {
|
||||
userListenerCtx, userListenerCancel = context.WithCancel(context.Background())
|
||||
userListenerDone = make(chan struct{})
|
||||
|
||||
go func() {
|
||||
pubsub := db.Redis.Subscribe(userListenerCtx, oauthUserInvalidationChannel)
|
||||
listenerCtx := userListenerCtx
|
||||
defer close(userListenerDone)
|
||||
|
||||
pubsub := db.Redis.Subscribe(listenerCtx, oauthUserInvalidationChannel)
|
||||
defer func() {
|
||||
_ = pubsub.Close()
|
||||
}()
|
||||
|
||||
go func() {
|
||||
<-userListenerCtx.Done()
|
||||
<-listenerCtx.Done()
|
||||
_ = pubsub.Close()
|
||||
}()
|
||||
|
||||
@@ -214,13 +224,21 @@ func InvalidateCachedUser(ctx context.Context, userID uint64) {
|
||||
func StopOauthCacheListener() {
|
||||
if tokenListenerCancel != nil {
|
||||
tokenListenerCancel()
|
||||
if tokenListenerDone != nil {
|
||||
<-tokenListenerDone
|
||||
}
|
||||
tokenListenerCancel = nil
|
||||
tokenListenerDone = nil
|
||||
}
|
||||
tokenListenerOnce = sync.Once{}
|
||||
|
||||
if userListenerCancel != nil {
|
||||
userListenerCancel()
|
||||
if userListenerDone != nil {
|
||||
<-userListenerDone
|
||||
}
|
||||
userListenerCancel = nil
|
||||
userListenerDone = nil
|
||||
}
|
||||
userListenerOnce = sync.Once{}
|
||||
}
|
||||
|
||||
@@ -381,6 +381,11 @@ func resetOIDCProviderCacheForTest() {
|
||||
}
|
||||
|
||||
func setupTestRouter(dbConn *gorm.DB, mockRedis *mockRedisClient, mockClient *http.Client) *gin.Engine {
|
||||
// 停止各层 Pub/Sub 监听 goroutine(可能在先前测试的 API 调用中随 sync.Once
|
||||
// 启动),否则它们在 db.Redis 被替换时仍读取旧值,产生数据竞争。
|
||||
StopOauthCacheListener()
|
||||
repository.StopAuthSourceCacheListener()
|
||||
repository.StopSystemConfigCacheListener()
|
||||
resetOIDCProviderCacheForTest()
|
||||
|
||||
r := testhelper.NewTestGinEngine(gin.Recovery())
|
||||
|
||||
@@ -50,10 +50,19 @@ func TestApplyCertificateReturnsApplying(t *testing.T) {
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
obtainDone := make(chan struct{})
|
||||
restore := SetObtainCertificateFuncForTest(func(ctx context.Context, cert *model.TLSCertificate) error {
|
||||
defer close(obtainDone)
|
||||
return updateCertError(ctx, cert, "dns challenge failed")
|
||||
})
|
||||
defer restore()
|
||||
defer func() {
|
||||
select {
|
||||
case <-obtainDone:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("async certificate obtain did not finish")
|
||||
}
|
||||
restore()
|
||||
}()
|
||||
|
||||
cert, err := ApplyCertificate(ctx, ApplyInput{
|
||||
Name: "Test ACME Cert",
|
||||
|
||||
@@ -197,12 +197,17 @@ func ApplyCertificate(ctx context.Context, input ApplyInput) (*model.TLSCertific
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 先取响应快照再启动异步续签:sanitize 会整体拷贝 cert,若与异步 goroutine
|
||||
// 的字段写入并发会构成数据竞争(生产真实问题)。
|
||||
returned := sanitizeCertificateForResponse(cert)
|
||||
|
||||
obtainFn := obtainTLSCertificate // 捕获当前实现,避免 goroutine 内读可变包变量(测试热替换)
|
||||
go func(c *model.TLSCertificate) {
|
||||
asyncCtx := context.WithoutCancel(ctx)
|
||||
_ = obtainTLSCertificate(asyncCtx, c)
|
||||
_ = obtainFn(asyncCtx, c)
|
||||
}(cert)
|
||||
|
||||
return sanitizeCertificateForResponse(cert), nil
|
||||
return returned, nil
|
||||
}
|
||||
|
||||
// UpdateACMECertificate 更新 ACME 证书配置。
|
||||
@@ -225,12 +230,15 @@ func UpdateACMECertificate(ctx context.Context, id uint, input ApplyInput) (*mod
|
||||
return nil, err
|
||||
}
|
||||
|
||||
returned := sanitizeCertificateForResponse(cert)
|
||||
|
||||
obtainFn := obtainTLSCertificate // 捕获当前实现,避免 goroutine 内读可变包变量(测试热替换)
|
||||
go func(c *model.TLSCertificate) {
|
||||
asyncCtx := context.WithoutCancel(ctx)
|
||||
_ = obtainTLSCertificate(asyncCtx, c)
|
||||
_ = obtainFn(asyncCtx, c)
|
||||
}(cert)
|
||||
|
||||
return sanitizeCertificateForResponse(cert), nil
|
||||
return returned, nil
|
||||
}
|
||||
|
||||
// ConvertCertificateToACME 将上传证书转为 ACME 管理。
|
||||
@@ -257,9 +265,10 @@ func ConvertCertificateToACME(ctx context.Context, id uint, input ApplyInput) (*
|
||||
return nil, err
|
||||
}
|
||||
|
||||
obtainFn := obtainTLSCertificate // 捕获当前实现,避免 goroutine 内读可变包变量(测试热替换)
|
||||
go func(c *model.TLSCertificate) {
|
||||
asyncCtx := context.WithoutCancel(ctx)
|
||||
if err := obtainTLSCertificate(asyncCtx, c); err != nil {
|
||||
if err := obtainFn(asyncCtx, c); err != nil {
|
||||
return
|
||||
}
|
||||
latest, err := repository.GetTLSCertificateByID(asyncCtx, c.ID)
|
||||
|
||||
@@ -7,7 +7,7 @@ import (
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"syscall"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -307,17 +307,17 @@ func TestSupervisorGenerationInterrupt(t *testing.T) {
|
||||
t.Error("expected old process killed and new command started")
|
||||
}
|
||||
|
||||
// Verify old process is actually killed
|
||||
var cmd1Finished int32
|
||||
go func() {
|
||||
_ = cmd1.Wait()
|
||||
atomic.StoreInt32(&cmd1Finished, 1)
|
||||
}()
|
||||
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
if atomic.LoadInt32(&cmd1Finished) != 1 {
|
||||
t.Error("expected first process to be killed")
|
||||
// Verify old process is actually killed:不要对受管 Cmd 调用 Wait(旧 supervise
|
||||
// goroutine 拥有 Wait 权,并发 Wait 会与 os/exec 内部状态竞争),改为探测
|
||||
// 进程是否已被收割(Signal(0) 在 Wait 后即报错)。
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
if err := cmd1.Process.Signal(syscall.Signal(0)); err != nil {
|
||||
return
|
||||
}
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
}
|
||||
t.Error("expected first process to be killed")
|
||||
}
|
||||
|
||||
func TestUpdateConfigKillsOrphanProcessBeforeRestart(t *testing.T) {
|
||||
|
||||
+3
-2
@@ -52,12 +52,13 @@ func ensureAccessCacheListener() {
|
||||
}
|
||||
|
||||
func startAccessCacheInvalidationListener() {
|
||||
if db.Redis == nil {
|
||||
redis := db.Redis // 调用方 goroutine 上捕获,避免 goroutine 内读可变全局(测试会替换 db.Redis)
|
||||
if redis == nil {
|
||||
return
|
||||
}
|
||||
|
||||
go func() {
|
||||
pubsub := db.Redis.Subscribe(
|
||||
pubsub := redis.Subscribe(
|
||||
context.Background(),
|
||||
objectstore.ConfigInvalidationChannel,
|
||||
fileAccessInvalidationChannel,
|
||||
|
||||
@@ -48,6 +48,7 @@ var (
|
||||
authSourceListenerOnce sync.Once
|
||||
authSourceListenerCtx context.Context
|
||||
authSourceListenerCancel context.CancelFunc
|
||||
authSourceListenerDone chan struct{}
|
||||
)
|
||||
|
||||
func cloneAuthSources(sources []model.AuthSource) []model.AuthSource {
|
||||
@@ -116,15 +117,19 @@ func ensureAuthSourceCacheListener() {
|
||||
|
||||
func startAuthSourceCacheInvalidationListener() {
|
||||
authSourceListenerCtx, authSourceListenerCancel = context.WithCancel(context.Background())
|
||||
authSourceListenerDone = make(chan struct{})
|
||||
|
||||
go func() {
|
||||
pubsub := db.Redis.Subscribe(authSourceListenerCtx, authSourceInvalidationChannel)
|
||||
listenerCtx := authSourceListenerCtx
|
||||
defer close(authSourceListenerDone)
|
||||
|
||||
pubsub := db.Redis.Subscribe(listenerCtx, authSourceInvalidationChannel)
|
||||
defer func() {
|
||||
_ = pubsub.Close()
|
||||
}()
|
||||
|
||||
go func() {
|
||||
<-authSourceListenerCtx.Done()
|
||||
<-listenerCtx.Done()
|
||||
_ = pubsub.Close()
|
||||
}()
|
||||
|
||||
@@ -257,7 +262,11 @@ func InvalidateAuthSourceCache(ctx context.Context) error {
|
||||
func StopAuthSourceCacheListener() {
|
||||
if authSourceListenerCancel != nil {
|
||||
authSourceListenerCancel()
|
||||
if authSourceListenerDone != nil {
|
||||
<-authSourceListenerDone
|
||||
}
|
||||
authSourceListenerCancel = nil
|
||||
authSourceListenerDone = nil
|
||||
}
|
||||
authSourceListenerOnce = sync.Once{}
|
||||
}
|
||||
|
||||
@@ -90,6 +90,7 @@ var (
|
||||
systemConfigListenerOnce sync.Once
|
||||
systemConfigListenerCtx context.Context
|
||||
systemConfigListenerCancel context.CancelFunc
|
||||
systemConfigListenerDone chan struct{}
|
||||
)
|
||||
|
||||
func ensureSystemConfigCacheListener() {
|
||||
@@ -102,15 +103,19 @@ func startSystemConfigCacheInvalidationListener() {
|
||||
}
|
||||
|
||||
systemConfigListenerCtx, systemConfigListenerCancel = context.WithCancel(context.Background())
|
||||
systemConfigListenerDone = make(chan struct{})
|
||||
|
||||
go func() {
|
||||
pubsub := db.Redis.Subscribe(systemConfigListenerCtx, SystemConfigBroadcastChannel)
|
||||
listenerCtx := systemConfigListenerCtx
|
||||
defer close(systemConfigListenerDone)
|
||||
|
||||
pubsub := db.Redis.Subscribe(listenerCtx, SystemConfigBroadcastChannel)
|
||||
defer func() {
|
||||
_ = pubsub.Close()
|
||||
}()
|
||||
|
||||
go func() {
|
||||
<-systemConfigListenerCtx.Done()
|
||||
<-listenerCtx.Done()
|
||||
_ = pubsub.Close()
|
||||
}()
|
||||
|
||||
@@ -135,7 +140,11 @@ func startSystemConfigCacheInvalidationListener() {
|
||||
func StopSystemConfigCacheListener() {
|
||||
if systemConfigListenerCancel != nil {
|
||||
systemConfigListenerCancel()
|
||||
if systemConfigListenerDone != nil {
|
||||
<-systemConfigListenerDone
|
||||
}
|
||||
systemConfigListenerCancel = nil
|
||||
systemConfigListenerDone = nil
|
||||
}
|
||||
systemConfigListenerOnce = sync.Once{}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user