fix: 修复应用日志异常膨胀

This commit is contained in:
ryan
2026-06-21 10:21:18 +08:00
parent a343c7a605
commit b12a9b0185
12 changed files with 358 additions and 53 deletions
+20
View File
@@ -10,6 +10,7 @@ import (
"log/slog"
"sort"
"strings"
"sync"
openrestyrender "github.com/Rain-kl/Wavelet/pkg/render/openresty"
@@ -49,6 +50,7 @@ type Service struct {
nginxManager NginxManager
stateStore *state.Store
pagesDir string
syncMu sync.Mutex
}
// SetPagesDir sets the local directory used for pages deployment packages.
@@ -76,6 +78,9 @@ func (s *Service) SyncOnStartup(ctx context.Context, target *protocol.ActiveConf
}
func (s *Service) sync(ctx context.Context, startup bool, target *protocol.ActiveConfigMeta) error {
s.syncMu.Lock()
defer s.syncMu.Unlock()
mode := syncMode(startup)
snapshot, currentChecksum, err := s.loadSyncState()
if err != nil {
@@ -94,6 +99,9 @@ func (s *Service) sync(ctx context.Context, startup bool, target *protocol.Activ
// ForceSyncOnce clears any blocked target state then unconditionally fetches and applies the active config.
func (s *Service) ForceSyncOnce(ctx context.Context, target *protocol.ActiveConfigMeta) error {
s.syncMu.Lock()
defer s.syncMu.Unlock()
snapshot, err := s.stateStore.Load()
if err != nil {
return err
@@ -161,12 +169,24 @@ func (s *Service) applyRenderedConfig(ctx context.Context, mode string, snapshot
mainConfigChecksum := checksumString(rendered.mainConfig)
routeConfigChecksum := checksumString(rendered.routeConfig)
slog.Info("applying new openresty config", "mode", mode, "from_version", snapshot.CurrentVersion, "to_version", config.Version, "old_checksum", currentChecksum, "new_checksum", config.Checksum)
alreadySynced := snapshotMatchesTarget(snapshot, config.Version, config.Checksum)
outcome, message := normalizeApplyOutcome(s.nginxManager.Apply(ctx, rendered.mainConfig, rendered.routeConfig, rendered.supportFiles))
applyResult := updateSnapshotFromApplyOutcome(mode, snapshot, config, outcome, message)
if err := s.stateStore.Save(snapshot); err != nil {
return err
}
if !shouldReportApplyLog(alreadySynced, applyResult.reportResult) {
slog.Debug("skipping duplicate apply log report", "version", config.Version, "checksum", config.Checksum, "result", applyResult.reportResult)
if applyResult.reportResult == ApplyResultFailed {
return outcomeError(config.Version, applyResult.message)
}
if err := s.syncReferencedWAFIPGroups(ctx, rendered.supportFiles); err != nil {
slog.Error("sync referenced waf ip groups failed", "version", config.Version, "error", err)
return err
}
return nil
}
if err := s.client.ReportApplyLog(ctx, protocol.ApplyLogPayload{
NodeID: snapshot.NodeID,
Version: config.Version,
+38
View File
@@ -479,6 +479,44 @@ func TestSyncOnceReportsNoopWhenVersionChangesButChecksumMatches(t *testing.T) {
}
}
func TestSyncOnStartupSkipsDuplicateSuccessReportWhenStateAlreadySynced(t *testing.T) {
client := &fakeClient{
config: protocol.ActiveConfigResponse{
Version: "20260309-003",
Checksum: "checksum-3",
SourceConfigJSON: testSourceConfigJSON("auto", 80),
CreatedAt: time.Now().Format(time.RFC3339),
},
}
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
nodeID, err := stateStore.EnsureNodeID()
if err != nil {
t.Fatalf("EnsureNodeID failed: %v", err)
}
if err = stateStore.Save(&state.Snapshot{
NodeID: nodeID,
CurrentVersion: "20260309-003",
CurrentChecksum: "checksum-3",
}); err != nil {
t.Fatalf("failed to seed state: %v", err)
}
manager := &fakeManager{currentChecksum: "checksum-3"}
service := New(client, manager, stateStore)
if err = service.SyncOnStartup(context.Background(), &protocol.ActiveConfigMeta{
Version: "20260309-003",
Checksum: "checksum-3",
}); err != nil {
t.Fatalf("SyncOnStartup failed: %v", err)
}
if len(client.reports) != 0 {
t.Fatalf("expected startup sync to skip duplicate success report, got %+v", client.reports)
}
if len(manager.applyMainContents) != 1 {
t.Fatal("expected startup sync to still refresh local config once")
}
}
func TestSyncOnceDoesNotRepeatNoopReportWhenStateAlreadyMatches(t *testing.T) {
client := &fakeClient{}
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
+15
View File
@@ -203,3 +203,18 @@ func updateSnapshotFromApplyOutcome(mode string, snapshot *state.Snapshot, confi
}
return result
}
func snapshotMatchesTarget(snapshot *state.Snapshot, version string, checksum string) bool {
if snapshot == nil {
return false
}
return strings.TrimSpace(snapshot.CurrentVersion) == strings.TrimSpace(version) &&
strings.TrimSpace(snapshot.CurrentChecksum) == strings.TrimSpace(checksum)
}
func shouldReportApplyLog(alreadySynced bool, result string) bool {
if result != ApplyResultSuccess {
return true
}
return !alreadySynced
}
+6 -5
View File
@@ -95,12 +95,13 @@ func (m *Manager) GetCurrentConfigChecksum() string {
}
// UpdateConfig reconciles running frpc processes with the latest tunnel configuration.
func (m *Manager) UpdateConfig(ctx context.Context, newConfig *service.FlaredTunnelConfigResponse) error {
// The returned bool indicates whether the active config version or checksum changed.
func (m *Manager) UpdateConfig(ctx context.Context, newConfig *service.FlaredTunnelConfigResponse) (bool, error) {
m.mu.Lock()
defer m.mu.Unlock()
if newConfig == nil {
return nil
return false, nil
}
versionChanged := newConfig.Version != m.currentVersion || newConfig.Checksum != m.currentChecksum
@@ -111,7 +112,7 @@ func (m *Manager) UpdateConfig(ctx context.Context, newConfig *service.FlaredTun
}
if err := os.MkdirAll(m.cfg.DataDir, dataDirPerm); err != nil {
return fmt.Errorf("create data dir failed: %w", err)
return false, fmt.Errorf("create data dir failed: %w", err)
}
activeRelays := make(map[string]struct{})
@@ -155,9 +156,9 @@ func (m *Manager) UpdateConfig(ctx context.Context, newConfig *service.FlaredTun
if versionChanged {
m.currentVersion = newConfig.Version
m.currentChecksum = newConfig.Checksum
return m.saveState()
return true, m.saveState()
}
return nil
return false, nil
}
func (m *Manager) restartProcess(ctx context.Context, relayID string, configPath string) {
+6 -6
View File
@@ -119,7 +119,7 @@ func TestStartProcessSuccess(t *testing.T) {
Proxies: nil,
}
err := m.UpdateConfig(context.Background(), newConfig)
_, err := m.UpdateConfig(context.Background(), newConfig)
if err != nil {
t.Fatalf("failed to UpdateConfig: %v", err)
}
@@ -160,7 +160,7 @@ func TestStartProcessFailureAndBackoff(t *testing.T) {
Proxies: nil,
}
_ = m.UpdateConfig(context.Background(), newConfig)
_, _ = m.UpdateConfig(context.Background(), newConfig)
assertStatusEventually(t, m, "relay-1", "error", 4*time.Second)
@@ -208,7 +208,7 @@ func TestUnexpectedExit0CPUProtection(t *testing.T) {
Proxies: nil,
}
_ = m.UpdateConfig(context.Background(), newConfig)
_, _ = m.UpdateConfig(context.Background(), newConfig)
assertStatusEventually(t, m, "relay-1", "stopped", 4*time.Second)
@@ -249,7 +249,7 @@ func TestBackoffReset(t *testing.T) {
Proxies: nil,
}
_ = m.UpdateConfig(context.Background(), newConfig)
_, _ = m.UpdateConfig(context.Background(), newConfig)
// Wait to crash
assertStatusEventually(t, m, "relay-1", "error", 4*time.Second)
@@ -323,7 +323,7 @@ func TestUpdateConfigKillsOrphanProcessBeforeRestart(t *testing.T) {
},
}
if err := m.UpdateConfig(context.Background(), newConfig); err != nil {
if _, err := m.UpdateConfig(context.Background(), newConfig); err != nil {
t.Fatalf("failed to UpdateConfig: %v", err)
}
@@ -361,7 +361,7 @@ func TestStopCancelsRunningProcesses(t *testing.T) {
},
}
if err := m.UpdateConfig(context.Background(), newConfig); err != nil {
if _, err := m.UpdateConfig(context.Background(), newConfig); err != nil {
t.Fatalf("failed to UpdateConfig: %v", err)
}
+14 -9
View File
@@ -68,19 +68,24 @@ func (s *Service) doSync(ctx context.Context) {
// 不在 sync 层做版本早退,由 frpcManager.UpdateConfig 负责判断。
// 原因:重启后进程全部消失,即使版本/checksum 未变,仍需重新拉起 frpc 进程。
err = s.frpcManager.UpdateConfig(ctx, configResp)
result := "success"
message := "apply success"
configChanged, err := s.frpcManager.UpdateConfig(ctx, configResp)
if err != nil {
result = "failed"
message = err.Error()
slog.Error("failed to apply tunnel config", "error", err)
} else {
slog.Info("tunnel config applied successfully", "version", configResp.Version)
s.reportApplyLog(ctx, configResp, "failed", err.Error())
return
}
if configChanged {
slog.Info("tunnel config applied successfully", "version", configResp.Version)
s.reportApplyLog(ctx, configResp, "success", "apply success")
return
}
slog.Debug("tunnel config unchanged, skipping apply log report", "version", configResp.Version)
}
// Report apply log
func (s *Service) reportApplyLog(ctx context.Context, configResp *service.FlaredTunnelConfigResponse, result string, message string) {
if configResp == nil {
return
}
logPayload := service.ApplyLogPayload{
Version: configResp.Version,
Result: result,
+39 -14
View File
@@ -181,6 +181,17 @@ func ReportApplyLog(ctx context.Context, payload ApplyLogPayload) (*model.OpenFl
return nil, errors.New(errInvalidApplyResult)
}
latest, err := model.GetLatestOpenFlareApplyLogByNodeID(ctx, payload.NodeID)
if err != nil {
return nil, err
}
if model.IsRepeatSuccessApplyLog(latest, payload.Version, payload.Checksum, payload.Result) {
if err := updateNodeFromApplyLog(ctx, payload, now); err != nil {
return nil, err
}
return latest, nil
}
log := &model.OpenFlareApplyLog{
NodeID: payload.NodeID,
Version: payload.Version,
@@ -198,23 +209,11 @@ func ReportApplyLog(ctx context.Context, payload ApplyLogPayload) (*model.OpenFl
return nil, errors.New("database not initialized")
}
err := conn.Transaction(func(tx *gorm.DB) error {
record := &model.OpenFlareNode{}
if err := tx.Where("node_id = ?", payload.NodeID).First(record).Error; err != nil {
return err
}
record.Status = nodeStatusOnline
record.LastSeenAt = &now
if payload.Result == applyResultOK {
record.CurrentVersion = payload.Version
record.LastError = ""
} else {
record.LastError = payload.Message
}
err = conn.Transaction(func(tx *gorm.DB) error {
if err := tx.Create(log).Error; err != nil {
return err
}
return tx.Model(record).Select("status", "last_seen_at", "current_version", "last_error").Updates(record).Error
return updateNodeFromApplyLogTx(tx, payload, now)
})
if err != nil {
return nil, err
@@ -222,6 +221,32 @@ func ReportApplyLog(ctx context.Context, payload ApplyLogPayload) (*model.OpenFl
return log, nil
}
func updateNodeFromApplyLog(ctx context.Context, payload ApplyLogPayload, now time.Time) error {
conn := db.DB(ctx)
if conn == nil {
return errors.New("database not initialized")
}
return conn.Transaction(func(tx *gorm.DB) error {
return updateNodeFromApplyLogTx(tx, payload, now)
})
}
func updateNodeFromApplyLogTx(tx *gorm.DB, payload ApplyLogPayload, now time.Time) error {
record := &model.OpenFlareNode{}
if err := tx.Where("node_id = ?", payload.NodeID).First(record).Error; err != nil {
return err
}
record.Status = nodeStatusOnline
record.LastSeenAt = &now
if payload.Result == applyResultOK {
record.CurrentVersion = payload.Version
record.LastError = ""
} else {
record.LastError = payload.Message
}
return tx.Model(record).Select("status", "last_seen_at", "current_version", "last_error").Updates(record).Error
}
// ValidateDiscoveryToken delegates to the node package discovery token helper.
func ValidateDiscoveryToken(ctx context.Context, token string) error {
return node.ValidateDiscoveryToken(ctx, token)
+40 -15
View File
@@ -170,6 +170,17 @@ func ReportApplyLog(ctx context.Context, payload ApplyLogPayload) (*model.OpenFl
return nil, errors.New("result 仅支持 success、warning 或 failed")
}
latest, err := model.GetLatestOpenFlareApplyLogByNodeID(ctx, payload.NodeID)
if err != nil {
return nil, err
}
if model.IsRepeatSuccessApplyLog(latest, payload.Version, payload.Checksum, payload.Result) {
if err := updateFlaredNodeFromApplyLog(ctx, payload, now); err != nil {
return nil, err
}
return latest, nil
}
log := &model.OpenFlareApplyLog{
NodeID: payload.NodeID,
Version: payload.Version,
@@ -182,24 +193,11 @@ func ReportApplyLog(ctx context.Context, payload ApplyLogPayload) (*model.OpenFl
CreatedAt: now,
}
err := db.DB(ctx).Transaction(func(tx *gorm.DB) error {
var node model.OpenFlareNode
if err := tx.Where("node_id = ?", payload.NodeID).First(&node).Error; err != nil {
return err
}
node.Status = nodeStatusOnline
lastSeen := now
node.LastSeenAt = &lastSeen
if payload.Result == applyResultOK {
node.CurrentVersion = payload.Version
node.LastError = ""
} else {
node.LastError = payload.Message
}
err = db.DB(ctx).Transaction(func(tx *gorm.DB) error {
if err := tx.Create(log).Error; err != nil {
return err
}
return tx.Model(&node).Select("status", "last_seen_at", "current_version", "last_error").Updates(&node).Error
return updateFlaredNodeFromApplyLogTx(tx, payload, now)
})
if err != nil {
return nil, err
@@ -207,6 +205,33 @@ func ReportApplyLog(ctx context.Context, payload ApplyLogPayload) (*model.OpenFl
return log, nil
}
func updateFlaredNodeFromApplyLog(ctx context.Context, payload ApplyLogPayload, now time.Time) error {
conn := db.DB(ctx)
if conn == nil {
return errors.New("database not initialized")
}
return conn.Transaction(func(tx *gorm.DB) error {
return updateFlaredNodeFromApplyLogTx(tx, payload, now)
})
}
func updateFlaredNodeFromApplyLogTx(tx *gorm.DB, payload ApplyLogPayload, now time.Time) error {
var node model.OpenFlareNode
if err := tx.Where("node_id = ?", payload.NodeID).First(&node).Error; err != nil {
return err
}
node.Status = nodeStatusOnline
lastSeen := now
node.LastSeenAt = &lastSeen
if payload.Result == applyResultOK {
node.CurrentVersion = payload.Version
node.LastError = ""
} else {
node.LastError = payload.Message
}
return tx.Model(&node).Select("status", "last_seen_at", "current_version", "last_error").Updates(&node).Error
}
func getActiveConfigVersion(ctx context.Context) (*configVersionRow, error) {
conn := db.DB(ctx)
if conn == nil {