From ccb2d99d28940c8f6736a98c2b1810442112c663 Mon Sep 17 00:00:00 2001 From: ZanzyTHEbar Date: Thu, 5 Mar 2026 21:17:27 +0000 Subject: [PATCH] refactor(session): update manager --- pkg/session/manager.go | 19 +++++++++++++++++-- 1 file changed, 17 insertions(+), 2 deletions(-) diff --git a/pkg/session/manager.go b/pkg/session/manager.go index e2d9b08e5..945df2777 100644 --- a/pkg/session/manager.go +++ b/pkg/session/manager.go @@ -361,7 +361,7 @@ func (sm *SessionManager) AddFullMessage(sessionKey string, msg messages.Message sm.touchLRU(sessionKey) 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) { ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) defer cancel() @@ -581,7 +594,9 @@ func (sm *SessionManager) Flush() { return } barrier := make(chan struct{}) - sm.msgChan <- msgPersistItem{barrier: barrier} + if !sm.enqueuePersistItem(msgPersistItem{barrier: barrier}) { + return + } <-barrier }