mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-11 03:13:24 +00:00
* 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.
208 lines
9.7 KiB
Go
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"
|
|
}
|