refactor(session): update manager

This commit is contained in:
ZanzyTHEbar 2026-03-05 21:17:27 +00:00
parent ab4655901a
commit ccb2d99d28

View file

@ -361,7 +361,7 @@ func (sm *SessionManager) AddFullMessage(sessionKey string, msg messages.Message
sm.touchLRU(sessionKey) sm.touchLRU(sessionKey)
if sm.msgChan != nil { if sm.msgChan != nil {
sm.msgChan <- msgPersistItem{sessionKey: sessionKey, msg: msg} _ = sm.enqueuePersistItem(msgPersistItem{sessionKey: sessionKey, msg: msg})
} }
} }
@ -378,6 +378,19 @@ func (sm *SessionManager) msgPersistWorker() {
} }
} }
func (sm *SessionManager) enqueuePersistItem(item msgPersistItem) (ok bool) {
if sm.msgChan == nil {
return false
}
defer func() {
if recover() != nil {
ok = false
}
}()
sm.msgChan <- item
return true
}
func (sm *SessionManager) persistMessageToDelegate(sessionKey string, msg messages.Message) { func (sm *SessionManager) persistMessageToDelegate(sessionKey string, msg messages.Message) {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel() defer cancel()
@ -581,7 +594,9 @@ func (sm *SessionManager) Flush() {
return return
} }
barrier := make(chan struct{}) barrier := make(chan struct{})
sm.msgChan <- msgPersistItem{barrier: barrier} if !sm.enqueuePersistItem(msgPersistItem{barrier: barrier}) {
return
}
<-barrier <-barrier
} }