From d8683deb5765d518f064f18c02134248d4a62e23 Mon Sep 17 00:00:00 2001 From: mingmxren Date: Sun, 1 Mar 2026 12:27:23 +0800 Subject: [PATCH] fix(session): warn invalid backlog and defer file delete retries --- pkg/config/config.go | 15 +++++ pkg/config/config_test.go | 10 ++++ pkg/session/manager.go | 110 ++++++++++++++++++++++++++++++++++-- pkg/session/manager_test.go | 70 +++++++++++++++++++++++ 4 files changed, 201 insertions(+), 4 deletions(-) diff --git a/pkg/config/config.go b/pkg/config/config.go index 04e7d69ce..794efa50c 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -3,6 +3,7 @@ package config import ( "encoding/json" "fmt" + "io" "os" "sync/atomic" @@ -14,6 +15,15 @@ import ( // rrCounter is a global counter for round-robin load balancing across models. var rrCounter atomic.Uint64 +var warningWriter io.Writer = os.Stderr + +func warnf(format string, args ...any) { + if warningWriter == nil { + return + } + _, _ = fmt.Fprintf(warningWriter, "warning: "+format+"\n", args...) +} + // FlexibleStringSlice is a []string that also accepts JSON numbers, // so allow_from can contain both "123" and 123. type FlexibleStringSlice []string @@ -620,6 +630,11 @@ func LoadConfig(path string) (*Config, error) { } if cfg.Session.BacklogLimit < 1 { + warnf( + "invalid session.backlog_limit=%d, fallback to default=%d", + cfg.Session.BacklogLimit, + DefaultSessionBacklogLimit, + ) cfg.Session.BacklogLimit = DefaultSessionBacklogLimit } diff --git a/pkg/config/config_test.go b/pkg/config/config_test.go index fea057c45..1dd6d960a 100644 --- a/pkg/config/config_test.go +++ b/pkg/config/config_test.go @@ -1,6 +1,7 @@ package config import ( + "bytes" "encoding/json" "os" "path/filepath" @@ -458,6 +459,12 @@ func TestDefaultConfig_SessionBacklogLimit(t *testing.T) { func TestLoadConfig_InvalidBacklogLimitFallsBackToDefault(t *testing.T) { tempDir := t.TempDir() configPath := filepath.Join(tempDir, "config.json") + var warnings bytes.Buffer + oldWarningWriter := warningWriter + warningWriter = &warnings + t.Cleanup(func() { + warningWriter = oldWarningWriter + }) configJSON := `{ "agents": {"defaults":{"workspace":"./workspace","model":"gpt4","max_tokens":8192,"max_tool_iterations":20}}, @@ -479,4 +486,7 @@ func TestLoadConfig_InvalidBacklogLimitFallsBackToDefault(t *testing.T) { DefaultSessionBacklogLimit, ) } + if got := warnings.String(); !strings.Contains(got, "invalid session.backlog_limit=0") { + t.Fatalf("expected warning about invalid backlog_limit, got: %q", got) + } } diff --git a/pkg/session/manager.go b/pkg/session/manager.go index 974ce3fb6..2e1375822 100644 --- a/pkg/session/manager.go +++ b/pkg/session/manager.go @@ -2,7 +2,9 @@ package session import ( "encoding/json" + "errors" "fmt" + "io" "os" "path/filepath" "strconv" @@ -23,6 +25,11 @@ type Session struct { const sessionIndexFilename = "index.json" +var ( + removeFile = os.Remove + warningWriter io.Writer = os.Stderr +) + type scopeIndex struct { ActiveSessionKey string `json:"active_session_key"` OrderedSessions []string `json:"ordered_sessions"` @@ -30,8 +37,9 @@ type scopeIndex struct { } type sessionIndex struct { - Version int `json:"version"` - Scopes map[string]*scopeIndex `json:"scopes"` + Version int `json:"version"` + Scopes map[string]*scopeIndex `json:"scopes"` + PendingDeletes []string `json:"pending_deletes,omitempty"` } type SessionMeta struct { @@ -273,8 +281,31 @@ func (sm *SessionManager) DeleteSession(sessionKey string) error { sm.mu.Unlock() if err := sm.deleteSessionFile(sessionKey); err != nil { - return err + sm.warnf("failed to delete session file for %q, deferred retry on startup: %v", sessionKey, err) + sm.mu.Lock() + pendingChanged := sm.addPendingDeleteLocked(sessionKey) + if pendingChanged { + if err := sm.saveIndexLocked(); err != nil { + sm.mu.Unlock() + sm.warnf("failed to persist deferred delete for %q: %v", sessionKey, err) + return nil + } + } + sm.mu.Unlock() + return nil } + + sm.mu.Lock() + pendingChanged := sm.removePendingDeleteLocked(sessionKey) + if pendingChanged { + if err := sm.saveIndexLocked(); err != nil { + sm.mu.Unlock() + sm.warnf("failed to persist cleanup of deferred delete for %q: %v", sessionKey, err) + return nil + } + } + sm.mu.Unlock() + return nil } @@ -508,6 +539,40 @@ func (sm *SessionManager) loadIndex() error { } changed := false + seenPending := make(map[string]struct{}, len(loaded.PendingDeletes)) + retryPending := make([]string, 0, len(loaded.PendingDeletes)) + for _, sessionKey := range loaded.PendingDeletes { + if sessionKey == "" { + changed = true + continue + } + if _, dup := seenPending[sessionKey]; dup { + changed = true + continue + } + seenPending[sessionKey] = struct{}{} + + // Deferred-delete sessions should not be visible even if stale files remain. + delete(sm.sessions, sessionKey) + + if err := sm.deleteSessionFile(sessionKey); err != nil { + // Invalid paths are unrecoverable; drop them from retry queue. + if errors.Is(err, os.ErrInvalid) { + changed = true + sm.warnf("dropping invalid deferred delete key %q: %v", sessionKey, err) + continue + } + sm.warnf("retry deferred session delete failed for %q: %v", sessionKey, err) + retryPending = append(retryPending, sessionKey) + continue + } + changed = true + } + if len(retryPending) != len(loaded.PendingDeletes) { + changed = true + } + loaded.PendingDeletes = retryPending + for scopeKey, scope := range loaded.Scopes { if scope == nil { delete(loaded.Scopes, scopeKey) @@ -677,6 +742,43 @@ func cloneSession(stored *Session) Session { return snapshot } +func (sm *SessionManager) warnf(format string, args ...any) { + if warningWriter == nil { + return + } + _, _ = fmt.Fprintf(warningWriter, "warning: "+format+"\n", args...) +} + +func (sm *SessionManager) addPendingDeleteLocked(sessionKey string) bool { + if sessionKey == "" { + return false + } + for _, existing := range sm.index.PendingDeletes { + if existing == sessionKey { + return false + } + } + sm.index.PendingDeletes = append(sm.index.PendingDeletes, sessionKey) + return true +} + +func (sm *SessionManager) removePendingDeleteLocked(sessionKey string) bool { + if len(sm.index.PendingDeletes) == 0 { + return false + } + filtered := sm.index.PendingDeletes[:0] + removed := false + for _, existing := range sm.index.PendingDeletes { + if existing == sessionKey { + removed = true + continue + } + filtered = append(filtered, existing) + } + sm.index.PendingDeletes = filtered + return removed +} + func (sm *SessionManager) saveSessionLocked(key string) error { if sm.storage == "" { return nil @@ -758,7 +860,7 @@ func (sm *SessionManager) deleteSessionFile(sessionKey string) error { } sessionPath := filepath.Join(sm.storage, filename+".json") - if err := os.Remove(sessionPath); err != nil { + if err := removeFile(sessionPath); err != nil { if os.IsNotExist(err) { return nil } diff --git a/pkg/session/manager_test.go b/pkg/session/manager_test.go index 60d87214d..60af0d6cc 100644 --- a/pkg/session/manager_test.go +++ b/pkg/session/manager_test.go @@ -1,9 +1,12 @@ package session import ( + "bytes" "encoding/json" + "errors" "os" "path/filepath" + "strings" "testing" ) @@ -383,3 +386,70 @@ func TestLoadIndex_SelfHealsStaleReferences(t *testing.T) { t.Fatalf("list[0]=%+v, want active newest", list[0]) } } + +func TestDeleteSession_FileDeleteFailureIsDeferredAndRetriedOnStartup(t *testing.T) { + dir := t.TempDir() + sm := NewSessionManager(dir) + scope := "agent:main:telegram:direct:user1" + + if _, err := sm.ResolveActive(scope); err != nil { + t.Fatal(err) + } + sessionKey, err := sm.StartNew(scope) + if err != nil { + t.Fatal(err) + } + + oldRemoveFile := removeFile + oldWarningWriter := warningWriter + var warnings bytes.Buffer + warningWriter = &warnings + removeFile = func(path string) error { + if strings.HasSuffix(path, sanitizeFilename(sessionKey)+".json") { + return errors.New("permission denied") + } + return oldRemoveFile(path) + } + t.Cleanup(func() { + removeFile = oldRemoveFile + warningWriter = oldWarningWriter + }) + + if err := sm.DeleteSession(sessionKey); err != nil { + t.Fatalf("DeleteSession(%q) returned unexpected error: %v", sessionKey, err) + } + if got := warnings.String(); !strings.Contains(got, "deferred retry on startup") { + t.Fatalf("expected deferred-delete warning, got: %q", got) + } + + sessionPath := filepath.Join(dir, sanitizeFilename(sessionKey)+".json") + if _, err := os.Stat(sessionPath); err != nil { + t.Fatalf("expected deferred file %q to still exist, err=%v", sessionPath, err) + } + + list, err := sm.List(scope) + if err != nil { + t.Fatal(err) + } + for _, item := range list { + if item.SessionKey == sessionKey { + t.Fatalf("deleted session key %q should not remain in index list", sessionKey) + } + } + + if len(sm.index.PendingDeletes) != 1 || sm.index.PendingDeletes[0] != sessionKey { + t.Fatalf("pending_deletes=%v, want [%q]", sm.index.PendingDeletes, sessionKey) + } + + removeFile = oldRemoveFile + reloaded := NewSessionManager(dir) + if len(reloaded.index.PendingDeletes) != 0 { + t.Fatalf("pending_deletes should be drained on startup retry, got %v", reloaded.index.PendingDeletes) + } + if _, err := os.Stat(sessionPath); !os.IsNotExist(err) { + t.Fatalf("expected %q to be removed by startup retry, stat err=%v", sessionPath, err) + } + if _, exists := reloaded.sessions[sessionKey]; exists { + t.Fatalf("session %q should not be present in memory after startup retry", sessionKey) + } +}