Files
DangTinh311 af43c952ed feat(bitrix24): agent activity indicator via InputAction.notify [B24:2794]
Show an ephemeral "agent is working" indicator (thinking/searching/
generating/analyzing…) in Bitrix24 chat while the agent processes, so
users on this non-streaming channel aren't left staring at silence
until the final reply. Send behavior is unchanged; no LLM call, no DB.

Generic layer (internal/channels):
- New optional ActivityIndicatorChannel interface.
- HandleAgentEvent routes run.started->THINKING, tool.call->mapped
  status, tool.result->ANALYZING; static tool->status mapping.
- Conditional heartbeat ticker fills LLM-inference gaps (re-sends only
  when idle), 5s per-run throttle caps call rate; ticker stopped on
  terminal events AND in UnregisterRun (safety net for missed terminals).

Bitrix24 (internal/channels/bitrix24):
- OnActivityEvent calls imbot.v2.Chat.InputAction.notify best-effort,
  drop-on-limit (raw Client.Call, no retry) so cosmetic notifies never
  steal leaky-bucket capacity from real message sends.
- Per-channel activity_indicator toggle (default on).

Web UI + docs: dashboard toggle, i18n en/vi/zh, channels doc section.
Tests: 26 unit tests incl. UnregisterRun-stops-ticker regression.
2026-07-14 15:28:14 +07:00

145 lines
5.2 KiB
Go

package channels
import (
"github.com/google/uuid"
"github.com/nextlevelbuilder/goclaw/internal/config"
)
// --- Run tracking for streaming/reaction event forwarding ---
// RegisterRun associates a run ID with a channel context so agent events
// (chunks, tool calls, completion) can be forwarded to the originating channel.
func (m *Manager) RegisterRun(runID, channelName, chatID, messageID string, metadata map[string]string, tenantID uuid.UUID, streaming, blockReply, toolStatus bool) {
m.RegisterRunWithBehavior(runID, channelName, chatID, messageID, metadata, tenantID, streaming, blockReply, toolStatus, ResolvedChatBehavior{})
}
// RegisterRunWithBehavior associates a run ID with channel context and
// resolved delivery behavior so event handlers do not read mutable config mid-run.
func (m *Manager) RegisterRunWithBehavior(runID, channelName, chatID, messageID string, metadata map[string]string, tenantID uuid.UUID, streaming, blockReply, toolStatus bool, chatBehavior ResolvedChatBehavior, reasoningDelivery ...ResolvedReasoningDelivery) {
m.RegisterRunWithDelivery(runID, channelName, chatID, messageID, metadata, tenantID, streaming, blockReply, toolStatus, chatBehavior, DeliveryRuntime{}, reasoningDelivery...)
}
func (m *Manager) RegisterRunWithDelivery(runID, channelName, chatID, messageID string, metadata map[string]string, tenantID uuid.UUID, streaming, blockReply, toolStatus bool, chatBehavior ResolvedChatBehavior, deliveryRuntime DeliveryRuntime, reasoningDelivery ...ResolvedReasoningDelivery) {
delivery := ResolveReasoningDelivery("", nil)
if len(reasoningDelivery) > 0 {
delivery = reasoningDelivery[0]
}
m.runs.Store(runID, &RunContext{
ChannelName: channelName,
ChatID: chatID,
MessageID: messageID,
Metadata: metadata,
TenantID: tenantID,
Streaming: streaming,
BlockReplyEnabled: blockReply,
ToolStatusEnabled: toolStatus,
ChatBehavior: chatBehavior,
Delivery: deliveryRuntime,
ReasoningDelivery: delivery,
})
}
// UnregisterRun removes a run tracking entry.
func (m *Manager) UnregisterRun(runID string) {
if val, ok := m.runs.LoadAndDelete(runID); ok {
if rc, ok := val.(*RunContext); ok {
m.cancelQuickAck(rc)
// Also stop the activity heartbeat here (not only in the terminal-event
// branch) — terminal AgentEvents are not guaranteed to reach
// HandleAgentEvent, and without this the heartbeat goroutine leaks and
// keeps firing InputAction.notify forever. stopActivityTicker is idempotent.
m.stopActivityTicker(rc)
m.flushReasoningBubblesForContext(rc)
}
}
}
func (m *Manager) InterimDeliverySnapshot(runID string) (int, string) {
val, ok := m.runs.Load(runID)
if ok {
rc, ok := val.(*RunContext)
if !ok {
return 0, ""
}
rc.mu.Lock()
defer rc.mu.Unlock()
return rc.interimDelivered, rc.lastInterimReply
}
return 0, ""
}
// IsStreamingChannel checks if a named channel implements StreamingChannel
// AND has streaming currently enabled for the given chat type.
// isGroup: true for group chats, false for DMs.
func (m *Manager) IsStreamingChannel(channelName string, isGroup bool) bool {
m.mu.RLock()
ch, exists := m.channels[channelName]
m.mu.RUnlock()
if !exists {
return false
}
sc, ok := ch.(StreamingChannel)
if !ok {
return false
}
return sc.StreamEnabled(isGroup)
}
// ResolveBlockReply checks per-channel override, falls back to gateway default.
// Returns true only if block.reply delivery should be enabled for this channel.
func (m *Manager) ResolveBlockReply(channelName string, globalDefault *bool) bool {
m.mu.RLock()
ch, exists := m.channels[channelName]
m.mu.RUnlock()
if exists {
if bc, ok := ch.(BlockReplyChannel); ok {
if v := bc.BlockReplyEnabled(); v != nil {
return *v
}
}
}
return globalDefault != nil && *globalDefault
}
// ResolveChatBehavior checks per-channel override, then falls back to gateway config.
func (m *Manager) ResolveChatBehavior(channelName string, globalDefault *config.ChatBehaviorConfig) ResolvedChatBehavior {
return m.ResolveChatBehaviorWithAgent(channelName, globalDefault, nil)
}
func (m *Manager) ResolveChatBehaviorWithAgent(channelName string, globalDefault, agentOverride *config.ChatBehaviorConfig) ResolvedChatBehavior {
var override *config.ChatBehaviorConfig
var channelBlockReply *bool
m.mu.RLock()
ch, exists := m.channels[channelName]
m.mu.RUnlock()
if exists {
if bc, ok := ch.(BlockReplyChannel); ok {
channelBlockReply = bc.BlockReplyEnabled()
}
if bc, ok := ch.(ChatBehaviorChannel); ok {
override = bc.ChatBehaviorConfig()
}
}
override = ChatBehaviorConfigWithIntermediateDefault(override, channelBlockReply)
return ResolveChatBehaviorWithAgent(globalDefault, agentOverride, override)
}
func (m *Manager) ResolveReasoningDelivery(channelName string) ResolvedReasoningDelivery {
m.mu.RLock()
ch, exists := m.channels[channelName]
m.mu.RUnlock()
if !exists {
return ResolveReasoningDelivery("", nil)
}
if rc, ok := ch.(ReasoningDeliveryChannel); ok {
mode, legacy := rc.ReasoningDeliveryConfig()
return ResolveReasoningDelivery(mode, legacy)
}
if sc, ok := ch.(StreamingChannel); ok {
legacy := sc.ReasoningStreamEnabled()
return ResolveReasoningDelivery("", &legacy)
}
return ResolveReasoningDelivery("", nil)
}