diff --git a/cmd/gateway_consumer.go b/cmd/gateway_consumer.go index 6f67a1fd..0e960c31 100644 --- a/cmd/gateway_consumer.go +++ b/cmd/gateway_consumer.go @@ -170,6 +170,20 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents runID := fmt.Sprintf("inbound-%s-%s-%s", msg.Channel, msg.ChatID, uuid.NewString()[:8]) + // Build outbound metadata for reply-to + thread routing BEFORE RegisterRun + // so block.reply handler can use it for routing intermediate messages. + outMeta := make(map[string]string) + if isGroup { + if mid := msg.Metadata["message_id"]; mid != "" { + outMeta["reply_to_message_id"] = mid + } + } + for _, k := range []string{"message_thread_id", "local_key", "placeholder_key", "group_id"} { + if v := msg.Metadata[k]; v != "" { + outMeta[k] = v + } + } + // Register run with channel manager for streaming/reaction event forwarding. // Use localKey (composite key with topic suffix) so streaming/reaction events // route to the correct per-topic state in the channel. @@ -178,8 +192,9 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents if lk := msg.Metadata["local_key"]; lk != "" { chatIDForRun = lk } + blockReply := channelMgr != nil && channelMgr.ResolveBlockReply(msg.Channel, cfg.Gateway.BlockReply) if channelMgr != nil { - channelMgr.RegisterRun(runID, msg.Channel, chatIDForRun, messageID) + channelMgr.RegisterRun(runID, msg.Channel, chatIDForRun, messageID, outMeta, enableStream, blockReply) } // Group-aware system prompt: help the LLM adapt tone and behavior for group chats. @@ -237,23 +252,8 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents MaxConcurrent: maxConcurrent, }) - // Build outbound metadata for reply-to + thread routing. - // Groups: reply to user's message so context is clear in busy chats. - // DMs: no reply needed — response edits the placeholder or sends inline. - outMeta := make(map[string]string) - if isGroup { - if mid := msg.Metadata["message_id"]; mid != "" { - outMeta["reply_to_message_id"] = mid - } - } - for _, k := range []string{"message_thread_id", "local_key", "placeholder_key", "group_id"} { - if v := msg.Metadata[k]; v != "" { - outMeta[k] = v - } - } - // Handle result asynchronously to not block the flush callback. - go func(channel, chatID, session, rID string, meta map[string]string) { + go func(channel, chatID, session, rID string, meta map[string]string, blockReplyEnabled bool) { outcome := <-outCh // Clean up run tracking (in case HandleAgentEvent didn't fire for terminal events) @@ -301,6 +301,21 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents return } + // Dedup: if block replies were delivered and the final content matches the last + // block reply, suppress the final message to avoid duplicate delivery. + // Only applies when blockReply is enabled (otherwise nothing was delivered). + if blockReplyEnabled && outcome.Result.BlockReplies > 0 && outcome.Result.Content == outcome.Result.LastBlockReply && len(outcome.Result.Media) == 0 { + slog.Debug("inbound: dedup final message (matches last block reply)", + "channel", channel, "run_id", rID) + msgBus.PublishOutbound(bus.OutboundMessage{ + Channel: channel, + ChatID: chatID, + Content: "", + Metadata: meta, + }) + return + } + // Publish response back to the channel outMsg := bus.OutboundMessage{ Channel: channel, @@ -324,7 +339,7 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents } msgBus.PublishOutbound(outMsg) - }(msg.Channel, msg.ChatID, sessionKey, runID, outMeta) + }(msg.Channel, msg.ChatID, sessionKey, runID, outMeta, blockReply) } // Inbound debounce: merge rapid messages from the same sender before processing. diff --git a/internal/agent/loop.go b/internal/agent/loop.go index 50d4d37e..59aa9d7d 100644 --- a/internal/agent/loop.go +++ b/internal/agent/loop.go @@ -369,6 +369,8 @@ func (l *Loop) runLoop(ctx context.Context, req RunRequest) (*RunResult, error) var asyncToolCalls []string // track async spawn tool names for fallback var mediaResults []MediaResult // media files from tool MEDIA: results var deliverables []string // actual content from tool outputs (for team task results) + var blockReplies int // count of block.reply events emitted (for dedup in consumer) + var lastBlockReply string // last block reply content // Mid-loop compaction: summarize in-memory messages when context exceeds threshold. // Uses same config as maybeSummarize (contextWindow * historyShare). @@ -577,6 +579,22 @@ func (l *Loop) runLoop(ctx context.Context, req RunRequest) (*RunResult, error) messages = append(messages, assistantMsg) pendingMsgs = append(pendingMsgs, assistantMsg) + // Emit block.reply for intermediate assistant content during tool iterations. + // Non-streaming channels (Zalo, Discord, WhatsApp) would otherwise lose this text. + if resp.Content != "" { + sanitized := SanitizeAssistantContent(resp.Content) + if sanitized != "" && !IsSilentReply(sanitized) { + blockReplies++ + lastBlockReply = sanitized + l.emit(AgentEvent{ + Type: protocol.AgentEventBlockReply, + AgentID: l.id, + RunID: req.RunID, + Payload: map[string]string{"content": sanitized}, + }) + } + } + // Track team_tasks create for orphan detection (argument-based, pre-execution). // Spawn counting is done post-execution so failed spawns don't get counted. for _, tc := range resp.ToolCalls { @@ -905,12 +923,14 @@ func (l *Loop) runLoop(ctx context.Context, req RunRequest) (*RunResult, error) } return &RunResult{ - Content: finalContent, - RunID: req.RunID, - Iterations: iteration, - Usage: &totalUsage, - Media: mediaResults, - Deliverables: deliverables, + Content: finalContent, + RunID: req.RunID, + Iterations: iteration, + Usage: &totalUsage, + Media: mediaResults, + Deliverables: deliverables, + BlockReplies: blockReplies, + LastBlockReply: lastBlockReply, }, nil } diff --git a/internal/agent/loop_types.go b/internal/agent/loop_types.go index 13ff6b82..65991b2d 100644 --- a/internal/agent/loop_types.go +++ b/internal/agent/loop_types.go @@ -291,8 +291,10 @@ type RunResult struct { RunID string `json:"runId"` Iterations int `json:"iterations"` Usage *providers.Usage `json:"usage,omitempty"` - Media []MediaResult `json:"media,omitempty"` // media files from tool results (MEDIA: prefix) - Deliverables []string `json:"deliverables,omitempty"` // actual content from tool outputs (for team task results) + Media []MediaResult `json:"media,omitempty"` // media files from tool results (MEDIA: prefix) + Deliverables []string `json:"deliverables,omitempty"` // actual content from tool outputs (for team task results) + BlockReplies int `json:"blockReplies,omitempty"` // number of block.reply events emitted + LastBlockReply string `json:"lastBlockReply,omitempty"` // last block reply content (for dedup) } // MediaResult represents a media file produced by a tool during the agent run. diff --git a/internal/channels/channel.go b/internal/channels/channel.go index a3ca4035..a2526466 100644 --- a/internal/channels/channel.go +++ b/internal/channels/channel.go @@ -89,6 +89,12 @@ type StreamingChannel interface { OnStreamEnd(ctx context.Context, chatID string, finalText string) 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 +} + // 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 diff --git a/internal/channels/discord/discord.go b/internal/channels/discord/discord.go index 68a068c1..53fca85f 100644 --- a/internal/channels/discord/discord.go +++ b/internal/channels/discord/discord.go @@ -95,6 +95,9 @@ func (c *Channel) Start(_ context.Context) error { return nil } +// BlockReplyEnabled returns the per-channel block_reply override (nil = inherit gateway default). +func (c *Channel) BlockReplyEnabled() *bool { return c.config.BlockReply } + // Stop closes the Discord gateway connection. func (c *Channel) Stop(_ context.Context) error { slog.Info("stopping discord bot") diff --git a/internal/channels/discord/factory.go b/internal/channels/discord/factory.go index a19edd7a..1c463381 100644 --- a/internal/channels/discord/factory.go +++ b/internal/channels/discord/factory.go @@ -22,6 +22,7 @@ type discordInstanceConfig struct { AllowFrom []string `json:"allow_from,omitempty"` RequireMention *bool `json:"require_mention,omitempty"` HistoryLimit int `json:"history_limit,omitempty"` + BlockReply *bool `json:"block_reply,omitempty"` } // Factory creates a Discord channel from DB instance data. @@ -53,6 +54,7 @@ func Factory(name string, creds json.RawMessage, cfg json.RawMessage, GroupPolicy: ic.GroupPolicy, RequireMention: ic.RequireMention, HistoryLimit: ic.HistoryLimit, + BlockReply: ic.BlockReply, } // DB instances default to "pairing" for groups (secure by default). diff --git a/internal/channels/feishu/factory.go b/internal/channels/feishu/factory.go index 034c80eb..756539f6 100644 --- a/internal/channels/feishu/factory.go +++ b/internal/channels/feishu/factory.go @@ -36,6 +36,7 @@ type feishuInstanceConfig struct { Streaming *bool `json:"streaming,omitempty"` ReactionLevel string `json:"reaction_level,omitempty"` HistoryLimit int `json:"history_limit,omitempty"` + BlockReply *bool `json:"block_reply,omitempty"` } // Factory creates a Feishu/Lark channel from DB instance data. @@ -81,6 +82,7 @@ func Factory(name string, creds json.RawMessage, cfg json.RawMessage, Streaming: ic.Streaming, ReactionLevel: ic.ReactionLevel, HistoryLimit: ic.HistoryLimit, + BlockReply: ic.BlockReply, } // DB instances default to "pairing" for groups (secure by default). diff --git a/internal/channels/feishu/feishu.go b/internal/channels/feishu/feishu.go index 0aab4a5b..1a701872 100644 --- a/internal/channels/feishu/feishu.go +++ b/internal/channels/feishu/feishu.go @@ -116,6 +116,9 @@ func (c *Channel) Start(ctx context.Context) error { } } +// BlockReplyEnabled returns the per-channel block_reply override (nil = inherit gateway default). +func (c *Channel) BlockReplyEnabled() *bool { return c.cfg.BlockReply } + // Stop shuts down the Feishu channel. func (c *Channel) Stop(_ context.Context) error { slog.Info("stopping feishu/lark bot") diff --git a/internal/channels/manager.go b/internal/channels/manager.go index 47eda10a..5506b8b0 100644 --- a/internal/channels/manager.go +++ b/internal/channels/manager.go @@ -19,6 +19,9 @@ type RunContext struct { ChannelName string ChatID string MessageID string // platform message ID (string to support Feishu "om_xxx", Telegram "12345", etc.) + Metadata map[string]string // outbound routing metadata (thread_id, local_key, group_id) + Streaming bool // whether run uses streaming (to avoid double-delivery of block replies) + BlockReplyEnabled bool // whether block.reply delivery is enabled for this run (resolved at RegisterRun time) mu sync.Mutex streamBuffer string // accumulated streaming text (chunks are deltas) inToolPhase bool // true after tool.call, reset on next chunk (new LLM iteration) @@ -260,11 +263,14 @@ func (m *Manager) SendToChannel(ctx context.Context, channelName, chatID, conten // 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) { +func (m *Manager) RegisterRun(runID, channelName, chatID, messageID string, metadata map[string]string, streaming, blockReply bool) { m.runs.Store(runID, &RunContext{ - ChannelName: channelName, - ChatID: chatID, - MessageID: messageID, + ChannelName: channelName, + ChatID: chatID, + MessageID: messageID, + Metadata: metadata, + Streaming: streaming, + BlockReplyEnabled: blockReply, }) } @@ -290,6 +296,22 @@ func (m *Manager) IsStreamingChannel(channelName string, isGroup bool) bool { 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 +} + // HandleAgentEvent routes agent lifecycle events to streaming/reaction channels. // Called from the bus event subscriber — must be non-blocking. // eventType: "run.started", "chunk", "tool.call", "tool.result", "run.completed", "run.failed" @@ -365,6 +387,49 @@ func (m *Manager) HandleAgentEvent(eventType, runID string, payload interface{}) } } + // Handle block.reply: deliver intermediate assistant text to non-streaming channels. + // Gated by BlockReplyEnabled (resolved from gateway + per-channel config at RegisterRun time). + // Streaming channels already deliver via chunks, so skip to avoid double-delivery. + if eventType == protocol.AgentEventBlockReply { + if !rc.BlockReplyEnabled { + return + } + content := extractPayloadString(payload, "content") + if content == "" { + return + } + rc.mu.Lock() + streaming := rc.Streaming + rc.mu.Unlock() + + if streaming { + return // streaming already delivered via chunks + } + + // Build outbound metadata: copy routing fields but strip reply_to_message_id + // (block replies are standalone) and placeholder_key (reserve for final message). + var outMeta map[string]string + if rc.Metadata != nil { + outMeta = make(map[string]string) + for _, k := range []string{"message_thread_id", "local_key", "group_id"} { + if v := rc.Metadata[k]; v != "" { + outMeta[k] = v + } + } + if len(outMeta) == 0 { + outMeta = nil + } + } + + m.bus.PublishOutbound(bus.OutboundMessage{ + Channel: rc.ChannelName, + ChatID: rc.ChatID, + Content: content, + Metadata: outMeta, + }) + return + } + // Handle LLM retry: update placeholder to notify user if eventType == protocol.AgentEventRunRetrying { attempt := extractPayloadString(payload, "attempt") diff --git a/internal/channels/telegram/channel.go b/internal/channels/telegram/channel.go index 64c25baa..41a4402a 100644 --- a/internal/channels/telegram/channel.go +++ b/internal/channels/telegram/channel.go @@ -202,6 +202,9 @@ func (c *Channel) StreamEnabled(isGroup bool) bool { return c.config.DMStream != nil && *c.config.DMStream } +// BlockReplyEnabled returns the per-channel block_reply override (nil = inherit gateway default). +func (c *Channel) BlockReplyEnabled() *bool { return c.config.BlockReply } + // Stop shuts down the Telegram bot by cancelling the long polling context // and waiting for the polling goroutine to exit. func (c *Channel) Stop(_ context.Context) error { diff --git a/internal/channels/telegram/factory.go b/internal/channels/telegram/factory.go index 1ec4c5e8..cbda42e6 100644 --- a/internal/channels/telegram/factory.go +++ b/internal/channels/telegram/factory.go @@ -27,6 +27,7 @@ type telegramInstanceConfig struct { ReactionLevel string `json:"reaction_level,omitempty"` MediaMaxBytes int64 `json:"media_max_bytes,omitempty"` LinkPreview *bool `json:"link_preview,omitempty"` + BlockReply *bool `json:"block_reply,omitempty"` AllowFrom []string `json:"allow_from,omitempty"` } @@ -79,6 +80,7 @@ func buildChannel(name string, creds json.RawMessage, cfg json.RawMessage, ReactionLevel: ic.ReactionLevel, MediaMaxBytes: ic.MediaMaxBytes, LinkPreview: ic.LinkPreview, + BlockReply: ic.BlockReply, } // DB instances default to "pairing" for groups (secure by default). diff --git a/internal/channels/whatsapp/factory.go b/internal/channels/whatsapp/factory.go index e9016f6d..3a039f33 100644 --- a/internal/channels/whatsapp/factory.go +++ b/internal/channels/whatsapp/factory.go @@ -20,6 +20,7 @@ type whatsappInstanceConfig struct { DMPolicy string `json:"dm_policy,omitempty"` GroupPolicy string `json:"group_policy,omitempty"` AllowFrom []string `json:"allow_from,omitempty"` + BlockReply *bool `json:"block_reply,omitempty"` } // Factory creates a WhatsApp channel from DB instance data. @@ -49,6 +50,7 @@ func Factory(name string, creds json.RawMessage, cfg json.RawMessage, AllowFrom: ic.AllowFrom, DMPolicy: ic.DMPolicy, GroupPolicy: ic.GroupPolicy, + BlockReply: ic.BlockReply, } // DB instances default to "pairing" for groups (secure by default). diff --git a/internal/channels/whatsapp/whatsapp.go b/internal/channels/whatsapp/whatsapp.go index bd35bdf0..bca0c348 100644 --- a/internal/channels/whatsapp/whatsapp.go +++ b/internal/channels/whatsapp/whatsapp.go @@ -68,6 +68,9 @@ func (c *Channel) Start(ctx context.Context) error { return nil } +// BlockReplyEnabled returns the per-channel block_reply override (nil = inherit gateway default). +func (c *Channel) BlockReplyEnabled() *bool { return c.config.BlockReply } + // Stop gracefully shuts down the WhatsApp channel. func (c *Channel) Stop(_ context.Context) error { slog.Info("stopping whatsapp channel") diff --git a/internal/channels/zalo/factory.go b/internal/channels/zalo/factory.go index 5ed5d5c7..661afd85 100644 --- a/internal/channels/zalo/factory.go +++ b/internal/channels/zalo/factory.go @@ -22,6 +22,7 @@ type zaloInstanceConfig struct { WebhookURL string `json:"webhook_url,omitempty"` MediaMaxMB int `json:"media_max_mb,omitempty"` AllowFrom []string `json:"allow_from,omitempty"` + BlockReply *bool `json:"block_reply,omitempty"` } // Factory creates a Zalo OA channel from DB instance data. @@ -53,6 +54,7 @@ func Factory(name string, creds json.RawMessage, cfg json.RawMessage, WebhookURL: ic.WebhookURL, WebhookSecret: c.WebhookSecret, MediaMaxMB: ic.MediaMaxMB, + BlockReply: ic.BlockReply, } ch, err := New(zCfg, msgBus, pairingSvc) diff --git a/internal/channels/zalo/personal/channel.go b/internal/channels/zalo/personal/channel.go index 54e2fe31..312f6c5e 100644 --- a/internal/channels/zalo/personal/channel.go +++ b/internal/channels/zalo/personal/channel.go @@ -76,6 +76,9 @@ func New(cfg config.ZaloPersonalConfig, msgBus *bus.MessageBus, pairingSvc store }, nil } +// BlockReplyEnabled returns the per-channel block_reply override (nil = inherit gateway default). +func (c *Channel) BlockReplyEnabled() *bool { return c.config.BlockReply } + // session returns the current session snapshot (thread-safe). func (c *Channel) session() *protocol.Session { c.mu.RLock() diff --git a/internal/channels/zalo/personal/factory.go b/internal/channels/zalo/personal/factory.go index a31fa566..57f51139 100644 --- a/internal/channels/zalo/personal/factory.go +++ b/internal/channels/zalo/personal/factory.go @@ -25,6 +25,7 @@ type zaloInstanceConfig struct { GroupPolicy string `json:"group_policy,omitempty"` RequireMention *bool `json:"require_mention,omitempty"` AllowFrom []string `json:"allow_from,omitempty"` + BlockReply *bool `json:"block_reply,omitempty"` } // Factory creates a Zalo Personal channel from DB instance data. @@ -58,6 +59,7 @@ func Factory(name string, creds json.RawMessage, cfg json.RawMessage, DMPolicy: ic.DMPolicy, GroupPolicy: ic.GroupPolicy, RequireMention: ic.RequireMention, + BlockReply: ic.BlockReply, } ch, err := New(zaloCfg, msgBus, pairingSvc) diff --git a/internal/channels/zalo/zalo.go b/internal/channels/zalo/zalo.go index 1bd1984f..ffd8e653 100644 --- a/internal/channels/zalo/zalo.go +++ b/internal/channels/zalo/zalo.go @@ -38,6 +38,7 @@ type Channel struct { token string dmPolicy string mediaMaxMB int + blockReply *bool pairingService store.PairingStore pairingDebounce sync.Map // senderID → time.Time stopCh chan struct{} @@ -68,12 +69,16 @@ func New(cfg config.ZaloConfig, msgBus *bus.MessageBus, pairingSvc store.Pairing token: cfg.Token, dmPolicy: dmPolicy, mediaMaxMB: mediaMax, + blockReply: cfg.BlockReply, pairingService: pairingSvc, stopCh: make(chan struct{}), client: &http.Client{Timeout: 60 * time.Second}, }, nil } +// BlockReplyEnabled returns the per-channel block_reply override (nil = inherit gateway default). +func (c *Channel) BlockReplyEnabled() *bool { return c.blockReply } + // Start begins polling for Zalo updates. func (c *Channel) Start(ctx context.Context) error { slog.Info("starting zalo bot (polling mode)") diff --git a/internal/config/config_channels.go b/internal/config/config_channels.go index a16cb195..5b927066 100644 --- a/internal/config/config_channels.go +++ b/internal/config/config_channels.go @@ -25,6 +25,7 @@ type TelegramConfig struct { ReactionLevel string `json:"reaction_level,omitempty"` // "off" (default), "minimal", "full" — status emoji reactions MediaMaxBytes int64 `json:"media_max_bytes,omitempty"` // max media download size in bytes (default 20MB) LinkPreview *bool `json:"link_preview,omitempty"` // enable URL previews in messages (default true) + BlockReply *bool `json:"block_reply,omitempty"` // override gateway block_reply (nil = inherit) // Optional STT (Speech-to-Text) pipeline for voice/audio inbound messages. // When stt_proxy_url is set, audio/voice messages are transcribed before being forwarded to the agent. @@ -76,6 +77,7 @@ type DiscordConfig struct { GroupPolicy string `json:"group_policy,omitempty"` // "open" (default), "allowlist", "disabled" RequireMention *bool `json:"require_mention,omitempty"` // require @bot mention in groups (default true) HistoryLimit int `json:"history_limit,omitempty"` // max pending group messages for context (default 50, 0=disabled) + BlockReply *bool `json:"block_reply,omitempty"` // override gateway block_reply (nil = inherit) } type SlackConfig struct { @@ -94,6 +96,7 @@ type WhatsAppConfig struct { AllowFrom FlexibleStringSlice `json:"allow_from"` DMPolicy string `json:"dm_policy,omitempty"` // "open" (default), "allowlist", "disabled" GroupPolicy string `json:"group_policy,omitempty"` // "open" (default), "allowlist", "disabled" + BlockReply *bool `json:"block_reply,omitempty"` // override gateway block_reply (nil = inherit) } type ZaloConfig struct { @@ -104,6 +107,7 @@ type ZaloConfig struct { WebhookURL string `json:"webhook_url,omitempty"` WebhookSecret string `json:"webhook_secret,omitempty"` MediaMaxMB int `json:"media_max_mb,omitempty"` // default 5 + BlockReply *bool `json:"block_reply,omitempty"` // override gateway block_reply (nil = inherit) } type ZaloPersonalConfig struct { @@ -113,6 +117,7 @@ type ZaloPersonalConfig struct { GroupPolicy string `json:"group_policy,omitempty"` // "open" (default), "allowlist", "disabled" RequireMention *bool `json:"require_mention,omitempty"` // require @bot mention in groups (default true) CredentialsPath string `json:"credentials_path,omitempty"` // path to saved cookies JSON + BlockReply *bool `json:"block_reply,omitempty"` // override gateway block_reply (nil = inherit) } type FeishuConfig struct { @@ -137,6 +142,7 @@ type FeishuConfig struct { Streaming *bool `json:"streaming,omitempty"` // default true ReactionLevel string `json:"reaction_level,omitempty"` // "off" (default), "minimal", "full" — typing emoji reactions HistoryLimit int `json:"history_limit,omitempty"` + BlockReply *bool `json:"block_reply,omitempty"` // override gateway block_reply (nil = inherit) } // ProvidersConfig maps provider name to its config. @@ -211,6 +217,7 @@ type GatewayConfig struct { InjectionAction string `json:"injection_action,omitempty"` // prompt injection action: "log", "warn" (default), "block", "off" InboundDebounceMs int `json:"inbound_debounce_ms,omitempty"` // merge rapid messages from same sender (default 1000ms, -1 = disabled) Quota *QuotaConfig `json:"quota,omitempty"` // per-user/group request quotas + BlockReply *bool `json:"block_reply,omitempty"` // deliver intermediate text during tool iterations (default false) } // ToolsConfig controls tool availability, policy, and web search. diff --git a/pkg/protocol/events.go b/pkg/protocol/events.go index 4c2bd9e8..8c80535e 100644 --- a/pkg/protocol/events.go +++ b/pkg/protocol/events.go @@ -72,6 +72,7 @@ const ( AgentEventRunRetrying = "run.retrying" AgentEventToolCall = "tool.call" AgentEventToolResult = "tool.result" + AgentEventBlockReply = "block.reply" ) // Chat event subtypes (in payload.type) diff --git a/ui/web/src/api/protocol.ts b/ui/web/src/api/protocol.ts index f75417fe..c55ff0b7 100644 --- a/ui/web/src/api/protocol.ts +++ b/ui/web/src/api/protocol.ts @@ -221,6 +221,7 @@ export const AgentEventTypes = { RUN_FAILED: "run.failed", TOOL_CALL: "tool.call", TOOL_RESULT: "tool.result", + BLOCK_REPLY: "block.reply", } as const; // Chat event subtypes (in payload.type) diff --git a/ui/web/src/pages/channels/channel-instance-form-dialog.tsx b/ui/web/src/pages/channels/channel-instance-form-dialog.tsx index 276c1b3c..b8503d53 100644 --- a/ui/web/src/pages/channels/channel-instance-form-dialog.tsx +++ b/ui/web/src/pages/channels/channel-instance-form-dialog.tsx @@ -20,7 +20,7 @@ import { import type { ChannelInstanceData, ChannelInstanceInput } from "./hooks/use-channel-instances"; import type { AgentData } from "@/types/agent"; import { slugify, isValidSlug } from "@/lib/slug"; -import { credentialsSchema, configSchema, wizardConfig } from "./channel-schemas"; +import { credentialsSchema, configSchema, wizardConfig, type FieldDef } from "./channel-schemas"; import { ChannelFields } from "./channel-fields"; import { wizardAuthSteps, wizardConfigSteps, wizardEditConfigs } from "./channel-wizard-registry"; import { TelegramGroupOverrides } from "./telegram-group-overrides"; @@ -89,7 +89,16 @@ export function ChannelInstanceFormDialog({ for (const f of schema) { if (f.defaultValue !== undefined) defaults[f.key] = f.defaultValue; } - setConfigValues({ ...defaults, ...(instance?.config ?? {}) }); + const merged: Record = { ...defaults, ...(instance?.config ?? {}) }; + // Convert boolean values to strings for select fields that use "true"/"false" options + const boolSelectKeys = new Set( + schema.filter((f) => f.type === "select" && f.options?.some((o) => o.value === "true")).map((f) => f.key), + ); + for (const key of boolSelectKeys) { + if (typeof merged[key] === "boolean") merged[key] = String(merged[key]); + else if (merged[key] === undefined || merged[key] === null) merged[key] = "inherit"; + } + setConfigValues(merged); setEnabled(instance?.enabled ?? true); setError(""); setStep("form"); @@ -117,6 +126,20 @@ export function ChannelInstanceFormDialog({ setConfigValues((prev) => ({ ...prev, [key]: value })); }, []); + // Convert select fields with "true"/"false"/"inherit" values to proper JSON types. + // "inherit" → remove key (nil on Go side), "true"/"false" → boolean. + const coerceBoolSelects = (cfg: Record, schema: FieldDef[]) => { + const boolSelectKeys = new Set( + schema.filter((f) => f.type === "select" && f.options?.some((o) => o.value === "true")).map((f) => f.key), + ); + for (const key of boolSelectKeys) { + const v = cfg[key]; + if (v === "true") cfg[key] = true; + else if (v === "false") cfg[key] = false; + else delete cfg[key]; // "inherit" or unset + } + }; + const handleSubmit = async () => { if (!name.trim()) { setError("Name is required"); return; } if (!isValidSlug(name.trim())) { @@ -137,6 +160,7 @@ export function ChannelInstanceFormDialog({ const cleanConfig = Object.fromEntries( Object.entries(configValues).filter(([, v]) => v !== undefined && v !== "" && v !== null), ); + coerceBoolSelects(cleanConfig, configSchema[channelType] ?? []); const cleanCreds = Object.fromEntries( Object.entries(credsValues).filter(([, v]) => v !== undefined && v !== "" && v !== null), ); @@ -180,6 +204,7 @@ export function ChannelInstanceFormDialog({ const cleanConfig = Object.fromEntries( Object.entries(configValues).filter(([, v]) => v !== undefined && v !== "" && v !== null), ); + coerceBoolSelects(cleanConfig, configSchema[channelType] ?? []); setLoading(true); setError(""); try { diff --git a/ui/web/src/pages/channels/channel-schemas.ts b/ui/web/src/pages/channels/channel-schemas.ts index aba10420..04537e2e 100644 --- a/ui/web/src/pages/channels/channel-schemas.ts +++ b/ui/web/src/pages/channels/channel-schemas.ts @@ -14,6 +14,12 @@ export interface FieldDef { // --- Shared option lists --- +const blockReplyOptions = [ + { value: "inherit", label: "Inherit from gateway" }, + { value: "true", label: "Enabled" }, + { value: "false", label: "Disabled" }, +]; + const dmPolicyOptions = [ { value: "pairing", label: "Pairing (require code)" }, { value: "open", label: "Open (accept all)" }, @@ -68,6 +74,7 @@ export const configSchema: Record = { { key: "media_max_bytes", label: "Max Media Size (bytes)", type: "number", defaultValue: 20971520, help: "Default: 20MB" }, { key: "link_preview", label: "Link Preview", type: "boolean", defaultValue: true }, { key: "allow_from", label: "Allowed Users", type: "tags", help: "User IDs or @usernames, one per line" }, + { key: "block_reply", label: "Block Reply", type: "select", options: blockReplyOptions, defaultValue: "inherit", help: "Deliver intermediate text during tool iterations" }, ], discord: [ { key: "dm_policy", label: "DM Policy", type: "select", options: dmPolicyOptions, defaultValue: "open" }, @@ -75,6 +82,7 @@ export const configSchema: Record = { { key: "require_mention", label: "Require @mention in groups", type: "boolean", defaultValue: true }, { key: "history_limit", label: "Group History Limit", type: "number", defaultValue: 50, help: "Max pending group messages for context (0 = disabled)" }, { key: "allow_from", label: "Allowed Users", type: "tags", help: "Discord user IDs" }, + { key: "block_reply", label: "Block Reply", type: "select", options: blockReplyOptions, defaultValue: "inherit", help: "Deliver intermediate text during tool iterations" }, ], feishu: [ { key: "domain", label: "Domain", type: "select", options: [{ value: "lark", label: "Lark (global) — webhook only" }, { value: "feishu", label: "Feishu (China)" }], defaultValue: "lark" }, @@ -92,23 +100,27 @@ export const configSchema: Record = { { key: "reaction_level", label: "Reaction Level", type: "select", options: [{ value: "off", label: "Off" }, { value: "minimal", label: "Minimal" }, { value: "full", label: "Full" }], defaultValue: "off", help: "Typing emoji reaction on user messages while bot is processing" }, { key: "allow_from", label: "Allowed Users", type: "tags", help: "Lark open_ids (ou_...)" }, { key: "group_allow_from", label: "Group Allowed Users", type: "tags", help: "Separate allowlist for group senders" }, + { key: "block_reply", label: "Block Reply", type: "select", options: blockReplyOptions, defaultValue: "inherit", help: "Deliver intermediate text during tool iterations" }, ], zalo_oa: [ { key: "dm_policy", label: "DM Policy", type: "select", options: dmPolicyOptions, defaultValue: "pairing" }, { key: "webhook_url", label: "Webhook URL", type: "text", placeholder: "https://..." }, { key: "media_max_mb", label: "Max Media Size (MB)", type: "number", defaultValue: 5 }, { key: "allow_from", label: "Allowed Users", type: "tags", help: "Zalo user IDs" }, + { key: "block_reply", label: "Block Reply", type: "select", options: blockReplyOptions, defaultValue: "inherit", help: "Deliver intermediate text during tool iterations" }, ], zalo_personal: [ { key: "dm_policy", label: "DM Policy", type: "select", options: dmPolicyOptions, defaultValue: "allowlist" }, { key: "group_policy", label: "Group Policy", type: "select", options: groupPolicyOptions, defaultValue: "allowlist" }, { key: "require_mention", label: "Require @mention in groups", type: "boolean", defaultValue: true }, { key: "allow_from", label: "Allowed Users", type: "tags", help: "Zalo user IDs or group IDs" }, + { key: "block_reply", label: "Block Reply", type: "select", options: blockReplyOptions, defaultValue: "inherit", help: "Deliver intermediate text during tool iterations" }, ], whatsapp: [ { key: "dm_policy", label: "DM Policy", type: "select", options: dmPolicyOptions, defaultValue: "open" }, { key: "group_policy", label: "Group Policy", type: "select", options: groupPolicyOptions, defaultValue: "open" }, { key: "allow_from", label: "Allowed Users", type: "tags", help: "WhatsApp user IDs" }, + { key: "block_reply", label: "Block Reply", type: "select", options: blockReplyOptions, defaultValue: "inherit", help: "Deliver intermediate text during tool iterations" }, ], }; diff --git a/ui/web/src/pages/config/sections/gateway-section.tsx b/ui/web/src/pages/config/sections/gateway-section.tsx index 33b46673..98c82e1b 100644 --- a/ui/web/src/pages/config/sections/gateway-section.tsx +++ b/ui/web/src/pages/config/sections/gateway-section.tsx @@ -2,6 +2,7 @@ import { useState, useEffect } from "react"; import { Save } from "lucide-react"; import { Button } from "@/components/ui/button"; import { Input } from "@/components/ui/input"; +import { Switch } from "@/components/ui/switch"; import { Card, CardContent, CardDescription, CardHeader, CardTitle } from "@/components/ui/card"; import { Select, SelectContent, SelectItem, SelectTrigger, SelectValue } from "@/components/ui/select"; import { InfoLabel } from "@/components/shared/info-label"; @@ -17,6 +18,7 @@ interface GatewayData { rate_limit_rpm?: number; injection_action?: string; inbound_debounce_ms?: number; + block_reply?: boolean; } const DEFAULT: GatewayData = {}; @@ -159,6 +161,12 @@ export function GatewaySection({ data, onSave, saving }: Props) { +
+ Block Reply +
+ update({ block_reply: v })} /> +
+
{dirty && (