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

846 lines
32 KiB
Go

// Package channels provides the channel abstraction layer for multi-platform messaging.
// Channels connect external platforms (Telegram, Discord, Slack, etc.) to the agent runtime
// via the message bus.
//
// Adapted from PicoClaw's pkg/channels with GoClaw-specific additions:
// - DM/Group policies (pairing, allowlist, open, disabled)
// - Mention gating for group chats
// - Rich MsgContext metadata
package channels
import (
"context"
"log/slog"
"net/http"
"path/filepath"
"strings"
"sync"
"time"
"github.com/google/uuid"
"github.com/nextlevelbuilder/goclaw/internal/bus"
"github.com/nextlevelbuilder/goclaw/internal/config"
"github.com/nextlevelbuilder/goclaw/internal/store"
"github.com/nextlevelbuilder/goclaw/internal/systemmessages"
)
// PolicyResult is returned by BaseChannel policy checks.
type PolicyResult int
const (
// PolicyAllow means the message should be processed.
PolicyAllow PolicyResult = iota
// PolicyDeny means the message should be silently dropped.
PolicyDeny
// PolicyNeedsPairing means the sender is unpaired; caller should send its platform-specific pairing reply.
PolicyNeedsPairing
)
// InternalChannels are system channels excluded from outbound dispatch.
// "browser" uses WebSocket directly — no outbound channel routing needed.
var InternalChannels = map[string]bool{
"cli": true,
"system": true,
"subagent": true,
"browser": true,
"ws": true, // WebSocket — responses delivered via events/RPC, not outbound dispatch
}
// IsInternalChannel checks if a channel name is internal.
func IsInternalChannel(name string) bool {
return InternalChannels[name]
}
// DMPolicy controls how DMs from unknown senders are handled.
type DMPolicy string
const (
DMPolicyPairing DMPolicy = "pairing" // Require pairing code
DMPolicyAllowlist DMPolicy = "allowlist" // Only whitelisted senders
DMPolicyOpen DMPolicy = "open" // Accept all
DMPolicyDisabled DMPolicy = "disabled" // Reject all DMs
)
// GroupPolicy controls how group messages are handled.
type GroupPolicy string
const (
GroupPolicyOpen GroupPolicy = "open" // Accept all groups
GroupPolicyAllowlist GroupPolicy = "allowlist" // Only whitelisted groups
GroupPolicyDisabled GroupPolicy = "disabled" // No group messages
)
// Channel type constants used across channel packages and gateway wiring.
const (
TypeBitrix24 = "bitrix24"
TypeDiscord = "discord"
TypeFacebook = "facebook"
TypeFeishu = "feishu"
TypePancake = "pancake"
TypeSlack = "slack"
TypeTelegram = "telegram"
TypeWhatsApp = "whatsapp"
TypeZaloOA = "zalo_oa"
TypeZaloPersonal = "zalo_personal"
)
// Channel defines the interface that all channel implementations must satisfy.
type Channel interface {
// Name returns the channel instance name (e.g., "telegram", "discord", "slack").
Name() string
// Type returns the platform type (e.g., "telegram", "zalo_personal").
// For config-based channels this equals Name(); for DB instances it may differ.
Type() string
// Start begins listening for messages. Should be non-blocking after setup.
Start(ctx context.Context) error
// Stop gracefully shuts down the channel.
Stop(ctx context.Context) error
// Send delivers an outbound message to the channel.
Send(ctx context.Context, msg bus.OutboundMessage) error
// IsRunning returns whether the channel is actively processing messages.
IsRunning() bool
// IsAllowed checks if a sender is permitted by the channel's allowlist.
IsAllowed(senderID string) bool
}
// StreamingChannel extends Channel with real-time streaming preview support.
// Channels that implement this interface can show incremental response updates
// (e.g., editing a Telegram message as chunks arrive) instead of waiting for the full response.
type StreamingChannel interface {
Channel
// StreamEnabled reports whether the channel currently wants LLM streaming.
// When false the agent loop uses non-streaming Chat() instead of ChatStream(),
// which gives more accurate token usage from providers that don't support
// stream_options (e.g. MiniMax). The channel still implements the interface
// so it can be toggled at runtime via config.
//
// isGroup indicates whether this is a group chat (true) or DM (false).
// Channels may choose to always stream for DMs while gating group streaming
// behind config (e.g. Telegram uses sendMessageDraft for DMs).
StreamEnabled(isGroup bool) bool
// CreateStream creates a new per-run streaming handle for the given chatID.
// The returned ChannelStream is stored on RunContext so each concurrent run
// gets its own stream — eliminates the chatID-keyed sync.Map collision bug.
// firstStream: true for the first stream in a run (may become reasoning lane —
// must use message transport so it persists as a real message). false for
// subsequent streams (answer lane — may use draft transport for stealth preview).
CreateStream(ctx context.Context, chatID string, firstStream bool) (ChannelStream, error)
// FinalizeStream is called after the stream has been stopped to hand off
// the stream's messageID (if any) back to the channel's placeholder map
// so that Send() can edit it with the final formatted response.
FinalizeStream(ctx context.Context, chatID string, stream ChannelStream)
// ReasoningStreamEnabled returns whether reasoning should be shown as a
// separate message. Default: true. Channels that don't support lanes can
// return false to skip reasoning routing.
ReasoningStreamEnabled() bool
}
// ChannelDestroyer extends Channel with a deletion hook. Channels that
// implement this are notified BEFORE their channel_instance row is removed
// from the DB so they can release external resources that won't survive a
// normal Stop() — e.g. Bitrix24 channels call imbot.unregister to delete
// the bot on the portal; without this hook the bot lingers as a zombie.
//
// Implementation must be best-effort: handlers log Destroy failures and
// proceed with DB deletion regardless. Blocking the delete on a permanently
// dead upstream would leave the row stuck forever with no recovery path.
//
// Channels without external state (Telegram, Discord, Slack — the channel
// itself IS the bot, identified by a token stored locally) don't implement
// this interface; their Stop() already handles all cleanup.
type ChannelDestroyer interface {
Channel
Destroy(ctx context.Context) error
}
// BlockReplyChannel is optionally implemented by channels that override
// the gateway-level block_reply setting. Returns nil to inherit the gateway default.
type BlockReplyChannel interface {
BlockReplyEnabled() *bool
}
// ChatBehaviorChannel is optionally implemented by channels that override
// gateway-level human-like delivery behavior. Nil means inherit the gateway default.
type ChatBehaviorChannel interface {
ChatBehaviorConfig() *config.ChatBehaviorConfig
}
// ReasoningDeliveryChannel is optionally implemented by channels that expose
// how model reasoning should be delivered to end users.
type ReasoningDeliveryChannel interface {
ReasoningDeliveryConfig() (mode string, legacyReasoningStream *bool)
}
// WebhookChannel extends Channel with an HTTP handler that can be mounted
// on the main gateway mux instead of starting a separate HTTP server.
// This allows webhook-based channels (e.g. Feishu/Lark) to share the main
// server port, avoiding the need to expose additional ports in Docker.
type WebhookChannel interface {
Channel
// WebhookHandler returns the HTTP handler and the path it should be mounted on.
// Returns ("", nil) if the channel doesn't use webhook mode.
WebhookHandler() (path string, handler http.Handler)
}
// ReactionChannel extends Channel with status reaction support.
// Channels that implement this interface can show emoji reactions on user messages
// to indicate agent status (thinking, tool call, done, error, stall).
// messageID is a string to support platforms with non-integer IDs (e.g., Feishu "om_xxx").
type ReactionChannel interface {
Channel
OnReactionEvent(ctx context.Context, chatID string, messageID string, status string) error
ClearReaction(ctx context.Context, chatID string, messageID string) error
}
// ActivityIndicatorChannel is optionally implemented by channels that can show an
// ephemeral "agent is working" indicator (e.g. Bitrix24 imbot InputAction.notify)
// while the agent thinks or runs tools — without persisting a chat message.
//
// statusCode is a platform-native activity code (see activity_indicator.go); an empty
// string means "keep the current / default working indicator". Implementations MUST be
// best-effort and non-blocking on the caller's side: the indicator is cosmetic, so a
// failed or rate-limited call must be dropped silently and never affect the real reply.
type ActivityIndicatorChannel interface {
Channel
OnActivityEvent(ctx context.Context, chatID string, statusCode string) error
}
// BaseChannel provides shared functionality for all channel implementations.
// Channel implementations should embed this struct.
type BaseChannel struct {
name string
channelType string // platform type; defaults to name if unset
bus *bus.MessageBus
running bool
stateMu sync.RWMutex
health ChannelHealth
allowList []string
agentID string // for DB instances: routes to specific agent (empty = use resolveAgentRoute)
tenantID uuid.UUID // for DB instances: tenant scope (zero = master tenant fallback)
contactCollector *store.ContactCollector // optional: auto-collect contacts from channel messages
// Shared policy + pairing fields (set via setters after construction).
pairingService store.PairingStore
systemMessages *systemmessages.Resolver
groupHistory *PendingHistory
historyLimit int
approvedGroups sync.Map // chatID → true (in-memory cache for paired group approval)
pairingDebounce sync.Map // senderID → time.Time (debounce pairing reply sends)
requireMention bool
}
// NewBaseChannel creates a new BaseChannel with the given parameters.
func NewBaseChannel(name string, msgBus *bus.MessageBus, allowList []string) *BaseChannel {
return &BaseChannel{
name: name,
bus: msgBus,
health: NewChannelHealthForType(name, ChannelHealthStateRegistered, "Configured", "", ChannelFailureKindUnknown, false),
allowList: allowList,
}
}
// Name returns the channel instance name.
func (c *BaseChannel) Name() string { return c.name }
// Type returns the platform type. Falls back to name if unset (config-based channels).
func (c *BaseChannel) Type() string {
if c.channelType != "" {
return c.channelType
}
return c.name
}
// SetName overrides the channel name (used by InstanceLoader for DB instances).
func (c *BaseChannel) SetName(name string) {
c.name = name
if c.groupHistory != nil {
c.groupHistory.SetChannelName(name)
}
}
// SetType sets the platform type (used by InstanceLoader for DB instances).
func (c *BaseChannel) SetType(t string) { c.channelType = t }
// AgentID returns the explicit agent ID for this channel (empty = use resolveAgentRoute).
func (c *BaseChannel) AgentID() string { return c.agentID }
// SetAgentID sets the explicit agent ID for routing (used by InstanceLoader for DB instances).
func (c *BaseChannel) SetAgentID(id string) { c.agentID = id }
// TenantID returns the tenant UUID for this channel (zero = master tenant fallback).
func (c *BaseChannel) TenantID() uuid.UUID { return c.tenantID }
// SetTenantID sets the tenant scope (used by InstanceLoader for DB instances).
func (c *BaseChannel) SetTenantID(id uuid.UUID) { c.tenantID = id }
// SetContactCollector sets the contact collector for auto-collecting contacts from messages.
func (c *BaseChannel) SetContactCollector(cc *store.ContactCollector) { c.contactCollector = cc }
// ContactCollector returns the contact collector (may be nil).
func (c *BaseChannel) ContactCollector() *store.ContactCollector { return c.contactCollector }
// SetPairingService sets the pairing store used for policy checks and code generation.
func (c *BaseChannel) SetPairingService(ps store.PairingStore) { c.pairingService = ps }
// PairingService returns the configured pairing store (may be nil).
func (c *BaseChannel) PairingService() store.PairingStore { return c.pairingService }
// SetSystemMessages sets the resolver used for operator/system messages.
func (c *BaseChannel) SetSystemMessages(r *systemmessages.Resolver) { c.systemMessages = r }
// SystemMessage renders a configurable operator/system message.
func (c *BaseChannel) SystemMessage(locale, key string, vars systemmessages.Vars) string {
if c.systemMessages != nil {
return c.systemMessages.Render(locale, key, vars)
}
return systemmessages.Render(locale, key, vars)
}
// SetGroupHistory sets the pending group history tracker.
func (c *BaseChannel) SetGroupHistory(gh *PendingHistory) {
c.groupHistory = gh
if c.groupHistory != nil {
c.groupHistory.SetChannelName(c.name)
}
}
// GroupHistory returns the pending group history tracker (may be nil).
func (c *BaseChannel) GroupHistory() *PendingHistory { return c.groupHistory }
// SetHistoryLimit sets the per-group message accumulation limit.
func (c *BaseChannel) SetHistoryLimit(n int) { c.historyLimit = n }
// HistoryLimit returns the per-group message accumulation limit.
func (c *BaseChannel) HistoryLimit() int { return c.historyLimit }
// SetRequireMention sets whether @mention is required in group chats.
func (c *BaseChannel) SetRequireMention(b bool) { c.requireMention = b }
// RequireMention returns whether @mention is required in group chats.
func (c *BaseChannel) RequireMention() bool { return c.requireMention }
// IsGroupApproved returns true if the group was already approved via pairing.
func (c *BaseChannel) IsGroupApproved(chatID string) bool {
_, ok := c.approvedGroups.Load(chatID)
return ok
}
// MarkGroupApproved caches a group as approved so future messages skip DB lookups.
func (c *BaseChannel) MarkGroupApproved(chatID string) {
c.approvedGroups.Store(chatID, true)
}
// ClearGroupApproval removes a group from the approval cache.
func (c *BaseChannel) ClearGroupApproval(chatID string) {
c.approvedGroups.Delete(chatID)
}
// CanSendPairingNotif returns true if debounce period has elapsed for senderID.
// debounce is the minimum interval between pairing replies to the same sender.
func (c *BaseChannel) CanSendPairingNotif(senderID string, debounce time.Duration) bool {
if lastSent, ok := c.pairingDebounce.Load(senderID); ok {
if time.Since(lastSent.(time.Time)) < debounce {
return false
}
}
return true
}
// MarkPairingNotifSent records the current time for senderID debounce tracking.
func (c *BaseChannel) MarkPairingNotifSent(senderID string) {
c.pairingDebounce.Store(senderID, time.Now())
}
// ClearPairingDebounce removes the debounce entry for a sender, allowing immediate pairing reply.
func (c *BaseChannel) ClearPairingDebounce(senderID string) {
c.pairingDebounce.Delete(senderID)
}
// CheckDMPolicy evaluates the DM policy for senderID.
// dmPolicy is one of: "pairing" (default), "open", "allowlist", "disabled".
// Returns PolicyAllow, PolicyDeny, or PolicyNeedsPairing.
// When PolicyNeedsPairing is returned, the caller should use its platform-specific
// pairing reply mechanism (BaseChannel has no knowledge of transport).
func (c *BaseChannel) CheckDMPolicy(ctx context.Context, senderID, dmPolicy string) PolicyResult {
if dmPolicy == "" {
dmPolicy = "pairing"
}
switch dmPolicy {
case "disabled":
return PolicyDeny
case "open":
return PolicyAllow
case "allowlist":
if c.IsAllowed(senderID) {
return PolicyAllow
}
return PolicyDeny
default: // "pairing"
if c.HasAllowList() && c.IsAllowed(senderID) {
return PolicyAllow
}
if c.pairingService != nil {
paired, err := c.pairingService.IsPaired(ctx, senderID, c.name)
if err != nil {
slog.Warn("security.pairing_check_failed, denying access (fail-closed)",
"sender_id", senderID, "channel", c.name, "error", err)
return PolicyDeny
}
if paired {
return PolicyAllow
}
}
return PolicyNeedsPairing
}
}
// CheckGroupPolicy evaluates the group policy for a message.
// groupPolicy is one of: "open" (default), "allowlist", "disabled", "pairing".
// chatID is the group chat identifier used for approval caching.
// Returns PolicyAllow, PolicyDeny, or PolicyNeedsPairing.
func (c *BaseChannel) CheckGroupPolicy(ctx context.Context, senderID, chatID, groupPolicy string) PolicyResult {
if groupPolicy == "" {
groupPolicy = "open"
}
switch groupPolicy {
case "disabled":
return PolicyDeny
case "allowlist":
if c.IsAllowed(senderID) {
return PolicyAllow
}
return PolicyDeny
case "pairing":
if c.HasAllowList() && c.IsAllowed(senderID) {
return PolicyAllow
}
if c.IsGroupApproved(chatID) {
return PolicyAllow
}
groupSenderID := "group:" + chatID
if c.pairingService != nil {
paired, err := c.pairingService.IsPaired(ctx, groupSenderID, c.name)
if err != nil {
slog.Warn("security.pairing_check_failed, denying access (fail-closed)",
"group_sender", groupSenderID, "channel", c.name, "error", err)
return PolicyDeny
}
if paired {
c.MarkGroupApproved(chatID)
return PolicyAllow
}
}
return PolicyNeedsPairing
default: // "open"
return PolicyAllow
}
}
// IsRunning returns whether the channel is running.
func (c *BaseChannel) IsRunning() bool {
c.stateMu.RLock()
defer c.stateMu.RUnlock()
return c.running
}
// SetRunning updates the running state.
func (c *BaseChannel) SetRunning(running bool) {
c.stateMu.Lock()
defer c.stateMu.Unlock()
next := c.health
next.ChannelType = c.Type()
next.Running = running
switch {
case running && (next.State == "" ||
next.State == ChannelHealthStateRegistered ||
next.State == ChannelHealthStateStarting ||
next.State == ChannelHealthStateStopped):
next.State = ChannelHealthStateHealthy
if next.Summary == "" ||
next.Summary == "Configured" ||
next.Summary == "Starting" ||
next.Summary == "Stopped" {
next.Summary = "Connected"
}
next.Detail = ""
next.FailureKind = ChannelFailureKindUnknown
next.Retryable = false
next.CheckedAt = time.Now().UTC()
case !running && next.State == ChannelHealthStateHealthy:
next.State = ChannelHealthStateStopped
next.Summary = "Stopped"
next.Detail = ""
next.FailureKind = ChannelFailureKindUnknown
next.Retryable = false
next.CheckedAt = time.Now().UTC()
default:
c.running = running
c.health.Running = running
return
}
next = mergeChannelHealth(c.health, next)
next.Running = running
c.running = running
c.health = next
}
// HealthSnapshot returns the current runtime health snapshot for the channel.
func (c *BaseChannel) HealthSnapshot() ChannelHealth {
c.stateMu.RLock()
defer c.stateMu.RUnlock()
snapshot := c.health
snapshot.ChannelType = c.Type()
snapshot.Enabled = true
snapshot.Running = c.running
return snapshot
}
// MarkRegistered records that the channel was configured and registered successfully.
func (c *BaseChannel) MarkRegistered(summary string) {
if summary == "" {
summary = "Configured"
}
c.setHealth(NewChannelHealth(ChannelHealthStateRegistered, summary, "", ChannelFailureKindUnknown, false))
}
// MarkStarting records that the channel is in startup validation / connection setup.
func (c *BaseChannel) MarkStarting(summary string) {
if summary == "" {
summary = "Starting"
}
c.setHealth(NewChannelHealth(ChannelHealthStateStarting, summary, "", ChannelFailureKindUnknown, true))
}
// MarkHealthy records a healthy connected state.
func (c *BaseChannel) MarkHealthy(summary string) {
if summary == "" {
summary = "Connected"
}
c.setHealth(NewChannelHealth(ChannelHealthStateHealthy, summary, "", ChannelFailureKindUnknown, false))
}
// MarkDegraded records a non-fatal warning state.
func (c *BaseChannel) MarkDegraded(summary, detail string, kind ChannelFailureKind, retryable bool) {
if summary == "" {
summary = "Running with warnings"
}
c.setHealth(NewChannelHealth(ChannelHealthStateDegraded, summary, detail, kind, retryable))
}
// MarkFailed records a startup or runtime failure.
func (c *BaseChannel) MarkFailed(summary, detail string, kind ChannelFailureKind, retryable bool) {
if summary == "" {
summary = "Channel failed"
}
c.setHealth(NewChannelHealth(ChannelHealthStateFailed, summary, detail, kind, retryable))
}
// MarkStopped records a cleanly stopped state.
func (c *BaseChannel) MarkStopped(summary string) {
if summary == "" {
summary = "Stopped"
}
c.setHealth(NewChannelHealth(ChannelHealthStateStopped, summary, "", ChannelFailureKindUnknown, false))
}
func (c *BaseChannel) setHealth(snapshot ChannelHealth) {
c.stateMu.Lock()
defer c.stateMu.Unlock()
snapshot.ChannelType = c.Type()
snapshot = mergeChannelHealth(c.health, snapshot)
snapshot.Running = snapshot.State == ChannelHealthStateHealthy || snapshot.State == ChannelHealthStateDegraded
c.running = snapshot.Running
c.health = snapshot
}
// Bus returns the message bus reference.
func (c *BaseChannel) Bus() *bus.MessageBus { return c.bus }
// HasAllowList returns true if an allowlist is configured (non-empty).
func (c *BaseChannel) HasAllowList() bool { return len(c.allowList) > 0 }
// IsAllowed checks if a sender is permitted by the allowlist.
// Supports compound senderID format: "123456|username".
// Empty allowlist means all senders are allowed.
func (c *BaseChannel) IsAllowed(senderID string) bool {
if len(c.allowList) == 0 {
return true
}
// Extract parts from compound senderID like "123456|username"
idPart := senderID
userPart := ""
if idx := strings.Index(senderID, "|"); idx > 0 {
idPart = senderID[:idx]
userPart = senderID[idx+1:]
}
for _, allowed := range c.allowList {
// Strip leading "@" from allowed value for username matching
trimmed := strings.TrimPrefix(allowed, "@")
allowedID := trimmed
allowedUser := ""
if idx := strings.Index(trimmed, "|"); idx > 0 {
allowedID = trimmed[:idx]
allowedUser = trimmed[idx+1:]
}
// Support either side using "id|username" compound form.
if senderID == allowed ||
idPart == allowed ||
senderID == trimmed ||
idPart == trimmed ||
idPart == allowedID ||
(allowedUser != "" && senderID == allowedUser) ||
(userPart != "" && (userPart == allowed || userPart == trimmed || userPart == allowedUser)) {
return true
}
}
return false
}
// CheckPolicy evaluates DM/Group policy for a message.
// Returns true if the message should be accepted, false if rejected.
// peerKind is "direct" or "group".
// dmPolicy/groupPolicy: "open" (default), "allowlist", "disabled".
func (c *BaseChannel) CheckPolicy(peerKind, dmPolicy, groupPolicy, senderID string) bool {
policy := dmPolicy
if peerKind == "group" {
policy = groupPolicy
}
if policy == "" {
policy = "open" // default for non-Telegram channels
}
switch policy {
case "disabled":
return false
case "allowlist":
return c.IsAllowed(senderID)
case "pairing":
// Channels with pairing handle this before CheckPolicy.
// If we reach here, no pairing service → still allow if in allowlist.
return c.IsAllowed(senderID)
default: // "open"
return true
}
}
// ValidatePolicy logs warnings for common policy misconfigurations.
// Should be called during channel initialization.
func (c *BaseChannel) ValidatePolicy(dmPolicy, groupPolicy string) {
if dmPolicy == "allowlist" && !c.HasAllowList() {
slog.Warn("channel policy misconfiguration: dmPolicy=allowlist but allowFrom is empty — all DMs will be rejected",
"channel", c.name)
}
if groupPolicy == "allowlist" && !c.HasAllowList() {
slog.Warn("channel policy misconfiguration: groupPolicy=allowlist but allowFrom is empty — all group messages will be rejected",
"channel", c.name)
}
}
// HandleMessage creates an InboundMessage and publishes it to the bus.
// This is the standard way for channels to forward received messages.
// peerKind should be "direct" or "group" (see sessions.PeerDirect, sessions.PeerGroup).
func (c *BaseChannel) HandleMessage(senderID, chatID, content string, media []string, metadata map[string]string, peerKind string) {
c.handleMessage(senderID, chatID, content, media, metadata, peerKind, false)
}
// HandleAuthorizedMessage publishes a message after the caller has already
// enforced the channel policy. It preserves the default direct-message safety
// net for adapters that do not have an explicit policy gate.
func (c *BaseChannel) HandleAuthorizedMessage(senderID, chatID, content string, media []string, metadata map[string]string, peerKind string) {
c.handleMessage(senderID, chatID, content, media, metadata, peerKind, true)
}
func (c *BaseChannel) handleMessage(senderID, chatID, content string, media []string, metadata map[string]string, peerKind string, policyChecked bool) {
// Convert string paths to MediaFile (legacy path-only callers).
// Use filepath.Base(p) as filename so persistMedia's sanitizer gets a
// meaningful stem instead of falling back to UUID. MimeType is left empty —
// persistMedia infers it from the file extension for these callers.
var mediaFiles []bus.MediaFile
for _, p := range media {
mediaFiles = append(mediaFiles, bus.MediaFile{Path: p, Filename: filepath.Base(p)})
}
c.handleMessageMedia(senderID, chatID, content, mediaFiles, metadata, peerKind, policyChecked)
}
// HandleMessageMedia is the richer sibling of HandleMessage: it accepts
// pre-built bus.MediaFile values so a channel can preserve the original MIME
// type and filename. This matters because the agent pipeline routes media by
// MIME (image vs document vs audio vs video) — the path-only HandleMessage
// loses that information. Channels that already know the content type at
// download time (e.g. Bitrix24 file events) should call this directly.
func (c *BaseChannel) HandleMessageMedia(senderID, chatID, content string, media []bus.MediaFile, metadata map[string]string, peerKind string) {
c.handleMessageMedia(senderID, chatID, content, media, metadata, peerKind, false)
}
// HandleAuthorizedMessageMedia is the media-preserving variant for callers
// that have already enforced the channel policy.
func (c *BaseChannel) HandleAuthorizedMessageMedia(senderID, chatID, content string, media []bus.MediaFile, metadata map[string]string, peerKind string) {
c.handleMessageMedia(senderID, chatID, content, media, metadata, peerKind, true)
}
func (c *BaseChannel) handleMessageMedia(senderID, chatID, content string, media []bus.MediaFile, metadata map[string]string, peerKind string, policyChecked bool) {
// For DMs, enforce the allowlist as a safety net.
// For group messages, skip this check — group access is already enforced
// by the channel-specific group policy (checkGroupPolicy / CheckPolicy).
// Re-checking the sender here would incorrectly block users who are not
// individually listed but are in an allowed (or open-policy) group.
if peerKind != "group" && !policyChecked && !c.IsAllowed(senderID) {
return
}
// Derive userID from senderID: strip "|username" suffix if present (legacy Slack compound format).
// All channels now pass plain senderID; kept for backward compat with stored compound IDs.
userID := senderID
if idx := strings.IndexByte(senderID, '|'); idx > 0 {
userID = senderID[:idx]
}
msg := bus.InboundMessage{
Channel: c.name,
SenderID: senderID,
ChatID: chatID,
Content: content,
Media: media,
PeerKind: peerKind,
UserID: userID,
Metadata: metadata,
TenantID: c.tenantID,
AgentID: c.agentID,
}
c.bus.PublishInbound(msg)
}
// GroupMember represents a member of a group chat.
type GroupMember struct {
MemberID string `json:"member_id"`
Name string `json:"name"`
}
// GroupMemberProvider is optionally implemented by channels that can list group members.
type GroupMemberProvider interface {
ListGroupMembers(ctx context.Context, chatID string) ([]GroupMember, error)
}
// GroupInfo represents a group/channel chat the account belongs to.
type GroupInfo struct {
GroupID string `json:"group_id"`
Name string `json:"name"`
TotalMember int `json:"total_member,omitempty"`
}
// GroupListProvider is optionally implemented by channels that can list the
// groups/chats the connected account belongs to — lets the agent resolve a
// group's display name to its real chat ID instead of guessing.
type GroupListProvider interface {
ListGroups(ctx context.Context) ([]GroupInfo, error)
}
// GroupTitleProvider is optionally implemented by channels that can resolve
// a platform group/channel ID to a human-readable title.
type GroupTitleProvider interface {
ResolveGroupTitle(ctx context.Context, chatID string) (string, error)
}
// MetadataRefreshFailure records one group whose presentation metadata could
// not be refreshed.
type MetadataRefreshFailure struct {
ChannelID string `json:"channel_id"`
Source string `json:"source"`
Reason string `json:"reason"`
}
// MetadataRefreshReport describes a manual channel metadata refresh. It keeps
// stable IDs intact while making refresh coverage and failures observable.
type MetadataRefreshReport struct {
OK bool `json:"ok"`
GroupsRefreshed int `json:"groups_refreshed"`
UsersRefreshed int `json:"users_refreshed"`
ContactTargets int `json:"contact_targets"`
PendingMessageTargets int `json:"pending_message_targets"`
LiveTargets int `json:"live_targets"`
DirectLookupAttempts int `json:"direct_lookup_attempts"`
DirectLookupResolved int `json:"direct_lookup_resolved"`
Errors []string `json:"errors,omitempty"`
Failures []MetadataRefreshFailure `json:"failures,omitempty"`
}
// GroupTitlesProvider is optionally implemented by channels that can resolve
// multiple platform group/channel IDs in one best-effort operation.
type GroupTitlesProvider interface {
ResolveGroupTitles(ctx context.Context, chatIDs []string) (map[string]string, error)
}
// GroupDisplayTitleProvider optionally resolves a presentation title for a
// group/channel ID. Unlike GroupTitleProvider, the result may include a
// platform hierarchy (for example, a Discord thread and its parent channel).
// Routing and persistence must continue to use the stable chat ID.
type GroupDisplayTitleProvider interface {
ResolveGroupDisplayTitle(ctx context.Context, chatID string) (string, error)
}
// TelegramManagerRequest describes a whitelisted Telegram Bot API management
// action requested by the agent-facing telegram_manager tool.
type TelegramManagerRequest struct {
Action string `json:"action"`
ChatID string `json:"chat_id,omitempty"`
MessageThreadID int `json:"message_thread_id,omitempty"`
MessageID int `json:"message_id,omitempty"`
UserID int64 `json:"user_id,omitempty"`
Name string `json:"name,omitempty"`
Text string `json:"text,omitempty"`
InviteLink string `json:"invite_link,omitempty"`
IconColor int `json:"icon_color,omitempty"`
IconCustomEmojiID string `json:"icon_custom_emoji_id,omitempty"`
ExpireDate int64 `json:"expire_date,omitempty"`
MemberLimit int `json:"member_limit,omitempty"`
DisableNotification bool `json:"disable_notification,omitempty"`
CreatesJoinRequest bool `json:"creates_join_request,omitempty"`
OnlyIfBanned bool `json:"only_if_banned,omitempty"`
RevokeMessages bool `json:"revoke_messages,omitempty"`
Params map[string]any `json:"params,omitempty"`
}
// TelegramManagerResult is a normalized response returned by Telegram-capable
// channels after executing a whitelisted management action.
type TelegramManagerResult struct {
Action string `json:"action"`
Result map[string]any `json:"result,omitempty"`
}
// TelegramManagerProvider is optionally implemented by channels that expose
// Telegram Bot API management actions through the channel manager.
type TelegramManagerProvider interface {
ManageTelegram(ctx context.Context, req TelegramManagerRequest) (TelegramManagerResult, error)
}
// PendingCompactable is optionally implemented by channels that have a PendingHistory
// supporting LLM-based compaction. InstanceLoader uses this to wire compaction config
// after channel creation.
type PendingCompactable interface {
SetPendingCompaction(cfg *CompactionConfig)
}
// Truncate shortens a string to maxLen, appending "..." if truncated.
func Truncate(s string, maxLen int) string {
if len(s) <= maxLen {
return s
}
return s[:maxLen] + "..."
}