Merge pull request #1390 from kiannidev/fix/1323-telegram-endless-typing
fix(telegram): stop typing indicator when LLM fails or hangs
This commit is contained in:
commit
bd4317f1f4
4 changed files with 113 additions and 47 deletions
|
|
@ -278,58 +278,64 @@ func (al *AgentLoop) Run(ctx context.Context) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
// Process message
|
// Process message
|
||||||
// TODO: Re-enable media cleanup after inbound media is properly consumed by the agent.
|
func() {
|
||||||
// Currently disabled because files are deleted before the LLM can access their content.
|
defer func() {
|
||||||
// defer func() {
|
if al.channelManager != nil {
|
||||||
// if al.mediaStore != nil && msg.MediaScope != "" {
|
al.channelManager.InvokeTypingStop(msg.Channel, msg.ChatID)
|
||||||
// if releaseErr := al.mediaStore.ReleaseAll(msg.MediaScope); releaseErr != nil {
|
}
|
||||||
// logger.WarnCF("agent", "Failed to release media", map[string]any{
|
}()
|
||||||
// "scope": msg.MediaScope,
|
// TODO: Re-enable media cleanup after inbound media is properly consumed by the agent.
|
||||||
// "error": releaseErr.Error(),
|
// Currently disabled because files are deleted before the LLM can access their content.
|
||||||
// })
|
// defer func() {
|
||||||
// }
|
// if al.mediaStore != nil && msg.MediaScope != "" {
|
||||||
// }
|
// if releaseErr := al.mediaStore.ReleaseAll(msg.MediaScope); releaseErr != nil {
|
||||||
// }()
|
// logger.WarnCF("agent", "Failed to release media", map[string]any{
|
||||||
|
// "scope": msg.MediaScope,
|
||||||
|
// "error": releaseErr.Error(),
|
||||||
|
// })
|
||||||
|
// }
|
||||||
|
// }
|
||||||
|
// }()
|
||||||
|
|
||||||
response, err := al.processMessage(ctx, msg)
|
response, err := al.processMessage(ctx, msg)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
response = fmt.Sprintf("Error processing message: %v", err)
|
response = fmt.Sprintf("Error processing message: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if response != "" {
|
if response != "" {
|
||||||
// Check if the message tool already sent a response during this round.
|
// Check if the message tool already sent a response during this round.
|
||||||
// If so, skip publishing to avoid duplicate messages to the user.
|
// If so, skip publishing to avoid duplicate messages to the user.
|
||||||
// Use default agent's tools to check (message tool is shared).
|
// Use default agent's tools to check (message tool is shared).
|
||||||
alreadySent := false
|
alreadySent := false
|
||||||
defaultAgent := al.GetRegistry().GetDefaultAgent()
|
defaultAgent := al.GetRegistry().GetDefaultAgent()
|
||||||
if defaultAgent != nil {
|
if defaultAgent != nil {
|
||||||
if tool, ok := defaultAgent.Tools.Get("message"); ok {
|
if tool, ok := defaultAgent.Tools.Get("message"); ok {
|
||||||
if mt, ok := tool.(*tools.MessageTool); ok {
|
if mt, ok := tool.(*tools.MessageTool); ok {
|
||||||
alreadySent = mt.HasSentInRound()
|
alreadySent = mt.HasSentInRound()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
if !alreadySent {
|
||||||
|
al.bus.PublishOutbound(ctx, bus.OutboundMessage{
|
||||||
if !alreadySent {
|
Channel: msg.Channel,
|
||||||
al.bus.PublishOutbound(ctx, bus.OutboundMessage{
|
ChatID: msg.ChatID,
|
||||||
Channel: msg.Channel,
|
Content: response,
|
||||||
ChatID: msg.ChatID,
|
|
||||||
Content: response,
|
|
||||||
})
|
|
||||||
logger.InfoCF("agent", "Published outbound response",
|
|
||||||
map[string]any{
|
|
||||||
"channel": msg.Channel,
|
|
||||||
"chat_id": msg.ChatID,
|
|
||||||
"content_len": len(response),
|
|
||||||
})
|
})
|
||||||
} else {
|
logger.InfoCF("agent", "Published outbound response",
|
||||||
logger.DebugCF(
|
map[string]any{
|
||||||
"agent",
|
"channel": msg.Channel,
|
||||||
"Skipped outbound (message tool already sent)",
|
"chat_id": msg.ChatID,
|
||||||
map[string]any{"channel": msg.Channel},
|
"content_len": len(response),
|
||||||
)
|
})
|
||||||
|
} else {
|
||||||
|
logger.DebugCF(
|
||||||
|
"agent",
|
||||||
|
"Skipped outbound (message tool already sent)",
|
||||||
|
map[string]any{"channel": msg.Channel},
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}()
|
||||||
default:
|
default:
|
||||||
time.Sleep(time.Microsecond * 200)
|
time.Sleep(time.Microsecond * 200)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -136,6 +136,19 @@ func (m *Manager) RecordTypingStop(channel, chatID string, stop func()) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// InvokeTypingStop invokes the registered typing stop function for the given channel and chatID.
|
||||||
|
// It is safe to call even when no typing indicator is active (no-op).
|
||||||
|
// Used by the agent loop to stop typing when processing completes (success, error, or panic),
|
||||||
|
// regardless of whether an outbound message is published.
|
||||||
|
func (m *Manager) InvokeTypingStop(channel, chatID string) {
|
||||||
|
key := channel + ":" + chatID
|
||||||
|
if v, loaded := m.typingStops.LoadAndDelete(key); loaded {
|
||||||
|
if entry, ok := v.(typingEntry); ok {
|
||||||
|
entry.stop()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// RecordReactionUndo registers a reaction undo function for later invocation.
|
// RecordReactionUndo registers a reaction undo function for later invocation.
|
||||||
// Implements PlaceholderRecorder.
|
// Implements PlaceholderRecorder.
|
||||||
func (m *Manager) RecordReactionUndo(channel, chatID string, undo func()) {
|
func (m *Manager) RecordReactionUndo(channel, chatID string, undo func()) {
|
||||||
|
|
|
||||||
|
|
@ -511,6 +511,43 @@ func TestPreSend_PlaceholderEditFails_FallsThrough(t *testing.T) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestInvokeTypingStop_CallsRegisteredStop(t *testing.T) {
|
||||||
|
m := newTestManager()
|
||||||
|
var stopCalled bool
|
||||||
|
|
||||||
|
m.RecordTypingStop("telegram", "chat123", func() {
|
||||||
|
stopCalled = true
|
||||||
|
})
|
||||||
|
|
||||||
|
m.InvokeTypingStop("telegram", "chat123")
|
||||||
|
|
||||||
|
if !stopCalled {
|
||||||
|
t.Fatal("expected typing stop func to be called")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestInvokeTypingStop_NoOpWhenNoEntry(t *testing.T) {
|
||||||
|
m := newTestManager()
|
||||||
|
// Should not panic
|
||||||
|
m.InvokeTypingStop("telegram", "nonexistent")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestInvokeTypingStop_Idempotent(t *testing.T) {
|
||||||
|
m := newTestManager()
|
||||||
|
var callCount int
|
||||||
|
|
||||||
|
m.RecordTypingStop("telegram", "chat123", func() {
|
||||||
|
callCount++
|
||||||
|
})
|
||||||
|
|
||||||
|
m.InvokeTypingStop("telegram", "chat123")
|
||||||
|
m.InvokeTypingStop("telegram", "chat123") // Second call: entry already removed, no-op
|
||||||
|
|
||||||
|
if callCount != 1 {
|
||||||
|
t.Fatalf("expected stop to be called once, got %d", callCount)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestPreSend_TypingStopCalled(t *testing.T) {
|
func TestPreSend_TypingStopCalled(t *testing.T) {
|
||||||
m := newTestManager()
|
m := newTestManager()
|
||||||
var stopCalled bool
|
var stopCalled bool
|
||||||
|
|
|
||||||
|
|
@ -302,10 +302,17 @@ func (c *TelegramChannel) sendChunk(
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// maxTypingDuration limits how long the typing indicator can run.
|
||||||
|
// Prevents endless typing when the LLM fails/hangs and preSend never invokes cancel.
|
||||||
|
// Matches channels.Manager's typingStopTTL (5 min) so behavior is consistent.
|
||||||
|
const maxTypingDuration = 5 * time.Minute
|
||||||
|
|
||||||
// StartTyping implements channels.TypingCapable.
|
// StartTyping implements channels.TypingCapable.
|
||||||
// It sends ChatAction(typing) immediately and then repeats every 4 seconds
|
// It sends ChatAction(typing) immediately and then repeats every 4 seconds
|
||||||
// (Telegram's typing indicator expires after ~5s) in a background goroutine.
|
// (Telegram's typing indicator expires after ~5s) in a background goroutine.
|
||||||
// The returned stop function is idempotent and cancels the goroutine.
|
// The returned stop function is idempotent and cancels the goroutine.
|
||||||
|
// The goroutine also exits automatically after maxTypingDuration if cancel is
|
||||||
|
// never called (e.g. when the LLM fails or times out without publishing).
|
||||||
func (c *TelegramChannel) StartTyping(ctx context.Context, chatID string) (func(), error) {
|
func (c *TelegramChannel) StartTyping(ctx context.Context, chatID string) (func(), error) {
|
||||||
cid, threadID, err := parseTelegramChatID(chatID)
|
cid, threadID, err := parseTelegramChatID(chatID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
@ -319,12 +326,15 @@ func (c *TelegramChannel) StartTyping(ctx context.Context, chatID string) (func(
|
||||||
_ = c.bot.SendChatAction(ctx, action)
|
_ = c.bot.SendChatAction(ctx, action)
|
||||||
|
|
||||||
typingCtx, cancel := context.WithCancel(ctx)
|
typingCtx, cancel := context.WithCancel(ctx)
|
||||||
|
// Cap lifetime so the goroutine cannot run indefinitely if cancel is never called
|
||||||
|
maxCtx, maxCancel := context.WithTimeout(typingCtx, maxTypingDuration)
|
||||||
go func() {
|
go func() {
|
||||||
|
defer maxCancel()
|
||||||
ticker := time.NewTicker(4 * time.Second)
|
ticker := time.NewTicker(4 * time.Second)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-typingCtx.Done():
|
case <-maxCtx.Done():
|
||||||
return
|
return
|
||||||
case <-ticker.C:
|
case <-ticker.C:
|
||||||
a := tu.ChatAction(tu.ID(cid), telego.ChatActionTyping)
|
a := tu.ChatAction(tu.ID(cid), telego.ChatActionTyping)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue