From 89b65cfb2d56f54880b4b8d3727fb14b2d781c81 Mon Sep 17 00:00:00 2001 From: Syngnat Date: Wed, 22 Jul 2026 09:00:33 +0800 Subject: [PATCH] =?UTF-8?q?=E2=9A=A1=EF=B8=8F=20perf(result-diff):=20?= =?UTF-8?q?=E4=B8=BB=E5=8A=A8=E5=9B=9E=E6=94=B6=E8=BF=87=E6=9C=9F=E6=AF=94?= =?UTF-8?q?=E5=AF=B9=E4=BC=9A=E8=AF=9D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/app/app.go | 1 + internal/app/app_keepalive.go | 3 + internal/app/methods_result_diff.go | 33 +++- internal/app/result_diff_lifecycle_test.go | 40 +++++ internal/resultdiff/session.go | 185 +++++++++++++++++++- internal/resultdiff/session_manager_test.go | 175 ++++++++++++++++++ 6 files changed, 422 insertions(+), 15 deletions(-) create mode 100644 internal/app/result_diff_lifecycle_test.go create mode 100644 internal/resultdiff/session_manager_test.go diff --git a/internal/app/app.go b/internal/app/app.go index 2d26c893..c7e2f209 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -405,6 +405,7 @@ func (a *App) Shutdown() { logger.Infof("应用开始关闭,准备释放资源") a.beginDatabaseShutdown() a.stopConnectionKeepAliveLoop() + a.closeResultDiffSessions() a.rollbackPendingSQLTransactionsOnShutdown() a.closeSQLAuditStore() a.closeCachedDatabasesForShutdown() diff --git a/internal/app/app_keepalive.go b/internal/app/app_keepalive.go index d8dcd6be..c5657cf6 100644 --- a/internal/app/app_keepalive.go +++ b/internal/app/app_keepalive.go @@ -160,6 +160,9 @@ func (a *App) runConnectionKeepAliveTickContext(ctx context.Context, now time.Ti if ctx == nil { ctx = context.Background() } + if a != nil && a.resultDiffManager != nil { + a.resultDiffManager.PruneExpired(now) + } targets := a.collectDueConnectionKeepAliveTargets(now) for index, target := range targets { if ctx.Err() != nil { diff --git a/internal/app/methods_result_diff.go b/internal/app/methods_result_diff.go index 7b94adc4..8fda66cf 100644 --- a/internal/app/methods_result_diff.go +++ b/internal/app/methods_result_diff.go @@ -51,7 +51,13 @@ func (a *App) ResultDiffStart(req ResultDiffStartRequest) connection.QueryResult MaxRowsPerSide: req.MaxRowsPerSide, IncludeSameRows: req.IncludeSameRows, } - session := a.resultDiffManager.Create(startReq) + session, releaseSession := a.resultDiffManager.CreateWithLease(startReq) + if session == nil { + return connection.QueryResult{Success: false, Message: a.appText("result_diff.backend.error.manager_unavailable", nil)} + } + defer func() { + releaseSession() + }() // 纯 rows 模式:仅创建会话,等待分块上传后 ResultDiffCompute if leftMode == resultdiff.DatasetModeRows && rightMode == resultdiff.DatasetModeRows { @@ -122,12 +128,14 @@ func (a *App) ResultDiffStart(req ResultDiffStartRequest) connection.QueryResult // 混合模式:rows 侧可能稍后上传 if leftMode == resultdiff.DatasetModeRows && len(leftRows) == 0 { - // 先装 right,left 等待 upload - _ = session.SetLoaded(nil, nil, rightCols, rightRows) - // SetLoaded 会把两侧都 mark done — 不适合混合。改用 Append。 - // 重新创建更简单: + // 先装 right,left 等待 upload。重新创建可避免 SetLoaded 把两侧都 + // 标记完成;切换 lease 前先结束旧会话的装载保护。 + releaseSession() a.resultDiffManager.Close(session.ID) - session = a.resultDiffManager.Create(startReq) + session, releaseSession = a.resultDiffManager.CreateWithLease(startReq) + if session == nil { + return connection.QueryResult{Success: false, Message: a.appText("result_diff.backend.error.manager_unavailable", nil)} + } if err := session.AppendRows("right", rightCols, rightRows, true); err != nil { a.resultDiffManager.Close(session.ID) return connection.QueryResult{Success: false, Message: err.Error()} @@ -139,8 +147,12 @@ func (a *App) ResultDiffStart(req ResultDiffStartRequest) connection.QueryResult } } if rightMode == resultdiff.DatasetModeRows && len(rightRows) == 0 { + releaseSession() a.resultDiffManager.Close(session.ID) - session = a.resultDiffManager.Create(startReq) + session, releaseSession = a.resultDiffManager.CreateWithLease(startReq) + if session == nil { + return connection.QueryResult{Success: false, Message: a.appText("result_diff.backend.error.manager_unavailable", nil)} + } if err := session.AppendRows("left", leftCols, leftRows, true); err != nil { a.resultDiffManager.Close(session.ID) return connection.QueryResult{Success: false, Message: err.Error()} @@ -228,6 +240,13 @@ func (a *App) ResultDiffClose(jobID string) connection.QueryResult { return connection.QueryResult{Success: true, Message: a.appText("result_diff.backend.result.closed", nil)} } +func (a *App) closeResultDiffSessions() int { + if a == nil || a.resultDiffManager == nil { + return 0 + } + return a.resultDiffManager.Shutdown() +} + func normalizeDatasetMode(mode resultdiff.DatasetMode) resultdiff.DatasetMode { switch strings.ToLower(strings.TrimSpace(string(mode))) { case "sql": diff --git a/internal/app/result_diff_lifecycle_test.go b/internal/app/result_diff_lifecycle_test.go new file mode 100644 index 00000000..c7fc07b9 --- /dev/null +++ b/internal/app/result_diff_lifecycle_test.go @@ -0,0 +1,40 @@ +package app + +import ( + "testing" + "time" + + "GoNavi-Wails/internal/resultdiff" +) + +func TestConnectionKeepAliveTickPrunesExpiredResultDiffSessions(t *testing.T) { + application := NewApp() + application.resultDiffManager = resultdiff.NewManager(time.Minute) + session := application.resultDiffManager.Create(resultdiff.StartRequest{KeyColumns: []string{"id"}}) + + application.runConnectionKeepAliveTick(session.CreatedAt.Add(time.Minute + time.Nanosecond)) + + if _, err := application.resultDiffManager.Get(session.ID); err == nil { + t.Fatal("maintenance tick left expired result diff session reachable") + } +} + +func TestCloseResultDiffSessionsReleasesManagerState(t *testing.T) { + application := NewApp() + session := application.resultDiffManager.Create(resultdiff.StartRequest{KeyColumns: []string{"id"}}) + + if closed := application.closeResultDiffSessions(); closed != 1 { + t.Fatalf("closeResultDiffSessions closed %d sessions, want 1", closed) + } + if _, err := application.resultDiffManager.Get(session.ID); err == nil { + t.Fatal("shutdown cleanup left result diff session reachable") + } + start := application.ResultDiffStart(ResultDiffStartRequest{ + Left: resultdiff.DatasetSpec{Mode: resultdiff.DatasetModeRows}, + Right: resultdiff.DatasetSpec{Mode: resultdiff.DatasetModeRows}, + KeyColumns: []string{"id"}, + }) + if start.Success { + t.Fatal("result diff manager accepted a new session after shutdown cleanup") + } +} diff --git a/internal/resultdiff/session.go b/internal/resultdiff/session.go index 3adc900b..d918f659 100644 --- a/internal/resultdiff/session.go +++ b/internal/resultdiff/session.go @@ -22,6 +22,15 @@ type Session struct { mu sync.RWMutex + // activityMu is independent from the data lock so pruning never waits for a + // potentially expensive diff to finish. time.Time retains its monotonic + // component, avoiding TTL decisions based on wall-clock jumps. + activityMu sync.Mutex + lastTouched time.Time + activeOperations int + now func() time.Time + owner *Manager + // 装载缓冲(rows 模式) leftColumns []string rightColumns []string @@ -42,16 +51,26 @@ type Manager struct { mu sync.Mutex sessions map[string]*Session ttl time.Duration + now func() time.Time + closed bool } // NewManager 创建会话管理器。 func NewManager(ttl time.Duration) *Manager { + return newManagerWithClock(ttl, time.Now) +} + +func newManagerWithClock(ttl time.Duration, now func() time.Time) *Manager { if ttl <= 0 { ttl = 30 * time.Minute } + if now == nil { + now = time.Now + } m := &Manager{ sessions: make(map[string]*Session), ttl: ttl, + now: now, } return m } @@ -60,8 +79,34 @@ func NewManager(ttl time.Duration) *Manager { func (m *Manager) Create(req StartRequest) *Session { m.mu.Lock() defer m.mu.Unlock() - m.gcLocked() + if m.closed { + return nil + } + now := m.now() + m.pruneExpiredLocked(now) + return m.createLocked(req, now) +} +// CreateWithLease creates a session and keeps it active until release is +// called. It is used while initial SQL datasets are loaded, before the first +// normal Get/Session operation can refresh the TTL. +func (m *Manager) CreateWithLease(req StartRequest) (*Session, func()) { + m.mu.Lock() + defer m.mu.Unlock() + if m.closed { + return nil, func() {} + } + now := m.now() + m.pruneExpiredLocked(now) + session := m.createLocked(req, now) + release := session.beginActivityAt(now) + var releaseOnce sync.Once + return session, func() { + releaseOnce.Do(release) + } +} + +func (m *Manager) createLocked(req StartRequest, now time.Time) *Session { id := strings.TrimSpace(req.JobID) if id == "" { id = "rdiff-" + uuid.NewString() @@ -72,7 +117,7 @@ func (m *Manager) Create(req StartRequest) *Session { } s := &Session{ ID: id, - CreatedAt: time.Now(), + CreatedAt: now, KeyColumns: normalizeColumnList(req.KeyColumns), CompareColumns: normalizeColumnList(req.CompareColumns), IgnoreColumns: normalizeColumnList(req.IgnoreColumns), @@ -81,7 +126,10 @@ func (m *Manager) Create(req StartRequest) *Session { IncludeSameRows: req.IncludeSameRows, leftRows: make([]map[string]interface{}, 0), rightRows: make([]map[string]interface{}, 0), + now: m.now, + owner: m, } + s.touch(now) m.sessions[id] = s return s } @@ -90,11 +138,16 @@ func (m *Manager) Create(req StartRequest) *Session { func (m *Manager) Get(jobID string) (*Session, error) { m.mu.Lock() defer m.mu.Unlock() - m.gcLocked() + if m.closed { + return nil, fmt.Errorf("result diff manager is closed") + } + now := m.now() + m.pruneExpiredLocked(now) s, ok := m.sessions[strings.TrimSpace(jobID)] if !ok { return nil, fmt.Errorf("result diff job not found: %s", jobID) } + s.touch(now) return s, nil } @@ -105,20 +158,128 @@ func (m *Manager) Close(jobID string) { delete(m.sessions, strings.TrimSpace(jobID)) } -func (m *Manager) gcLocked() { - if m.ttl <= 0 { - return +// CloseAll releases every managed session and returns the number removed. +func (m *Manager) CloseAll() int { + if m == nil { + return 0 } - now := time.Now() + m.mu.Lock() + defer m.mu.Unlock() + closed := len(m.sessions) + m.sessions = make(map[string]*Session) + return closed +} + +// Shutdown permanently closes the manager, releases every session, and rejects +// later creation or lookup attempts. +func (m *Manager) Shutdown() int { + if m == nil { + return 0 + } + m.mu.Lock() + defer m.mu.Unlock() + closed := len(m.sessions) + m.sessions = make(map[string]*Session) + m.closed = true + return closed +} + +// PruneExpired removes sessions idle for longer than the configured TTL. +// Passing a time explicitly keeps the maintenance path deterministic and lets +// callers share an existing ticker instead of creating another goroutine. +func (m *Manager) PruneExpired(now time.Time) int { + if m == nil { + return 0 + } + m.mu.Lock() + defer m.mu.Unlock() + if now.IsZero() { + now = m.now() + } + return m.pruneExpiredLocked(now) +} + +func (m *Manager) pruneExpiredLocked(now time.Time) int { + if m.ttl <= 0 { + return 0 + } + removed := 0 for id, s := range m.sessions { - if now.Sub(s.CreatedAt) > m.ttl { + lastTouched, activeOperations := s.activityState() + if activeOperations > 0 { + continue + } + if now.Sub(lastTouched) > m.ttl { delete(m.sessions, id) + removed++ } } + return removed +} + +func (s *Session) touch(now time.Time) { + if s == nil { + return + } + s.activityMu.Lock() + s.lastTouched = now + s.activityMu.Unlock() +} + +func (s *Session) activityState() (time.Time, int) { + if s == nil { + return time.Time{}, 0 + } + s.activityMu.Lock() + defer s.activityMu.Unlock() + lastTouched := s.lastTouched + if lastTouched.IsZero() { + lastTouched = s.CreatedAt + } + return lastTouched, s.activeOperations +} + +func (s *Session) beginActivity() func() { + if s == nil { + return func() {} + } + if s.owner == nil { + return s.beginActivityAt(s.currentTime()) + } + + s.owner.mu.Lock() + defer s.owner.mu.Unlock() + if s.owner.closed || s.owner.sessions[s.ID] != s { + return func() {} + } + return s.beginActivityAt(s.owner.now()) +} + +func (s *Session) beginActivityAt(now time.Time) func() { + s.activityMu.Lock() + s.lastTouched = now + s.activeOperations++ + s.activityMu.Unlock() + return func() { + finishedAt := s.currentTime() + s.activityMu.Lock() + s.lastTouched = finishedAt + s.activeOperations-- + s.activityMu.Unlock() + } +} + +func (s *Session) currentTime() time.Time { + if s != nil && s.now != nil { + return s.now() + } + return time.Now() } // AppendRows 向一侧追加行。 func (s *Session) AppendRows(side string, columns []string, rows []map[string]interface{}, done bool) error { + finishActivity := s.beginActivity() + defer finishActivity() s.mu.Lock() defer s.mu.Unlock() if s.computed { @@ -169,6 +330,8 @@ func (s *Session) SetLoaded( leftCols []string, leftRows []map[string]interface{}, rightCols []string, rightRows []map[string]interface{}, ) error { + finishActivity := s.beginActivity() + defer finishActivity() s.mu.Lock() defer s.mu.Unlock() if len(leftRows) > s.MaxRowsPerSide { @@ -188,6 +351,8 @@ func (s *Session) SetLoaded( // Compute 执行 diff。 func (s *Session) Compute() (Summary, error) { + finishActivity := s.beginActivity() + defer finishActivity() s.mu.Lock() defer s.mu.Unlock() if s.computed { @@ -213,6 +378,8 @@ func (s *Session) Compute() (Summary, error) { // Page 分页。 func (s *Session) Page(req PageRequest) (PageResult, error) { + finishActivity := s.beginActivity() + defer finishActivity() s.mu.RLock() defer s.mu.RUnlock() if !s.computed { @@ -226,6 +393,8 @@ func (s *Session) Page(req PageRequest) (PageResult, error) { // Summary 返回汇总。 func (s *Session) Summary() (Summary, bool) { + finishActivity := s.beginActivity() + defer finishActivity() s.mu.RLock() defer s.mu.RUnlock() return s.summary, s.computed diff --git a/internal/resultdiff/session_manager_test.go b/internal/resultdiff/session_manager_test.go new file mode 100644 index 00000000..389c99e3 --- /dev/null +++ b/internal/resultdiff/session_manager_test.go @@ -0,0 +1,175 @@ +package resultdiff + +import ( + "sync" + "testing" + "time" +) + +type managerTestClock struct { + now time.Time +} + +func (c *managerTestClock) Now() time.Time { + return c.now +} + +func (c *managerTestClock) Advance(delta time.Duration) { + c.now = c.now.Add(delta) +} + +func TestManagerPruneExpiredWithoutFollowupActivity(t *testing.T) { + clock := &managerTestClock{now: time.Date(2026, time.July, 22, 12, 0, 0, 0, time.UTC)} + manager := newManagerWithClock(time.Minute, clock.Now) + session := manager.Create(StartRequest{KeyColumns: []string{"id"}}) + + clock.Advance(time.Minute + time.Nanosecond) + if removed := manager.PruneExpired(clock.Now()); removed != 1 { + t.Fatalf("PruneExpired removed %d sessions, want 1", removed) + } + if removed := manager.PruneExpired(clock.Now()); removed != 0 { + t.Fatalf("second PruneExpired removed %d sessions, want 0", removed) + } + if _, err := manager.Get(session.ID); err == nil { + t.Fatal("expired session is still reachable after proactive prune") + } +} + +func TestManagerPruneExpiredKeepsRecentlyAccessedSession(t *testing.T) { + clock := &managerTestClock{now: time.Date(2026, time.July, 22, 12, 0, 0, 0, time.UTC)} + manager := newManagerWithClock(time.Minute, clock.Now) + session := manager.Create(StartRequest{KeyColumns: []string{"id"}}) + + clock.Advance(45 * time.Second) + if _, err := manager.Get(session.ID); err != nil { + t.Fatalf("Get active session: %v", err) + } + clock.Advance(45 * time.Second) + if removed := manager.PruneExpired(clock.Now()); removed != 0 { + t.Fatalf("PruneExpired removed %d recently accessed sessions, want 0", removed) + } + + clock.Advance(time.Minute + time.Nanosecond) + if removed := manager.PruneExpired(clock.Now()); removed != 1 { + t.Fatalf("PruneExpired removed %d idle sessions after refreshed TTL, want 1", removed) + } +} + +func TestManagerSessionActivityRefreshesTTL(t *testing.T) { + clock := &managerTestClock{now: time.Date(2026, time.July, 22, 12, 0, 0, 0, time.UTC)} + manager := newManagerWithClock(time.Minute, clock.Now) + session := manager.Create(StartRequest{KeyColumns: []string{"id"}, MaxRowsPerSide: 10}) + + clock.Advance(45 * time.Second) + if err := session.AppendRows("left", []string{"id"}, []map[string]interface{}{{"id": 1}}, false); err != nil { + t.Fatalf("AppendRows: %v", err) + } + clock.Advance(45 * time.Second) + if removed := manager.PruneExpired(clock.Now()); removed != 0 { + t.Fatalf("PruneExpired removed %d active upload sessions, want 0", removed) + } +} + +func TestManagerPruneExpiredSkipsOperationInProgress(t *testing.T) { + clock := &managerTestClock{now: time.Date(2026, time.July, 22, 12, 0, 0, 0, time.UTC)} + manager := newManagerWithClock(time.Minute, clock.Now) + session := manager.Create(StartRequest{KeyColumns: []string{"id"}}) + + finishActivity := session.beginActivity() + defer finishActivity() + clock.Advance(2 * time.Minute) + if removed := manager.PruneExpired(clock.Now()); removed != 0 { + t.Fatalf("PruneExpired removed %d sessions with an operation in progress, want 0", removed) + } +} + +func TestManagerCreateWithLeaseProtectsLongInitialLoad(t *testing.T) { + clock := &managerTestClock{now: time.Date(2026, time.July, 22, 12, 0, 0, 0, time.UTC)} + manager := newManagerWithClock(time.Minute, clock.Now) + session, release := manager.CreateWithLease(StartRequest{KeyColumns: []string{"id"}}) + if session == nil { + t.Fatal("CreateWithLease returned nil session before shutdown") + } + + clock.Advance(2 * time.Minute) + if removed := manager.PruneExpired(clock.Now()); removed != 0 { + t.Fatalf("PruneExpired removed %d sessions during initial load lease, want 0", removed) + } + release() + release() + if _, activeOperations := session.activityState(); activeOperations != 0 { + t.Fatalf("idempotent lease release left %d active operations, want 0", activeOperations) + } + clock.Advance(time.Minute + time.Nanosecond) + if removed := manager.PruneExpired(clock.Now()); removed != 1 { + t.Fatalf("PruneExpired removed %d sessions after initial load became idle, want 1", removed) + } +} + +func TestManagerCloseAllReleasesEverySession(t *testing.T) { + manager := NewManager(time.Hour) + first := manager.Create(StartRequest{KeyColumns: []string{"id"}}) + second := manager.Create(StartRequest{KeyColumns: []string{"id"}}) + + if closed := manager.CloseAll(); closed != 2 { + t.Fatalf("CloseAll closed %d sessions, want 2", closed) + } + if closed := manager.CloseAll(); closed != 0 { + t.Fatalf("second CloseAll closed %d sessions, want 0", closed) + } + for _, jobID := range []string{first.ID, second.ID} { + if _, err := manager.Get(jobID); err == nil { + t.Fatalf("session %s is still reachable after CloseAll", jobID) + } + } + if replacement := manager.Create(StartRequest{KeyColumns: []string{"id"}}); replacement == nil { + t.Fatal("CloseAll unexpectedly made the manager terminal") + } +} + +func TestManagerShutdownRejectsNewSessions(t *testing.T) { + manager := NewManager(time.Hour) + existing := manager.Create(StartRequest{KeyColumns: []string{"id"}}) + + if closed := manager.Shutdown(); closed != 1 { + t.Fatalf("Shutdown closed %d sessions, want 1", closed) + } + if closed := manager.Shutdown(); closed != 0 { + t.Fatalf("second Shutdown closed %d sessions, want 0", closed) + } + if manager.Create(StartRequest{KeyColumns: []string{"id"}}) != nil { + t.Fatal("Create accepted a session after shutdown") + } + if session, release := manager.CreateWithLease(StartRequest{KeyColumns: []string{"id"}}); session != nil { + release() + t.Fatal("CreateWithLease accepted a session after shutdown") + } + if _, err := manager.Get(existing.ID); err == nil { + t.Fatal("Get accepted a session after shutdown") + } +} + +func TestManagerShutdownPreventsConcurrentSessionResurrection(t *testing.T) { + manager := NewManager(time.Hour) + const workerCount = 32 + start := make(chan struct{}) + var workers sync.WaitGroup + workers.Add(workerCount) + for index := 0; index < workerCount; index++ { + go func() { + defer workers.Done() + <-start + manager.Create(StartRequest{KeyColumns: []string{"id"}}) + }() + } + + close(start) + manager.Shutdown() + workers.Wait() + if remaining := manager.Shutdown(); remaining != 0 { + t.Fatalf("shutdown left %d concurrently created sessions reachable", remaining) + } + if manager.Create(StartRequest{KeyColumns: []string{"id"}}) != nil { + t.Fatal("Create resurrected manager after concurrent shutdown") + } +}