From 2c699f31dead16153de7e34a3a9d86353f4802ca Mon Sep 17 00:00:00 2001 From: Duy /zuey/ Date: Sat, 23 May 2026 14:02:34 +0700 Subject: [PATCH] fix(chat): support zero debounce and agent overrides --- cmd/gateway_consumer.go | 13 ++-- cmd/gateway_consumer_debounce.go | 62 ++++++++++++++++ cmd/gateway_system_config_sync.go | 5 +- cmd/gateway_system_config_sync_test.go | 48 ++++++++++++ docs/05-channels-messaging.md | 2 +- docs/19-websocket-rpc.md | 2 +- docs/project-changelog.md | 2 + internal/bus/inbound_debounce.go | 36 ++++++--- internal/bus/inbound_debounce_test.go | 44 ++++++++++- internal/config/config_channels.go | 2 +- internal/gateway/methods/chat.go | 2 +- internal/gateway/methods/chat_debounce.go | 11 ++- .../gateway/methods/chat_debounce_test.go | 23 +++--- internal/gateway/methods/config_patch_test.go | 73 +++++++++++++++++++ internal/store/agent_inbound_debounce_test.go | 28 +++++++ internal/store/agent_store.go | 22 ++++++ ui/web/src/i18n/locales/en/agents.json | 10 +++ ui/web/src/i18n/locales/en/config.json | 2 +- ui/web/src/i18n/locales/vi/agents.json | 10 +++ ui/web/src/i18n/locales/vi/config.json | 2 +- ui/web/src/i18n/locales/zh/agents.json | 10 +++ ui/web/src/i18n/locales/zh/config.json | 2 +- .../agent-detail/agent-advanced-dialog.tsx | 14 +++- .../agent-advanced-state-utils.ts | 39 +++++++++- .../inbound-debounce-section.tsx | 63 ++++++++++++++++ .../agent-detail/config-sections/index.ts | 1 + .../config/sections/behavior-rate-card.tsx | 4 +- .../config/sections/behavior-section.tsx | 14 +--- ui/web/src/types/agent.ts | 6 ++ 29 files changed, 498 insertions(+), 54 deletions(-) create mode 100644 cmd/gateway_consumer_debounce.go create mode 100644 cmd/gateway_system_config_sync_test.go create mode 100644 internal/gateway/methods/config_patch_test.go create mode 100644 internal/store/agent_inbound_debounce_test.go create mode 100644 ui/web/src/pages/agents/agent-detail/config-sections/inbound-debounce-section.tsx diff --git a/cmd/gateway_consumer.go b/cmd/gateway_consumer.go index cbca58e3..aa6d8469 100644 --- a/cmd/gateway_consumer.go +++ b/cmd/gateway_consumer.go @@ -82,19 +82,17 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents // Inbound debounce: merge rapid messages from the same sender before processing. // Matching TS createInboundDebouncer from src/auto-reply/inbound-debounce.ts. - debounceMs := cfg.Gateway.InboundDebounceMs - if debounceMs == 0 { - debounceMs = 1000 // default: 1000ms - } - debouncer := bus.NewInboundDebouncer( - time.Duration(debounceMs)*time.Millisecond, + debouncer := bus.NewInboundDebouncerFunc( + func(msg bus.InboundMessage) time.Duration { + return resolveInboundDebounceDelay(ctx, msg, deps) + }, func(msg bus.InboundMessage) { processNormalMessage(ctx, msg, deps) }, ) defer debouncer.Stop() - slog.Info("inbound debounce configured", "debounce_ms", debounceMs) + slog.Info("inbound debounce configured", "global_debounce_ms", cfg.Gateway.InboundDebounceMs, "agent_override", true) // Track background goroutines (subagent announces, teammate messages) // so shutdown can wait for in-flight work to complete. @@ -139,6 +137,7 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents } // --- Normal messages: route through debouncer --- + prepareInboundDebounceMessage(&msg, deps) debouncer.Push(msg) } } diff --git a/cmd/gateway_consumer_debounce.go b/cmd/gateway_consumer_debounce.go new file mode 100644 index 00000000..a0c98419 --- /dev/null +++ b/cmd/gateway_consumer_debounce.go @@ -0,0 +1,62 @@ +package cmd + +import ( + "context" + "log/slog" + "time" + + "github.com/google/uuid" + + "github.com/nextlevelbuilder/goclaw/internal/bus" + "github.com/nextlevelbuilder/goclaw/internal/store" +) + +func prepareInboundDebounceMessage(msg *bus.InboundMessage, deps *ConsumerDeps) { + if msg == nil || deps == nil || deps.Cfg == nil || msg.AgentID != "" { + return + } + msg.AgentID = resolveAgentRoute(deps.Cfg, msg.Channel, msg.ChatID, msg.PeerKind) +} + +func resolveInboundDebounceDelay(ctx context.Context, msg bus.InboundMessage, deps *ConsumerDeps) time.Duration { + debounceMs := 0 + if deps != nil && deps.Cfg != nil { + debounceMs = deps.Cfg.Gateway.InboundDebounceMs + } + if deps == nil || deps.AgentStore == nil || msg.AgentID == "" { + return inboundDebounceDuration(debounceMs) + } + + agentCtx := ctx + if msg.TenantID != uuid.Nil { + agentCtx = store.WithTenantID(agentCtx, msg.TenantID) + } else { + agentCtx = store.WithTenantID(agentCtx, store.MasterTenantID) + } + + agentData, err := getInboundDebounceAgent(agentCtx, deps.AgentStore, msg.AgentID) + if err != nil || agentData == nil { + if err != nil { + slog.Debug("inbound debounce: agent config unavailable", "agent", msg.AgentID, "error", err) + } + return inboundDebounceDuration(debounceMs) + } + if overrideMs, ok := agentData.ParseInboundDebounceMs(); ok { + debounceMs = overrideMs + } + return inboundDebounceDuration(debounceMs) +} + +func getInboundDebounceAgent(ctx context.Context, agentStore store.AgentStore, agentID string) (*store.AgentData, error) { + if parsed, err := uuid.Parse(agentID); err == nil && parsed != uuid.Nil { + return agentStore.GetByID(ctx, parsed) + } + return agentStore.GetByKey(ctx, agentID) +} + +func inboundDebounceDuration(ms int) time.Duration { + if ms <= 0 { + return 0 + } + return time.Duration(ms) * time.Millisecond +} diff --git a/cmd/gateway_system_config_sync.go b/cmd/gateway_system_config_sync.go index dbff536f..bcfd5100 100644 --- a/cmd/gateway_system_config_sync.go +++ b/cmd/gateway_system_config_sync.go @@ -78,6 +78,9 @@ func seedConfigForContext(ctx context.Context, sc store.SystemConfigStore, cfg * set(key, fmt.Sprintf("%d", val)) } } + setIntAllowZero := func(key string, val int) { + set(key, fmt.Sprintf("%d", val)) + } setBool := func(key string, val *bool) { if val != nil { set(key, fmt.Sprintf("%t", *val)) @@ -102,7 +105,7 @@ func seedConfigForContext(ctx context.Context, sc store.SystemConfigStore, cfg * setInt("gateway.rate_limit_rpm", cfg.Gateway.RateLimitRPM) setInt("gateway.max_message_chars", cfg.Gateway.MaxMessageChars) set("gateway.injection_action", cfg.Gateway.InjectionAction) - setInt("gateway.inbound_debounce_ms", cfg.Gateway.InboundDebounceMs) + setIntAllowZero("gateway.inbound_debounce_ms", cfg.Gateway.InboundDebounceMs) setBool("gateway.block_reply", cfg.Gateway.BlockReply) setBool("gateway.tool_status", cfg.Gateway.ToolStatus) setInt("gateway.task_recovery_interval_sec", cfg.Gateway.TaskRecoveryIntervalSec) diff --git a/cmd/gateway_system_config_sync_test.go b/cmd/gateway_system_config_sync_test.go new file mode 100644 index 00000000..3b9e2b2c --- /dev/null +++ b/cmd/gateway_system_config_sync_test.go @@ -0,0 +1,48 @@ +package cmd + +import ( + "context" + "maps" + "testing" + + "github.com/nextlevelbuilder/goclaw/internal/config" + "github.com/nextlevelbuilder/goclaw/internal/store" +) + +type captureSystemConfigStore struct { + data map[string]string +} + +func (s *captureSystemConfigStore) Get(_ context.Context, key string) (string, error) { + return s.data[key], nil +} + +func (s *captureSystemConfigStore) Set(_ context.Context, key, value string) error { + s.data[key] = value + return nil +} + +func (s *captureSystemConfigStore) Delete(_ context.Context, key string) error { + delete(s.data, key) + return nil +} + +func (s *captureSystemConfigStore) List(_ context.Context) (map[string]string, error) { + out := make(map[string]string, len(s.data)) + maps.Copy(out, s.data) + return out, nil +} + +func TestSeedConfigForContextPersistsZeroInboundDebounce(t *testing.T) { + t.Parallel() + + sc := &captureSystemConfigStore{data: map[string]string{}} + cfg := config.Default() + cfg.Gateway.InboundDebounceMs = 0 + + seedConfigForContext(store.WithTenantID(context.Background(), store.MasterTenantID), sc, cfg, false) + + if got := sc.data["gateway.inbound_debounce_ms"]; got != "0" { + t.Fatalf("gateway.inbound_debounce_ms = %q, want 0", got) + } +} diff --git a/docs/05-channels-messaging.md b/docs/05-channels-messaging.md index c0137f1b..6e70ae68 100644 --- a/docs/05-channels-messaging.md +++ b/docs/05-channels-messaging.md @@ -69,7 +69,7 @@ The consumer routes system messages based on sender ID prefixes: ### Inbound Debounce -Normal channel messages pass through the shared inbound debouncer before agent execution. `gateway.inbound_debounce_ms` merges rapid text messages from the same `channel:chatID:senderID`; `0` uses the 1000ms default and `-1` disables debounce. Media messages bypass the wait window after flushing pending text, and command/control messages such as stop/reset and system escalations bypass debounce. +Normal channel messages pass through the shared inbound debouncer before agent execution. `gateway.inbound_debounce_ms` merges rapid text messages from the same `channel:chatID:senderID:agentID`; `0` means no debounce and positive values set the wait window. Agents can override the global value with `other_config.inbound_debounce_ms`; unset inherits the global config. Media messages bypass the wait window after flushing pending text, and command/control messages such as stop/reset and system escalations bypass debounce. --- diff --git a/docs/19-websocket-rpc.md b/docs/19-websocket-rpc.md index 25e9610e..c1d77d57 100644 --- a/docs/19-websocket-rpc.md +++ b/docs/19-websocket-rpc.md @@ -106,7 +106,7 @@ Send a message to an agent and trigger execution. When `stream: true`, intermediate events are emitted: `chunk`, `tool.call`, `tool.result`, `run.started`, `run.completed`. -Rapid text-only `chat.send` requests for the same user and session are debounced by `gateway.inbound_debounce_ms`: `0` uses the 1000ms default and `-1` disables debounce. The merged message keeps request params from the latest send and joins text with newlines. Cancel keywords bypass debounce and abort the active run immediately. Media sends bypass the wait window and drain any pending text into the same dispatch. +Rapid text-only `chat.send` requests for the same user and session are debounced by `gateway.inbound_debounce_ms`: `0` means no debounce and positive values set the wait window. Agents can override the global value with `other_config.inbound_debounce_ms`; unset inherits the global config. The merged message keeps request params from the latest send and joins text with newlines. Cancel keywords bypass debounce and abort the active run immediately. Media sends bypass the wait window and drain any pending text into the same dispatch. ### `chat.history` diff --git a/docs/project-changelog.md b/docs/project-changelog.md index bcdb256e..b1702a21 100644 --- a/docs/project-changelog.md +++ b/docs/project-changelog.md @@ -10,11 +10,13 @@ Significant changes, features, and fixes in reverse chronological order. **Features** +- Added per-agent inbound debounce override via `other_config.inbound_debounce_ms`; unset inherits the global gateway setting. - Added Web Chat debounce for rapid text-only `chat.send` calls using `gateway.inbound_debounce_ms`. - Clarified shared inbound debounce behavior in docs and Web UI config help text. **Fixes** +- Fixed inbound debounce semantics so `gateway.inbound_debounce_ms=0` means no debounce and positive values set the wait window. - Fixed Slack `debounce_delay: 0` so it disables per-thread batching instead of falling back to the default. **Tests** diff --git a/internal/bus/inbound_debounce.go b/internal/bus/inbound_debounce.go index 6506ce7e..8df92cce 100644 --- a/internal/bus/inbound_debounce.go +++ b/internal/bus/inbound_debounce.go @@ -16,10 +16,10 @@ import ( // InboundDebouncer buffers rapid inbound messages from the same sender // and merges them into a single message before calling flushFn. type InboundDebouncer struct { - debounceMs time.Duration - mu sync.Mutex - buffers map[string]*debounceBuffer - flushFn func(InboundMessage) + delayFn func(InboundMessage) time.Duration + mu sync.Mutex + buffers map[string]*debounceBuffer + flushFn func(InboundMessage) } type debounceBuffer struct { @@ -30,18 +30,30 @@ type debounceBuffer struct { // NewInboundDebouncer creates a debouncer with the given window and flush callback. // If debounceMs <= 0, messages are passed through immediately (debouncing disabled). func NewInboundDebouncer(debounceMs time.Duration, flushFn func(InboundMessage)) *InboundDebouncer { + return NewInboundDebouncerFunc(func(InboundMessage) time.Duration { + return debounceMs + }, flushFn) +} + +// NewInboundDebouncerFunc creates a debouncer whose window can vary per message. +func NewInboundDebouncerFunc(delayFn func(InboundMessage) time.Duration, flushFn func(InboundMessage)) *InboundDebouncer { + if delayFn == nil { + delayFn = func(InboundMessage) time.Duration { return 0 } + } return &InboundDebouncer{ - debounceMs: debounceMs, - buffers: make(map[string]*debounceBuffer), - flushFn: flushFn, + delayFn: delayFn, + buffers: make(map[string]*debounceBuffer), + flushFn: flushFn, } } // Push adds a message to the debounce buffer. // If debouncing is disabled or the message should bypass (media), it is flushed immediately. func (d *InboundDebouncer) Push(msg InboundMessage) { + debounceMs := d.delayFn(msg) + // Disabled: pass through immediately. - if d.debounceMs <= 0 { + if debounceMs <= 0 { d.flushFn(msg) return } @@ -70,13 +82,13 @@ func (d *InboundDebouncer) Push(msg InboundMessage) { if buf.timer != nil { buf.timer.Stop() } - buf.timer = time.AfterFunc(d.debounceMs, func() { + buf.timer = time.AfterFunc(debounceMs, func() { d.flushKey(key) }) if len(buf.messages) == 1 { slog.Debug("inbound debounce: buffering", - "key", key, "debounce_ms", d.debounceMs.Milliseconds()) + "key", key, "debounce_ms", debounceMs.Milliseconds()) } else { slog.Debug("inbound debounce: message appended", "key", key, "buffered", len(buf.messages)) @@ -127,9 +139,9 @@ func (d *InboundDebouncer) flushKey(key string) { d.flushFn(merged) } -// debounceKey builds the buffer key: channel:chatID:senderID. +// debounceKey builds the buffer key: channel:chatID:senderID:agentID. func debounceKey(msg InboundMessage) string { - return msg.Channel + ":" + msg.ChatID + ":" + msg.SenderID + return msg.Channel + ":" + msg.ChatID + ":" + msg.SenderID + ":" + msg.AgentID } // mergeInboundMessages combines multiple messages into one. diff --git a/internal/bus/inbound_debounce_test.go b/internal/bus/inbound_debounce_test.go index b68ac8f6..59b2b1cf 100644 --- a/internal/bus/inbound_debounce_test.go +++ b/internal/bus/inbound_debounce_test.go @@ -26,7 +26,7 @@ func TestInboundDebouncerMergesRapidText(t *testing.T) { func TestInboundDebouncerDisabledPassesThrough(t *testing.T) { out := make(chan InboundMessage, 2) - d := NewInboundDebouncer(-1, func(msg InboundMessage) { + d := NewInboundDebouncer(0, func(msg InboundMessage) { out <- msg }) @@ -41,6 +41,48 @@ func TestInboundDebouncerDisabledPassesThrough(t *testing.T) { } } +func TestInboundDebouncerDynamicDelay(t *testing.T) { + out := make(chan InboundMessage, 2) + d := NewInboundDebouncerFunc(func(msg InboundMessage) time.Duration { + if msg.AgentID == "instant" { + return 0 + } + return 20 * time.Millisecond + }, func(msg InboundMessage) { + out <- msg + }) + defer d.Stop() + + d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", AgentID: "instant", Content: "one"}) + d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", AgentID: "debounced", Content: "two"}) + d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", AgentID: "debounced", Content: "three"}) + + if got := waitInbound(t, out); got.Content != "one" { + t.Fatalf("instant content = %q", got.Content) + } + if got := waitInbound(t, out); got.Content != "two\nthree" { + t.Fatalf("debounced content = %q", got.Content) + } +} + +func TestInboundDebouncerSeparatesAgents(t *testing.T) { + out := make(chan InboundMessage, 2) + d := NewInboundDebouncer(20*time.Millisecond, func(msg InboundMessage) { + out <- msg + }) + defer d.Stop() + + d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", AgentID: "agent-a", Content: "a"}) + d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", AgentID: "agent-b", Content: "b"}) + + first := waitInbound(t, out) + second := waitInbound(t, out) + got := map[string]string{first.AgentID: first.Content, second.AgentID: second.Content} + if got["agent-a"] != "a" || got["agent-b"] != "b" { + t.Fatalf("agent buffers = %#v, want separate flushes", got) + } +} + func TestInboundDebouncerMediaFlushesPendingTextFirst(t *testing.T) { out := make(chan InboundMessage, 2) d := NewInboundDebouncer(time.Minute, func(msg InboundMessage) { diff --git a/internal/config/config_channels.go b/internal/config/config_channels.go index 9cd8167e..a4f3d401 100644 --- a/internal/config/config_channels.go +++ b/internal/config/config_channels.go @@ -370,7 +370,7 @@ type GatewayConfig struct { MaxMessageChars int `json:"max_message_chars,omitempty"` // max user message characters (default 32000) RateLimitRPM int `json:"rate_limit_rpm,omitempty"` // rate limit: requests per minute per user (default 20, 0 = disabled) InjectionAction string `json:"injection_action,omitempty"` // prompt injection action: "log", "warn" (default), "block", "off" - InboundDebounceMs int `json:"inbound_debounce_ms,omitempty"` // merge rapid channel/Web Chat messages from same sender/session (default 1000ms, -1 = disabled) + InboundDebounceMs int `json:"inbound_debounce_ms,omitempty"` // merge rapid channel/Web Chat messages from same sender/session (0 = no wait) Quota *QuotaConfig `json:"quota,omitempty"` // per-user/group request quotas BlockReply *bool `json:"block_reply,omitempty"` // deliver intermediate text during tool iterations (default false) ToolStatus *bool `json:"tool_status,omitempty"` // show tool name in streaming preview during tool execution (default true) diff --git a/internal/gateway/methods/chat.go b/internal/gateway/methods/chat.go index 1281a9c5..3ed1f2a1 100644 --- a/internal/gateway/methods/chat.go +++ b/internal/gateway/methods/chat.go @@ -217,7 +217,7 @@ func (m *ChatMethods) handleSend(ctx context.Context, client *gateway.Client, re m.dispatchChatSends(append(pending, item)) return } - if delay := chatDebounceDelay(m.cfg); delay > 0 { + if delay := chatDebounceDelay(m.cfg, loop.OtherConfig()); delay > 0 { m.debouncer.Push(debounceKey, delay, item) return } diff --git a/internal/gateway/methods/chat_debounce.go b/internal/gateway/methods/chat_debounce.go index 8a62bf3e..513a405c 100644 --- a/internal/gateway/methods/chat_debounce.go +++ b/internal/gateway/methods/chat_debounce.go @@ -2,6 +2,7 @@ package methods import ( "context" + "encoding/json" "strings" "sync" "time" @@ -9,6 +10,7 @@ import ( "github.com/nextlevelbuilder/goclaw/internal/agent" "github.com/nextlevelbuilder/goclaw/internal/config" "github.com/nextlevelbuilder/goclaw/internal/gateway" + "github.com/nextlevelbuilder/goclaw/internal/store" ) type chatSendRequest struct { @@ -127,13 +129,16 @@ func mergeChatSendRequests(items []chatSendRequest) chatSendParams { return last } -func chatDebounceDelay(cfg *config.Config) time.Duration { +func chatDebounceDelay(cfg *config.Config, agentOtherConfig json.RawMessage) time.Duration { debounceMs := 0 if cfg != nil { debounceMs = cfg.Gateway.InboundDebounceMs } - if debounceMs == 0 { - debounceMs = 1000 + if overrideMs, ok := store.ParseInboundDebounceMsFromOtherConfig(agentOtherConfig); ok { + debounceMs = overrideMs + } + if debounceMs <= 0 { + return 0 } return time.Duration(debounceMs) * time.Millisecond } diff --git a/internal/gateway/methods/chat_debounce_test.go b/internal/gateway/methods/chat_debounce_test.go index 1644e0dd..47226f10 100644 --- a/internal/gateway/methods/chat_debounce_test.go +++ b/internal/gateway/methods/chat_debounce_test.go @@ -71,18 +71,23 @@ func TestChatDebouncerDiscardDropsPendingBeforeCancel(t *testing.T) { assertNoChatDebounceFlush(t, out) } -func TestChatDebounceDelayDefaultAndDisabled(t *testing.T) { - if got := chatDebounceDelay(&config.Config{}); got != time.Second { - t.Fatalf("default debounce = %s, want 1s", got) +func TestChatDebounceDelayGlobalAndAgentOverride(t *testing.T) { + if got := chatDebounceDelay(&config.Config{}, nil); got != 0 { + t.Fatalf("default debounce = %s, want disabled", got) } cfg := &config.Config{} - cfg.Gateway.InboundDebounceMs = -1 - if got := chatDebounceDelay(cfg); got >= 0 { - t.Fatalf("disabled debounce = %s, want negative duration", got) - } cfg.Gateway.InboundDebounceMs = 250 - if got := chatDebounceDelay(cfg); got != 250*time.Millisecond { - t.Fatalf("custom debounce = %s, want 250ms", got) + if got := chatDebounceDelay(cfg, nil); got != 250*time.Millisecond { + t.Fatalf("global debounce = %s, want 250ms", got) + } + if got := chatDebounceDelay(cfg, []byte(`{"inbound_debounce_ms":0}`)); got != 0 { + t.Fatalf("agent disabled debounce = %s, want disabled", got) + } + if got := chatDebounceDelay(cfg, []byte(`{"inbound_debounce_ms":500}`)); got != 500*time.Millisecond { + t.Fatalf("agent custom debounce = %s, want 500ms", got) + } + if got := chatDebounceDelay(cfg, []byte(`{"other":true}`)); got != 250*time.Millisecond { + t.Fatalf("agent inherit debounce = %s, want 250ms", got) } } diff --git a/internal/gateway/methods/config_patch_test.go b/internal/gateway/methods/config_patch_test.go new file mode 100644 index 00000000..cd86e722 --- /dev/null +++ b/internal/gateway/methods/config_patch_test.go @@ -0,0 +1,73 @@ +package methods + +import ( + "bytes" + "context" + "encoding/json" + "os" + "path/filepath" + "testing" + "time" + + "github.com/nextlevelbuilder/goclaw/internal/config" + "github.com/nextlevelbuilder/goclaw/internal/gateway" + "github.com/nextlevelbuilder/goclaw/internal/permissions" + "github.com/nextlevelbuilder/goclaw/internal/store" + "github.com/nextlevelbuilder/goclaw/pkg/protocol" +) + +func TestConfigPatchPersistsInboundDebounceMs(t *testing.T) { + t.Parallel() + + cfg := config.Default() + cfgPath := filepath.Join(t.TempDir(), "config.json") + methods := NewConfigMethods(cfg, cfgPath, nil, nil) + client, responses := gateway.NewCapturingTestClient(permissions.RoleOwner, store.MasterTenantID, "owner", 1) + params, err := json.Marshal(map[string]string{ + "raw": `{"gateway":{"inbound_debounce_ms":500}}`, + }) + if err != nil { + t.Fatal(err) + } + + methods.handlePatch( + store.WithTenantID(context.Background(), store.MasterTenantID), + client, + &protocol.RequestFrame{ + Type: protocol.FrameTypeRequest, + ID: "patch-inbound-debounce", + Method: protocol.MethodConfigPatch, + Params: params, + }, + ) + + res := readConfigPatchResponse(t, responses) + if !res.OK { + t.Fatalf("config.patch failed: %#v", res.Error) + } + if cfg.Gateway.InboundDebounceMs != 500 { + t.Fatalf("in-memory inbound_debounce_ms = %d, want 500", cfg.Gateway.InboundDebounceMs) + } + data, err := os.ReadFile(cfgPath) + if err != nil { + t.Fatal(err) + } + if !bytes.Contains(data, []byte(`"inbound_debounce_ms": 500`)) { + t.Fatalf("saved config missing inbound_debounce_ms=500:\n%s", data) + } +} + +func readConfigPatchResponse(t *testing.T, responses <-chan []byte) protocol.ResponseFrame { + t.Helper() + select { + case raw := <-responses: + var res protocol.ResponseFrame + if err := json.Unmarshal(raw, &res); err != nil { + t.Fatal(err) + } + return res + case <-time.After(500 * time.Millisecond): + t.Fatal("timed out waiting for config.patch response") + return protocol.ResponseFrame{} + } +} diff --git a/internal/store/agent_inbound_debounce_test.go b/internal/store/agent_inbound_debounce_test.go new file mode 100644 index 00000000..20ec30f9 --- /dev/null +++ b/internal/store/agent_inbound_debounce_test.go @@ -0,0 +1,28 @@ +package store + +import "testing" + +func TestParseInboundDebounceMsFromOtherConfig(t *testing.T) { + t.Parallel() + + cases := []struct { + name string + raw []byte + want int + ok bool + }{ + {name: "missing", raw: []byte(`{"prompt_mode":"full"}`), ok: false}, + {name: "zero", raw: []byte(`{"inbound_debounce_ms":0}`), want: 0, ok: true}, + {name: "positive", raw: []byte(`{"inbound_debounce_ms":500}`), want: 500, ok: true}, + {name: "malformed", raw: []byte(`{"inbound_debounce_ms":`), ok: false}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + got, ok := ParseInboundDebounceMsFromOtherConfig(tc.raw) + if ok != tc.ok || got != tc.want { + t.Fatalf("ParseInboundDebounceMsFromOtherConfig() = (%d, %t), want (%d, %t)", got, ok, tc.want, tc.ok) + } + }) + } +} diff --git a/internal/store/agent_store.go b/internal/store/agent_store.go index a8ed8b04..c4409099 100644 --- a/internal/store/agent_store.go +++ b/internal/store/agent_store.go @@ -253,6 +253,28 @@ func (a *AgentData) ParseAllowImageGeneration() bool { return *bag.AllowImageGeneration } +// ParseInboundDebounceMs returns the per-agent inbound debounce override. +// Missing or malformed config means "inherit global gateway.inbound_debounce_ms". +func (a *AgentData) ParseInboundDebounceMs() (int, bool) { + if a == nil { + return 0, false + } + return ParseInboundDebounceMsFromOtherConfig(a.OtherConfig) +} + +func ParseInboundDebounceMsFromOtherConfig(raw json.RawMessage) (int, bool) { + if len(raw) <= 2 { + return 0, false + } + var bag struct { + InboundDebounceMs *int `json:"inbound_debounce_ms"` + } + if json.Unmarshal(raw, &bag) != nil || bag.InboundDebounceMs == nil { + return 0, false + } + return *bag.InboundDebounceMs, true +} + // validPromptModes is the set of allowed prompt_mode values. var validPromptModes = map[string]bool{ "full": true, "task": true, "minimal": true, "none": true, diff --git a/ui/web/src/i18n/locales/en/agents.json b/ui/web/src/i18n/locales/en/agents.json index a0d28566..547c5cbf 100644 --- a/ui/web/src/i18n/locales/en/agents.json +++ b/ui/web/src/i18n/locales/en/agents.json @@ -815,6 +815,16 @@ "memoryFlush": "Memory Flush", "memoryFlushTip": "Before compaction, the agent gets a turn to save important context to memory files. Also triggers Knowledge Graph extraction." }, + "inboundDebounce": { + "title": "Inbound Debounce", + "description": "Override rapid Web Chat and channel message batching for this agent.", + "mode": "Mode", + "modeTip": "Inherit uses the global Behavior setting. Custom stores an override on this agent.", + "inherit": "Inherit Global", + "custom": "Custom", + "debounceMs": "Debounce (ms)", + "debounceMsTip": "Set 0 for no debounce; positive values merge follow-ups after that wait." + }, "contextPruning": { "title": "Context Pruning", "description": "Trim old tool results to save context window", diff --git a/ui/web/src/i18n/locales/en/config.json b/ui/web/src/i18n/locales/en/config.json index c345a62c..e130312c 100644 --- a/ui/web/src/i18n/locales/en/config.json +++ b/ui/web/src/i18n/locales/en/config.json @@ -42,7 +42,7 @@ "gateway.rateLimitRpm": "Rate Limit (RPM)", "gateway.rateLimitRpmTip": "Maximum requests per minute per user. Set to 0 to disable rate limiting.", "gateway.inboundDebounceMs": "Inbound Debounce (ms)", - "gateway.inboundDebounceMsTip": "Delay in milliseconds before processing rapid channel or Web Chat follow-ups. Set 0 for the 1000ms default, or -1 to disable.", + "gateway.inboundDebounceMsTip": "Delay in milliseconds before processing rapid channel or Web Chat follow-ups. Set 0 for no debounce; positive values merge follow-ups after that wait.", "gateway.injectionAction": "Injection Detection Action", "agents.title": "Agent Defaults", diff --git a/ui/web/src/i18n/locales/vi/agents.json b/ui/web/src/i18n/locales/vi/agents.json index 8bb5c2a5..9bac61c0 100644 --- a/ui/web/src/i18n/locales/vi/agents.json +++ b/ui/web/src/i18n/locales/vi/agents.json @@ -800,6 +800,16 @@ "memoryFlush": "Ghi nhớ trước nén", "memoryFlushTip": "Trước khi nén, agent được một lượt để lưu ngữ cảnh quan trọng vào file bộ nhớ. Cũng kích hoạt trích xuất Knowledge Graph." }, + "inboundDebounce": { + "title": "Debounce đầu vào", + "description": "Override việc gộp tin nhắn nhanh liên tiếp từ Web Chat và channel cho agent này.", + "mode": "Chế độ", + "modeTip": "Inherit dùng cấu hình global trong Behavior. Custom lưu override trên agent này.", + "inherit": "Kế thừa global", + "custom": "Tùy chỉnh", + "debounceMs": "Debounce (ms)", + "debounceMsTip": "Đặt 0 để không debounce; giá trị dương sẽ gộp tin nhắn sau thời gian chờ đó." + }, "contextPruning": { "title": "Cắt bớt ngữ cảnh", "description": "Cắt bỏ kết quả công cụ cũ để tiết kiệm cửa sổ ngữ cảnh", diff --git a/ui/web/src/i18n/locales/vi/config.json b/ui/web/src/i18n/locales/vi/config.json index 87da884d..7a5ab1ba 100644 --- a/ui/web/src/i18n/locales/vi/config.json +++ b/ui/web/src/i18n/locales/vi/config.json @@ -42,7 +42,7 @@ "gateway.rateLimitRpm": "Giới hạn tốc độ (RPM)", "gateway.rateLimitRpmTip": "Số yêu cầu tối đa mỗi phút mỗi người dùng. Đặt 0 để tắt giới hạn.", "gateway.inboundDebounceMs": "Debounce đầu vào (ms)", - "gateway.inboundDebounceMsTip": "Độ trễ tính bằng mili giây trước khi xử lý các tin nhắn nhanh liên tiếp từ channel hoặc Web Chat. Đặt 0 để dùng mặc định 1000ms, hoặc -1 để tắt.", + "gateway.inboundDebounceMsTip": "Độ trễ tính bằng mili giây trước khi xử lý các tin nhắn nhanh liên tiếp từ channel hoặc Web Chat. Đặt 0 để không debounce; giá trị dương sẽ gộp tin nhắn sau thời gian chờ đó.", "gateway.injectionAction": "Hành động phát hiện injection", "agents.title": "Mặc định agent", diff --git a/ui/web/src/i18n/locales/zh/agents.json b/ui/web/src/i18n/locales/zh/agents.json index fef22a44..48c7ead0 100644 --- a/ui/web/src/i18n/locales/zh/agents.json +++ b/ui/web/src/i18n/locales/zh/agents.json @@ -800,6 +800,16 @@ "memoryFlush": "压缩前记忆", "memoryFlushTip": "压缩前,Agent可以将重要上下文保存到记忆文件。同时触发知识图谱提取。" }, + "inboundDebounce": { + "title": "入站防抖", + "description": "为此 Agent 覆盖 Web Chat 和频道快速连续消息的合并设置。", + "mode": "模式", + "modeTip": "继承会使用全局 Behavior 设置。自定义会在此 Agent 上保存覆盖值。", + "inherit": "继承全局", + "custom": "自定义", + "debounceMs": "防抖(毫秒)", + "debounceMsTip": "设为 0 表示不防抖;正数会在该等待时间后合并连续消息。" + }, "contextPruning": { "title": "上下文裁剪", "description": "裁剪旧工具结果以节省上下文窗口", diff --git a/ui/web/src/i18n/locales/zh/config.json b/ui/web/src/i18n/locales/zh/config.json index af808a73..d5541ea0 100644 --- a/ui/web/src/i18n/locales/zh/config.json +++ b/ui/web/src/i18n/locales/zh/config.json @@ -42,7 +42,7 @@ "gateway.rateLimitRpm": "速率限制(RPM)", "gateway.rateLimitRpmTip": "每用户每分钟最大请求数。设为 0 禁用速率限制。", "gateway.inboundDebounceMs": "入站防抖(毫秒)", - "gateway.inboundDebounceMsTip": "处理来自频道或 Web Chat 的快速连续消息前的延迟毫秒数。设为 0 使用默认 1000ms,设为 -1 禁用。", + "gateway.inboundDebounceMsTip": "处理来自频道或 Web Chat 的快速连续消息前的延迟毫秒数。设为 0 表示不防抖;正数会在该等待时间后合并连续消息。", "gateway.injectionAction": "注入检测动作", "agents.title": "Agent 默认值", diff --git a/ui/web/src/pages/agents/agent-detail/agent-advanced-dialog.tsx b/ui/web/src/pages/agents/agent-detail/agent-advanced-dialog.tsx index c1fd3196..96becf72 100644 --- a/ui/web/src/pages/agents/agent-detail/agent-advanced-dialog.tsx +++ b/ui/web/src/pages/agents/agent-detail/agent-advanced-dialog.tsx @@ -13,7 +13,7 @@ import type { } from "@/types/agent"; import { ChatGPTOAuthRoutingSection, ThinkingSection, WorkspaceSharingSection, CompactionSection, - ContextPruningSection, ModelFallbackSection, SandboxSection, + ContextPruningSection, InboundDebounceSection, ModelFallbackSection, SandboxSection, } from "./config-sections"; import { WorkspaceSection } from "./general-sections"; import { useProviders } from "@/pages/providers/hooks/use-providers"; @@ -57,6 +57,8 @@ export function AgentAdvancedDialog({ open, onOpenChange, agent, onUpdate }: Age const [chatgptRouting, setChatgptRouting] = useState(init.chatgptRouting); const [modelFallback, setModelFallback] = useState(init.modelFallback); const [comp, setComp] = useState(init.comp); + const [inboundDebounceMode, setInboundDebounceMode] = useState(init.inboundDebounceMode); + const [inboundDebounceMs, setInboundDebounceMs] = useState(init.inboundDebounceMs); const [pruneEnabled, setPruneEnabled] = useState(init.pruneEnabled); const [prune, setPrune] = useState(init.prune); const [sbEnabled, setSbEnabled] = useState(init.sbEnabled); @@ -76,6 +78,8 @@ export function AgentAdvancedDialog({ open, onOpenChange, agent, onUpdate }: Age setModelFallback(s.modelFallback); setWsSharing(s.wsSharing); setComp(s.comp); + setInboundDebounceMode(s.inboundDebounceMode); + setInboundDebounceMs(s.inboundDebounceMs); setPruneEnabled(s.pruneEnabled); setPrune(s.prune); setSbEnabled(s.sbEnabled); @@ -129,6 +133,8 @@ export function AgentAdvancedDialog({ open, onOpenChange, agent, onUpdate }: Age modelFallback, wsSharing, comp, + inboundDebounceMode, + inboundDebounceMs, pruneEnabled, prune, sbEnabled, @@ -229,6 +235,12 @@ export function AgentAdvancedDialog({ open, onOpenChange, agent, onUpdate }: Age />
+ ; + const rawInboundDebounceMs = otherConfig.inbound_debounce_ms; + const inboundDebounceMs = + typeof rawInboundDebounceMs === "number" && Number.isFinite(rawInboundDebounceMs) + ? Math.max(0, Math.trunc(rawInboundDebounceMs)) + : undefined; return { reasoningMode, @@ -100,6 +109,8 @@ export function deriveState( {} ) as WorkspaceSharingConfig, comp: agent.compaction_config ?? {}, + inboundDebounceMode: inboundDebounceMs === undefined ? "inherit" : "custom", + inboundDebounceMs: inboundDebounceMs ?? 0, pruneEnabled: agent.context_pruning?.mode === "cache-ttl", prune: agent.context_pruning ?? {}, sbEnabled: agent.sandbox_config != null, @@ -122,6 +133,8 @@ export interface BuildAdvancedUpdatePayloadParams { modelFallback: ModelFallbackConfig; wsSharing: WorkspaceSharingConfig; comp: CompactionConfig; + inboundDebounceMode: InboundDebounceOverrideMode; + inboundDebounceMs: number; pruneEnabled: boolean; prune: ContextPruningConfig; sbEnabled: boolean; @@ -135,7 +148,8 @@ export function buildAdvancedUpdatePayload( agent, currentProvider, providersLoading, providerModelsLoading, expertReasoningAvailable, reasoningMode, reasoningEffort, reasoningExpert, reasoningFallback, thinkingLevel, chatgptRouting, wsSharing, - modelFallback, comp, pruneEnabled, prune, sbEnabled, sb, + modelFallback, comp, inboundDebounceMode, inboundDebounceMs, + pruneEnabled, prune, sbEnabled, sb, } = params; const routingPayload = buildAgentOtherConfigWithChatGPTOAuthRouting( @@ -156,6 +170,11 @@ export function buildAdvancedUpdatePayload( model_fallback: normalizeModelFallbackForPayload(modelFallback), ...routingPayload, }; + updates.other_config = buildOtherConfigWithInboundDebounce( + updates.other_config, + inboundDebounceMode, + inboundDebounceMs, + ); // Build reasoning_config and thinking_level as top-level fields if (reasoningMode === "inherit") { @@ -192,6 +211,24 @@ export function buildAdvancedUpdatePayload( return updates; } +function buildOtherConfigWithInboundDebounce( + base: unknown, + mode: InboundDebounceOverrideMode, + debounceMs: number, +): Record | null { + const bag = isPlainObject(base) ? { ...base } : {}; + if (mode === "inherit") { + delete bag.inbound_debounce_ms; + } else { + bag.inbound_debounce_ms = Math.max(0, Math.trunc(Number.isFinite(debounceMs) ? debounceMs : 0)); + } + return Object.keys(bag).length > 0 ? bag : null; +} + +function isPlainObject(value: unknown): value is Record { + return Boolean(value) && typeof value === "object" && !Array.isArray(value); +} + function normalizeModelFallbackForPayload(config: ModelFallbackConfig): ModelFallbackConfig { const candidates = (config.candidates ?? []) .map((candidate) => ({ diff --git a/ui/web/src/pages/agents/agent-detail/config-sections/inbound-debounce-section.tsx b/ui/web/src/pages/agents/agent-detail/config-sections/inbound-debounce-section.tsx new file mode 100644 index 00000000..847d9090 --- /dev/null +++ b/ui/web/src/pages/agents/agent-detail/config-sections/inbound-debounce-section.tsx @@ -0,0 +1,63 @@ +import { useTranslation } from "react-i18next"; +import { Input } from "@/components/ui/input"; +import { + Select, + SelectContent, + SelectItem, + SelectTrigger, + SelectValue, +} from "@/components/ui/select"; +import type { InboundDebounceOverrideMode } from "@/types/agent"; +import { InfoLabel } from "./config-section"; + +interface InboundDebounceSectionProps { + mode: InboundDebounceOverrideMode; + debounceMs: number; + onModeChange: (mode: InboundDebounceOverrideMode) => void; + onDebounceMsChange: (value: number) => void; +} + +export function InboundDebounceSection({ + mode, + debounceMs, + onModeChange, + onDebounceMsChange, +}: InboundDebounceSectionProps) { + const { t } = useTranslation("agents"); + const s = "configSections.inboundDebounce"; + return ( +
+
+

{t(`${s}.title`)}

+

{t(`${s}.description`)}

+
+
+
+
+ {t(`${s}.mode`)} + +
+
+ {t(`${s}.debounceMs`)} + onDebounceMsChange(Math.max(0, Number(event.target.value)))} + placeholder="0" + /> +
+
+
+
+ ); +} diff --git a/ui/web/src/pages/agents/agent-detail/config-sections/index.ts b/ui/web/src/pages/agents/agent-detail/config-sections/index.ts index 40563194..2c2ed95b 100644 --- a/ui/web/src/pages/agents/agent-detail/config-sections/index.ts +++ b/ui/web/src/pages/agents/agent-detail/config-sections/index.ts @@ -9,3 +9,4 @@ export { ThinkingSection } from "./thinking-section"; export { WorkspaceSharingSection } from "./workspace-sharing-section"; export { ChatGPTOAuthRoutingSection } from "./chatgpt-oauth-routing-section"; export { ModelFallbackSection } from "./model-fallback-section"; +export { InboundDebounceSection } from "./inbound-debounce-section"; diff --git a/ui/web/src/pages/config/sections/behavior-rate-card.tsx b/ui/web/src/pages/config/sections/behavior-rate-card.tsx index 5c891733..b15331cc 100644 --- a/ui/web/src/pages/config/sections/behavior-rate-card.tsx +++ b/ui/web/src/pages/config/sections/behavior-rate-card.tsx @@ -59,8 +59,8 @@ export function BehaviorRateCard({ value, onChange }: Props) { type="number" value={value.inbound_debounce_ms ?? ""} onChange={(e) => update({ inbound_debounce_ms: Number(e.target.value) })} - placeholder="1000 (-1 = disabled)" - min={-1} + placeholder="0 (no debounce)" + min={0} />
diff --git a/ui/web/src/pages/config/sections/behavior-section.tsx b/ui/web/src/pages/config/sections/behavior-section.tsx index 94126856..37b63e18 100644 --- a/ui/web/src/pages/config/sections/behavior-section.tsx +++ b/ui/web/src/pages/config/sections/behavior-section.tsx @@ -83,13 +83,8 @@ export function BehaviorSection({ config, onPatch, saving }: Props) { (v: T) => { setter(v); setDirty(true); }; const handleSave = () => { - // Strip masked/secret values to avoid overwriting real secrets with "***" - const gwClean = Object.fromEntries( - Object.entries(gw).filter(([, v]) => typeof v !== "string" || v !== "***"), - ); onPatch({ gateway: { - ...gwClean, tool_status: ux.tool_status, block_reply: ux.block_reply, max_message_chars: rate.max_message_chars, @@ -98,12 +93,11 @@ export function BehaviorSection({ config, onPatch, saving }: Props) { injection_action: security.injection_action, }, agents: { - ...config.agents, - defaults: { ...ag, intent_classify: ux.intent_classify }, + defaults: { intent_classify: ux.intent_classify }, }, - tools: { ...tl, scrub_credentials: security.scrub_credentials }, - sessions: { ...ss, ...sessions }, - channels: { ...ch, pending_compaction: pendingCompaction }, + tools: { scrub_credentials: security.scrub_credentials }, + sessions, + channels: { pending_compaction: pendingCompaction }, }); }; diff --git a/ui/web/src/types/agent.ts b/ui/web/src/types/agent.ts index 9016c1fc..c4934b97 100644 --- a/ui/web/src/types/agent.ts +++ b/ui/web/src/types/agent.ts @@ -108,6 +108,7 @@ export type EffectiveChatGPTOAuthRoutingStrategy = export type ChatGPTOAuthRoutingOverrideMode = "inherit" | "custom"; export type ReasoningOverrideMode = "inherit" | "custom"; +export type InboundDebounceOverrideMode = "inherit" | "custom"; export interface AgentReasoningConfig { override_mode?: ReasoningOverrideMode; @@ -115,6 +116,11 @@ export interface AgentReasoningConfig { fallback?: "downgrade" | "provider_default" | "off"; } +export interface InboundDebounceConfig { + override_mode?: InboundDebounceOverrideMode; + inbound_debounce_ms?: number; +} + export interface ChatGPTOAuthRoutingConfig { override_mode?: ChatGPTOAuthRoutingOverrideMode; strategy?: ChatGPTOAuthRoutingStrategy;