Files
Duc Nguyen aa15317d5b fix: resolve Zalo group by name + notify origin on failed forward (#1395)
* fix: resolve Zalo group name to real chat ID, notify origin on failed forward

Two related root causes behind "forward message to group by name" silently
failing while the agent reports success:

1. No tool let the agent resolve a group's display name (e.g. "Ban Dieu
   Hanh") to the chat ID the message tool actually requires. sessions_list
   only exposes session keys (numeric group IDs), never human-readable
   names, so the agent had no reliable way to turn a name into a real
   target and ended up passing the display name itself as `target`.

   Adds zalo_list_groups, wrapping the already-used (dashboard picker)
   protocol.FetchGroups behind the same optional-interface pattern as
   list_group_members/GroupMemberProvider (GroupListProvider on
   channels.Manager, gated to zalo_personal via RequiredChannelTypes).

2. When the resulting send fails downstream (e.g. Zalo rejects a bad
   chat_id), dispatchOutbound only ever retried/notified media failures on
   the same (already-broken) destination, and dropped text-only failures
   entirely — even though message.go's own postCrossTargetNotice comment
   states forwards must never announce a fake delivery. Because the bus
   publish is fire-and-forget, the tool had already returned "sent" and
   announced success to the origin chat before the real send was even
   attempted.

   message.go now tags cross-target forwards with origin channel/chat in
   OutboundMessage.Metadata; dispatchOutbound uses it to notify the ORIGIN
   chat with the real failure instead of silently dropping it or retrying
   against the same invalid target.

* test: cover forward-origin metadata tagging and dispatch failure notice

Extracts dispatchOutbound's error branch into handleSendFailure so it can
be unit tested without driving the consumer loop/goroutine, and adds
coverage for: forward failures notifying the origin chat (not the broken
destination), pre-existing non-forward media/text-only behavior staying
unchanged, message.go tagging cross-target group forwards with origin
metadata alongside group_id, and the new zalo_list_groups tool/Manager
delegator.
2026-07-08 23:23:08 +07:00

208 lines
9.7 KiB
Go

package bus
import (
"context"
"encoding/json"
"strings"
"github.com/google/uuid"
)
// MediaFile represents an inbound media file with its MIME type.
// Used throughout the media pipeline to preserve content type from channel download to storage.
type MediaFile struct {
Path string `json:"path"`
MimeType string `json:"mime_type,omitempty"` // e.g. "application/pdf", "image/jpeg"
Filename string `json:"filename,omitempty"` // original user-provided filename, e.g. "Báo cáo Q4.pdf"; empty → UUID fallback in persistMedia
Caption string `json:"caption,omitempty"` // optional outbound caption attached to this file
}
// InboundMessage represents a message received from a channel (Telegram, Discord, etc.)
type InboundMessage struct {
Channel string `json:"channel"`
SenderID string `json:"sender_id"`
ChatID string `json:"chat_id"`
Content string `json:"content"`
Media []MediaFile `json:"media,omitempty"`
SessionKey string `json:"session_key"` // deprecated: gateway builds canonical key
PeerKind string `json:"peer_kind,omitempty"` // "direct" or "group" (used for session key)
TenantID uuid.UUID `json:"tenant_id,omitempty"` // tenant scope from channel instance
AgentID string `json:"agent_id,omitempty"` // target agent (for multi-agent routing)
UserID string `json:"user_id,omitempty"` // external user ID for per-user scoping (memory, bootstrap)
HistoryLimit int `json:"history_limit,omitempty"` // max turns to keep in context (0=unlimited, from channel config)
ToolAllow []string `json:"tool_allow,omitempty"` // per-group tool allow list (nil = no restriction)
TelegramManagerPermissions []string `json:"telegram_manager_permissions,omitempty"` // hidden Telegram management permission groups for this inbound run
Metadata map[string]string `json:"metadata,omitempty"`
}
// OutboundMessage represents a message to be sent to a channel.
type OutboundMessage struct {
Channel string `json:"channel"`
ChatID string `json:"chat_id"`
Content string `json:"content"`
Media []MediaAttachment `json:"media,omitempty"` // optional media attachments
Metadata map[string]string `json:"metadata,omitempty"` // channel-specific metadata
TenantID uuid.UUID `json:"tenant_id,omitempty"` // tenant scope for per-tenant TTS
AgentID uuid.UUID `json:"agent_id,omitempty"` // agent scope for per-agent TTS voice override
AgentOtherConfig []byte `json:"agent_other_config,omitempty"` // agent's other_config for TTS voice/model
}
// Metadata keys on OutboundMessage.Metadata used to track the origin chat of
// a cross-target forward (message tool, forward=true). The outbound dispatch
// consumer runs async with no path back to the tool call that queued it, so
// these let it notify the ORIGIN chat if delivery to the forward target
// actually fails — otherwise a bad target (e.g. a display name instead of a
// real chat ID) fails silently downstream while the tool already reported
// success back to the model.
const (
MetaForwardOriginChannel = "forward_origin_channel"
MetaForwardOriginChatID = "forward_origin_chat_id"
)
// MediaAttachment represents a media file to be sent with a message.
type MediaAttachment struct {
URL string `json:"url"` // file path or URL
ContentType string `json:"content_type,omitempty"` // MIME type (e.g. "image/jpeg", "video/mp4")
Caption string `json:"caption,omitempty"` // optional caption for media
}
// Event represents a server-side event to broadcast to WebSocket clients.
type Event struct {
Name string `json:"name"` // event name (e.g. "agent", "chat", "health")
Payload any `json:"payload,omitempty"`
TenantID uuid.UUID `json:"-"` // tenant scope for event filtering (not serialized to clients)
}
// Cache invalidation kind constants.
const (
CacheKindAgent = "agent"
CacheKindBootstrap = "bootstrap"
CacheKindSkills = "skills"
CacheKindCron = "cron"
CacheKindChannelInstances = "channel_instances"
CacheKindBuiltinTools = "builtin_tools"
CacheKindTeam = "team"
CacheKindUserWorkspace = "user_workspace"
CacheKindSkillGrants = "skill_grants"
CacheKindMCP = "mcp"
CacheKindProvider = "provider"
CacheKindAPIKeys = "api_keys"
CacheKindHeartbeat = "heartbeat"
CacheKindConfigPerms = "config_perms"
CacheKindTenantUsers = "tenant_users"
CacheKindAgentAccess = "agent_access"
CacheKindTeamAccess = "team_access"
CacheKindTenants = "tenants"
)
// Topic constants for msgBus.Subscribe() / Broadcast().
const (
TopicCacheBootstrap = "cache:bootstrap"
TopicCacheAgent = "cache:agent"
TopicCacheSkills = "cache:skills"
TopicCacheCron = "cache:cron"
TopicCacheBuiltinTools = "cache:builtin_tools"
TopicCacheTeam = "cache:team"
TopicCacheUserWorkspace = "cache:user_workspace"
TopicCacheChannelInstances = "cache:channel_instances"
TopicCacheSkillGrants = "cache:skill_grants"
TopicCacheMCP = "cache:mcp"
TopicCacheProvider = "cache:provider"
TopicCacheHeartbeat = "cache:heartbeat"
TopicCacheConfigPerms = "cache:config_perms"
TopicAudit = "audit"
TopicTeamTaskAudit = "team-task-audit"
TopicChannelStreaming = "channel-streaming"
TopicConfigChanged = "config:changed"
TopicSystemConfigChanged = "system_config:changed"
TopicPairingRevoked = "pairing:revoked"
TopicAgentStatusChanged = "agent:status_changed"
TopicAgentDeleted = "agent:deleted"
)
// EventPairingRevoked is the event name broadcast when a paired device is revoked.
const EventPairingRevoked = "pairing.revoked"
// PairingRevokedPayload identifies the revoked device.
type PairingRevokedPayload struct {
SenderID string `json:"sender_id"`
Channel string `json:"channel"`
}
// EventAgentStatusChanged is broadcast when an agent's status changes (e.g., active → inactive).
const EventAgentStatusChanged = "agent.status_changed"
// AgentStatusChangedPayload carries agent status transition info for cascade operations.
type AgentStatusChangedPayload struct {
AgentID string `json:"agent_id"`
OldStatus string `json:"old_status"`
NewStatus string `json:"new_status"`
}
// AgentDeletedPayload carries agent deletion info for async cleanup (e.g. orphaned provider removal).
type AgentDeletedPayload struct {
AgentKey string `json:"agent_key"`
Provider string `json:"provider,omitempty"` // provider name for orphan cleanup
TenantID uuid.UUID `json:"tenant_id,omitempty"`
}
// AuditEventPayload carries audit log data emitted by handlers.
// A single subscriber persists these to the activity_logs table.
type AuditEventPayload struct {
ActorType string `json:"actor_type"`
ActorID string `json:"actor_id"`
Action string `json:"action"`
EntityType string `json:"entity_type"`
EntityID string `json:"entity_id"`
IPAddress string `json:"ip_address,omitempty"`
Details json.RawMessage `json:"details,omitempty"`
TenantID uuid.UUID `json:"tenant_id,omitempty"` // for async subscriber tenant scoping
}
// CacheInvalidatePayload signals cache layers to evict stale entries.
// Used with protocol.EventCacheInvalidate events. Events are delivered
// in-process via MessageBus and never marshaled to the wire, so the json
// tags are documentation-only (and omitempty on uuid.UUID is a no-op
// because uuid.UUID is [16]byte — all-zero arrays don't count as empty).
type CacheInvalidatePayload struct {
Kind string `json:"kind"` // CacheKind* constants
Key string `json:"key"` // agent_key, agent_id, etc. Empty = invalidate all
// TenantID scopes the invalidation to a single tenant. uuid.Nil means
// global (master admin action) — subscribers treat it as "invalidate all".
TenantID uuid.UUID `json:"tenant_id"`
}
// MessageHandler handles an inbound message from a specific channel.
type MessageHandler func(InboundMessage) error
// EventHandler handles a broadcast event.
type EventHandler func(Event)
// EventPublisher abstracts event broadcast + subscription.
// Used by gateway server and agents to decouple from concrete MessageBus.
type EventPublisher interface {
Subscribe(id string, handler EventHandler)
Unsubscribe(id string)
Broadcast(event Event)
}
// MessageRouter abstracts inbound/outbound message routing between channels and the agent runtime.
type MessageRouter interface {
PublishInbound(msg InboundMessage)
ConsumeInbound(ctx context.Context) (InboundMessage, bool)
PublishOutbound(msg OutboundMessage)
SubscribeOutbound(ctx context.Context) (OutboundMessage, bool)
}
// IsInternalSender returns true if the senderID belongs to an internal system
// component (not a real channel user). These should not be stored as contacts
// and must be rejected by per-user permission checks in group contexts (#915).
func IsInternalSender(senderID string) bool {
return strings.HasPrefix(senderID, "system:") ||
strings.HasPrefix(senderID, "notification:") ||
strings.HasPrefix(senderID, "teammate:") ||
strings.HasPrefix(senderID, "ticker:") ||
strings.HasPrefix(senderID, "subagent:") ||
senderID == "session_send_tool"
}