fix: prevent duplicate heartbeat completion messages on Telegram
The streaming status bubble (IsStatus) and task completion notification (IsTaskStatus) were creating two separate messages. Now the streaming bubble's message ID is promoted to taskMsgIDs so the completion notification edits the existing bubble instead of creating a new one. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
e5f92c10dd
commit
c52144d75a
2 changed files with 21 additions and 9 deletions
|
|
@ -1354,17 +1354,16 @@ func (al *AgentLoop) runAgentLoop(ctx context.Context, agent *AgentInstance, opt
|
|||
task.Result = utils.Truncate(finalContent, 280)
|
||||
}
|
||||
|
||||
// 5e. Replace orphaned streaming status bubble for background tasks.
|
||||
// 5e. Promote streaming status bubble to task message for background tasks.
|
||||
// When SendResponse is false (e.g. heartbeat), no final non-status message
|
||||
// triggers cleanup, so the last streaming chunk persists on Telegram.
|
||||
// triggers preSend's statusMsgIDs.LoadAndDelete cleanup, so the last
|
||||
// streaming chunk persists on Telegram. Move the tracked status message ID
|
||||
// into taskMsgIDs so the defer's IsTaskStatus completion notification edits
|
||||
// the existing bubble instead of creating a duplicate message.
|
||||
if opts.Background && !opts.SendResponse && !constants.IsInternalChannel(opts.Channel) && task != nil &&
|
||||
task.streamedChunks {
|
||||
_ = al.bus.PublishOutbound(ctx, bus.OutboundMessage{
|
||||
Channel: opts.Channel,
|
||||
ChatID: opts.ChatID,
|
||||
Content: utils.Truncate(finalContent, 200),
|
||||
IsStatus: true,
|
||||
})
|
||||
task.streamedChunks && opts.TaskID != "" {
|
||||
statusKey := opts.Channel + ":" + opts.ChatID
|
||||
al.channelManager.PromoteStatusToTask(statusKey, opts.TaskID)
|
||||
}
|
||||
|
||||
// 6. Save final assistant message to session (deferred write-behind)
|
||||
|
|
|
|||
|
|
@ -839,6 +839,19 @@ func (m *Manager) runTTLJanitor(ctx context.Context) {
|
|||
}
|
||||
}
|
||||
|
||||
// PromoteStatusToTask moves the tracked streaming status message for the given
|
||||
// channel:chatID key into the task message map under taskID. This allows the
|
||||
// next IsTaskStatus publish to edit the streaming bubble instead of creating a
|
||||
// new message. Returns true if a status message was found and promoted.
|
||||
func (m *Manager) PromoteStatusToTask(statusKey, taskID string) bool {
|
||||
v, loaded := m.statusMsgIDs.LoadAndDelete(statusKey)
|
||||
if !loaded {
|
||||
return false
|
||||
}
|
||||
m.taskMsgIDs.Store(taskID, v)
|
||||
return true
|
||||
}
|
||||
|
||||
func (m *Manager) GetChannel(name string) (Channel, bool) {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue