fix(session): warn invalid backlog and defer file delete retries
This commit is contained in:
parent
49b69cda81
commit
d8683deb57
4 changed files with 201 additions and 4 deletions
|
|
@ -3,6 +3,7 @@ package config
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"os"
|
"os"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
|
|
||||||
|
|
@ -14,6 +15,15 @@ import (
|
||||||
// rrCounter is a global counter for round-robin load balancing across models.
|
// rrCounter is a global counter for round-robin load balancing across models.
|
||||||
var rrCounter atomic.Uint64
|
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,
|
// FlexibleStringSlice is a []string that also accepts JSON numbers,
|
||||||
// so allow_from can contain both "123" and 123.
|
// so allow_from can contain both "123" and 123.
|
||||||
type FlexibleStringSlice []string
|
type FlexibleStringSlice []string
|
||||||
|
|
@ -620,6 +630,11 @@ func LoadConfig(path string) (*Config, error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
if cfg.Session.BacklogLimit < 1 {
|
if cfg.Session.BacklogLimit < 1 {
|
||||||
|
warnf(
|
||||||
|
"invalid session.backlog_limit=%d, fallback to default=%d",
|
||||||
|
cfg.Session.BacklogLimit,
|
||||||
|
DefaultSessionBacklogLimit,
|
||||||
|
)
|
||||||
cfg.Session.BacklogLimit = DefaultSessionBacklogLimit
|
cfg.Session.BacklogLimit = DefaultSessionBacklogLimit
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,7 @@
|
||||||
package config
|
package config
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
|
@ -458,6 +459,12 @@ func TestDefaultConfig_SessionBacklogLimit(t *testing.T) {
|
||||||
func TestLoadConfig_InvalidBacklogLimitFallsBackToDefault(t *testing.T) {
|
func TestLoadConfig_InvalidBacklogLimitFallsBackToDefault(t *testing.T) {
|
||||||
tempDir := t.TempDir()
|
tempDir := t.TempDir()
|
||||||
configPath := filepath.Join(tempDir, "config.json")
|
configPath := filepath.Join(tempDir, "config.json")
|
||||||
|
var warnings bytes.Buffer
|
||||||
|
oldWarningWriter := warningWriter
|
||||||
|
warningWriter = &warnings
|
||||||
|
t.Cleanup(func() {
|
||||||
|
warningWriter = oldWarningWriter
|
||||||
|
})
|
||||||
|
|
||||||
configJSON := `{
|
configJSON := `{
|
||||||
"agents": {"defaults":{"workspace":"./workspace","model":"gpt4","max_tokens":8192,"max_tool_iterations":20}},
|
"agents": {"defaults":{"workspace":"./workspace","model":"gpt4","max_tokens":8192,"max_tool_iterations":20}},
|
||||||
|
|
@ -479,4 +486,7 @@ func TestLoadConfig_InvalidBacklogLimitFallsBackToDefault(t *testing.T) {
|
||||||
DefaultSessionBacklogLimit,
|
DefaultSessionBacklogLimit,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
if got := warnings.String(); !strings.Contains(got, "invalid session.backlog_limit=0") {
|
||||||
|
t.Fatalf("expected warning about invalid backlog_limit, got: %q", got)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,9 @@ package session
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
|
@ -23,6 +25,11 @@ type Session struct {
|
||||||
|
|
||||||
const sessionIndexFilename = "index.json"
|
const sessionIndexFilename = "index.json"
|
||||||
|
|
||||||
|
var (
|
||||||
|
removeFile = os.Remove
|
||||||
|
warningWriter io.Writer = os.Stderr
|
||||||
|
)
|
||||||
|
|
||||||
type scopeIndex struct {
|
type scopeIndex struct {
|
||||||
ActiveSessionKey string `json:"active_session_key"`
|
ActiveSessionKey string `json:"active_session_key"`
|
||||||
OrderedSessions []string `json:"ordered_sessions"`
|
OrderedSessions []string `json:"ordered_sessions"`
|
||||||
|
|
@ -30,8 +37,9 @@ type scopeIndex struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
type sessionIndex struct {
|
type sessionIndex struct {
|
||||||
Version int `json:"version"`
|
Version int `json:"version"`
|
||||||
Scopes map[string]*scopeIndex `json:"scopes"`
|
Scopes map[string]*scopeIndex `json:"scopes"`
|
||||||
|
PendingDeletes []string `json:"pending_deletes,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type SessionMeta struct {
|
type SessionMeta struct {
|
||||||
|
|
@ -273,8 +281,31 @@ func (sm *SessionManager) DeleteSession(sessionKey string) error {
|
||||||
sm.mu.Unlock()
|
sm.mu.Unlock()
|
||||||
|
|
||||||
if err := sm.deleteSessionFile(sessionKey); err != nil {
|
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
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -508,6 +539,40 @@ func (sm *SessionManager) loadIndex() error {
|
||||||
}
|
}
|
||||||
|
|
||||||
changed := false
|
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 {
|
for scopeKey, scope := range loaded.Scopes {
|
||||||
if scope == nil {
|
if scope == nil {
|
||||||
delete(loaded.Scopes, scopeKey)
|
delete(loaded.Scopes, scopeKey)
|
||||||
|
|
@ -677,6 +742,43 @@ func cloneSession(stored *Session) Session {
|
||||||
return snapshot
|
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 {
|
func (sm *SessionManager) saveSessionLocked(key string) error {
|
||||||
if sm.storage == "" {
|
if sm.storage == "" {
|
||||||
return nil
|
return nil
|
||||||
|
|
@ -758,7 +860,7 @@ func (sm *SessionManager) deleteSessionFile(sessionKey string) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
sessionPath := filepath.Join(sm.storage, filename+".json")
|
sessionPath := filepath.Join(sm.storage, filename+".json")
|
||||||
if err := os.Remove(sessionPath); err != nil {
|
if err := removeFile(sessionPath); err != nil {
|
||||||
if os.IsNotExist(err) {
|
if os.IsNotExist(err) {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,12 @@
|
||||||
package session
|
package session
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -383,3 +386,70 @@ func TestLoadIndex_SelfHealsStaleReferences(t *testing.T) {
|
||||||
t.Fatalf("list[0]=%+v, want active newest", list[0])
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue