mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-03 07:12:50 +00:00
feat(block-reply): deliver intermediate text during tool iterations with 2-tier config (#55)
Add block.reply event that delivers intermediate assistant text to non-streaming channels during multi-tool iterations. Includes 2-tier config toggle: gateway-level default (disabled) + per-channel override (inherit/on/off). Backend: - Emit block.reply events from agent loop between tool iterations - Add BlockReply *bool to GatewayConfig and all 6 channel config structs - Add BlockReplyChannel interface with ResolveBlockReply() resolution - Guard delivery in HandleAgentEvent by RunContext.BlockReplyEnabled - Resolve config at RegisterRun time, pass to consumer goroutine - Conditional dedup: skip final message if identical to last block reply UI: - Gateway settings: Switch toggle for global default - Per-channel: tri-state select (Inherit from gateway / Enabled / Disabled) - Protocol: BLOCK_REPLY constant in AgentEventTypes - Form: coerceBoolSelects for proper JSON boolean serialization
This commit is contained in:
1 parent
3dfdc8ca39
commit
faa47abfb6
23 files changed
+226
-32
No files matched your search
+33
-18
@@ -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.
|
||||
|
||||
+26
-6
@@ -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
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)")
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -72,6 +72,7 @@ const (
|
||||
AgentEventRunRetrying = "run.retrying"
|
||||
AgentEventToolCall = "tool.call"
|
||||
AgentEventToolResult = "tool.result"
|
||||
AgentEventBlockReply = "block.reply"
|
||||
)
|
||||
|
||||
// Chat event subtypes (in payload.type)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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<string, unknown> = { ...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<string, unknown>, 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 {
|
||||
|
||||
@@ -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<string, FieldDef[]> = {
|
||||
{ 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<string, FieldDef[]> = {
|
||||
{ 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<string, FieldDef[]> = {
|
||||
{ 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" },
|
||||
],
|
||||
};
|
||||
|
||||
|
||||
@@ -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) {
|
||||
</SelectContent>
|
||||
</Select>
|
||||
</div>
|
||||
<div className="grid gap-1.5">
|
||||
<InfoLabel tip="Deliver intermediate assistant text to non-streaming channels during tool iterations. Each assistant block is sent before tool execution.">Block Reply</InfoLabel>
|
||||
<div className="flex items-center h-9">
|
||||
<Switch checked={draft.block_reply ?? false} onCheckedChange={(v) => update({ block_reply: v })} />
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
{dirty && (
|
||||
|
||||
Reference in new issue
Block a user