mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-11 12:18:59 +00:00
feat(chat): debounce rapid web messages
This commit is contained in:
1 parent
8b661951c0
commit
8af0e4ae09
14 files changed
+498
-57
No files matched your search
@@ -67,6 +67,10 @@ The consumer routes system messages based on sender ID prefixes:
|
||||
| `delegate:` | Parent agent's original session (legacy session key format) | team |
|
||||
| `teammate:` | Target agent session | team |
|
||||
|
||||
### 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.
|
||||
|
||||
---
|
||||
|
||||
## 2. Channel Interfaces
|
||||
@@ -509,7 +513,7 @@ The Slack channel uses the `slack-go/slack` library to connect via Socket Mode (
|
||||
- **Mention gating**: `requireMention` default true; `<@botUserID>` stripped from content
|
||||
- **Thread participation cache**: After bot replies in a thread, subsequent messages in that thread auto-trigger response without @mention (24h TTL)
|
||||
- **Message dedup**: `channel+ts` key prevents duplicate processing on Socket Mode reconnect
|
||||
- **Message debounce**: Per-thread batching of rapid messages (300ms default, configurable)
|
||||
- **Message debounce**: Per-thread batching of rapid messages (300ms default, configurable; `debounce_delay: 0` disables)
|
||||
- **Dead socket classification**: Non-retryable auth errors (invalid_auth, token_revoked) fail fast instead of infinite reconnect
|
||||
- **Streaming**: Edit-in-place via `chat.update` with 1000ms throttle (Slack Tier 3 rate limit)
|
||||
- **Reactions**: Status emoji on user messages (thinking_face, hammer_and_wrench, white_check_mark, x, hourglass_flowing_sand)
|
||||
|
||||
@@ -106,6 +106,8 @@ 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.
|
||||
|
||||
### `chat.history`
|
||||
|
||||
Retrieve chat history for a session.
|
||||
|
||||
@@ -6,6 +6,21 @@ Significant changes, features, and fixes in reverse chronological order.
|
||||
|
||||
## 2026-05-22
|
||||
|
||||
### Messaging debounce hardening
|
||||
|
||||
**Features**
|
||||
|
||||
- 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 Slack `debounce_delay: 0` so it disables per-thread batching instead of falling back to the default.
|
||||
|
||||
**Tests**
|
||||
|
||||
- Added regression coverage for channel inbound debounce, Web Chat debounce, media/cancel bypass handling, and Slack debounce config defaults.
|
||||
|
||||
### CLI P6 backend API unblock
|
||||
|
||||
**Features**
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
package bus
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestInboundDebouncerMergesRapidText(t *testing.T) {
|
||||
out := make(chan InboundMessage, 1)
|
||||
d := NewInboundDebouncer(20*time.Millisecond, func(msg InboundMessage) {
|
||||
out <- msg
|
||||
})
|
||||
defer d.Stop()
|
||||
|
||||
d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "one"})
|
||||
d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "two", Metadata: map[string]string{"message_id": "m2"}})
|
||||
|
||||
got := waitInbound(t, out)
|
||||
if got.Content != "one\ntwo" {
|
||||
t.Fatalf("merged content = %q, want %q", got.Content, "one\ntwo")
|
||||
}
|
||||
if got.Metadata["message_id"] != "m2" {
|
||||
t.Fatalf("metadata should come from latest message, got %#v", got.Metadata)
|
||||
}
|
||||
}
|
||||
|
||||
func TestInboundDebouncerDisabledPassesThrough(t *testing.T) {
|
||||
out := make(chan InboundMessage, 2)
|
||||
d := NewInboundDebouncer(-1, func(msg InboundMessage) {
|
||||
out <- msg
|
||||
})
|
||||
|
||||
d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "one"})
|
||||
d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "two"})
|
||||
|
||||
if got := waitInbound(t, out); got.Content != "one" {
|
||||
t.Fatalf("first content = %q", got.Content)
|
||||
}
|
||||
if got := waitInbound(t, out); got.Content != "two" {
|
||||
t.Fatalf("second content = %q", got.Content)
|
||||
}
|
||||
}
|
||||
|
||||
func TestInboundDebouncerMediaFlushesPendingTextFirst(t *testing.T) {
|
||||
out := make(chan InboundMessage, 2)
|
||||
d := NewInboundDebouncer(time.Minute, func(msg InboundMessage) {
|
||||
out <- msg
|
||||
})
|
||||
defer d.Stop()
|
||||
|
||||
d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "pending"})
|
||||
d.Push(InboundMessage{
|
||||
Channel: "telegram",
|
||||
ChatID: "chat-1",
|
||||
SenderID: "user-1",
|
||||
Content: "with media",
|
||||
Media: []MediaFile{{Path: "/tmp/a.png", MimeType: "image/png"}},
|
||||
})
|
||||
|
||||
if got := waitInbound(t, out); got.Content != "pending" || len(got.Media) != 0 {
|
||||
t.Fatalf("first flush = %#v, want pending text without media", got)
|
||||
}
|
||||
if got := waitInbound(t, out); got.Content != "with media" || len(got.Media) != 1 {
|
||||
t.Fatalf("second flush = %#v, want media message", got)
|
||||
}
|
||||
}
|
||||
|
||||
func waitInbound(t *testing.T, ch <-chan InboundMessage) InboundMessage {
|
||||
t.Helper()
|
||||
select {
|
||||
case msg := <-ch:
|
||||
return msg
|
||||
case <-time.After(500 * time.Millisecond):
|
||||
t.Fatal("timed out waiting for debounced message")
|
||||
return InboundMessage{}
|
||||
}
|
||||
}
|
||||
@@ -101,8 +101,11 @@ func New(cfg config.SlackConfig, msgBus *bus.MessageBus, pairingSvc store.Pairin
|
||||
historyLimit = channels.DefaultGroupHistoryLimit
|
||||
}
|
||||
|
||||
debounceDelay := time.Duration(cfg.DebounceDelay) * time.Millisecond
|
||||
if cfg.DebounceDelay == 0 {
|
||||
debounceDelay := 300 * time.Millisecond
|
||||
if cfg.DebounceDelay != nil {
|
||||
debounceDelay = time.Duration(*cfg.DebounceDelay) * time.Millisecond
|
||||
}
|
||||
if debounceDelay < 0 {
|
||||
debounceDelay = 300 * time.Millisecond
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
package slack
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nextlevelbuilder/goclaw/internal/bus"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/config"
|
||||
)
|
||||
|
||||
func TestSlackDebounceDelayDefaultAndDisabled(t *testing.T) {
|
||||
defaultChannel := newTestSlackChannel(t, config.SlackConfig{})
|
||||
if defaultChannel.debounceDelay != 300*time.Millisecond {
|
||||
t.Fatalf("default debounce = %s, want 300ms", defaultChannel.debounceDelay)
|
||||
}
|
||||
|
||||
disabled := 0
|
||||
disabledChannel := newTestSlackChannel(t, config.SlackConfig{DebounceDelay: &disabled})
|
||||
if disabledChannel.debounceDelay != 0 {
|
||||
t.Fatalf("disabled debounce = %s, want 0", disabledChannel.debounceDelay)
|
||||
}
|
||||
|
||||
custom := 750
|
||||
customChannel := newTestSlackChannel(t, config.SlackConfig{DebounceDelay: &custom})
|
||||
if customChannel.debounceDelay != 750*time.Millisecond {
|
||||
t.Fatalf("custom debounce = %s, want 750ms", customChannel.debounceDelay)
|
||||
}
|
||||
}
|
||||
|
||||
func newTestSlackChannel(t *testing.T, cfg config.SlackConfig) *Channel {
|
||||
t.Helper()
|
||||
cfg.Enabled = true
|
||||
cfg.BotToken = "xoxb-test"
|
||||
cfg.AppToken = "xapp-test"
|
||||
ch, err := New(cfg, bus.New(), nil, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("New() error = %v", err)
|
||||
}
|
||||
return ch
|
||||
}
|
||||
@@ -29,7 +29,7 @@ type slackInstanceConfig struct {
|
||||
NativeStream *bool `json:"native_stream,omitempty"`
|
||||
ReactionLevel string `json:"reaction_level,omitempty"`
|
||||
BlockReply *bool `json:"block_reply,omitempty"`
|
||||
DebounceDelay int `json:"debounce_delay,omitempty"`
|
||||
DebounceDelay *int `json:"debounce_delay,omitempty"`
|
||||
ThreadTTL *int `json:"thread_ttl,omitempty"`
|
||||
}
|
||||
|
||||
|
||||
@@ -126,7 +126,7 @@ type SlackConfig struct {
|
||||
NativeStream *bool `json:"native_stream,omitempty"` // use Slack ChatStreamer API if available (default false)
|
||||
ReactionLevel string `json:"reaction_level,omitempty"` // "off" (default), "minimal", "full"
|
||||
BlockReply *bool `json:"block_reply,omitempty"` // override gateway block_reply (nil = inherit)
|
||||
DebounceDelay int `json:"debounce_delay,omitempty"` // ms delay before dispatching rapid messages (default 300, 0=disabled)
|
||||
DebounceDelay *int `json:"debounce_delay,omitempty"` // ms delay before dispatching rapid messages (default 300, 0=disabled)
|
||||
ThreadTTL *int `json:"thread_ttl,omitempty"` // hours before thread participation expires (default 24, 0=disabled — always require @mention)
|
||||
MediaMaxBytes int64 `json:"media_max_bytes,omitempty"` // max file download size in bytes (default 20MB)
|
||||
}
|
||||
@@ -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 messages from same sender (default 1000ms, -1 = disabled)
|
||||
InboundDebounceMs int `json:"inbound_debounce_ms,omitempty"` // merge rapid channel/Web Chat messages from same sender/session (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)
|
||||
ToolStatus *bool `json:"tool_status,omitempty"` // show tool name in streaming preview during tool execution (default true)
|
||||
|
||||
@@ -11,10 +11,10 @@ import (
|
||||
"github.com/nextlevelbuilder/goclaw/internal/agent"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/audio"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/bus"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/config"
|
||||
httpapi "github.com/nextlevelbuilder/goclaw/internal/http"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/channels/media"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/config"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/gateway"
|
||||
httpapi "github.com/nextlevelbuilder/goclaw/internal/http"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/i18n"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/providers"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/sessions"
|
||||
@@ -32,10 +32,13 @@ type ChatMethods struct {
|
||||
eventBus bus.EventPublisher
|
||||
postTurn tools.PostTurnProcessor
|
||||
audioMgr *audio.Manager // for TTS auto-apply on WS responses (nil = disabled)
|
||||
debouncer *chatDebouncer
|
||||
}
|
||||
|
||||
func NewChatMethods(agents *agent.Router, sess store.SessionStore, cfg *config.Config, rl *gateway.RateLimiter, eventBus bus.EventPublisher) *ChatMethods {
|
||||
return &ChatMethods{agents: agents, sessions: sess, cfg: cfg, rateLimiter: rl, eventBus: eventBus}
|
||||
m := &ChatMethods{agents: agents, sessions: sess, cfg: cfg, rateLimiter: rl, eventBus: eventBus}
|
||||
m.debouncer = newChatDebouncer(m.dispatchChatSends)
|
||||
return m
|
||||
}
|
||||
|
||||
// SetAudioManager sets the audio manager for TTS auto-apply on WS responses.
|
||||
@@ -102,11 +105,11 @@ type chatMediaItem struct {
|
||||
}
|
||||
|
||||
type chatSendParams struct {
|
||||
Message string `json:"message"`
|
||||
AgentID string `json:"agentId"`
|
||||
SessionKey string `json:"sessionKey"`
|
||||
Stream bool `json:"stream"`
|
||||
Media json.RawMessage `json:"media,omitempty"` // []string (legacy) or []chatMediaItem
|
||||
Message string `json:"message"`
|
||||
AgentID string `json:"agentId"`
|
||||
SessionKey string `json:"sessionKey"`
|
||||
Stream bool `json:"stream"`
|
||||
Media json.RawMessage `json:"media,omitempty"` // []string (legacy) or []chatMediaItem
|
||||
}
|
||||
|
||||
// parseMedia handles both legacy string paths and new {path,filename} objects.
|
||||
@@ -175,69 +178,106 @@ func (m *ChatMethods) handleSend(ctx context.Context, client *gateway.Client, re
|
||||
return
|
||||
}
|
||||
|
||||
runID := uuid.NewString()
|
||||
providedSessionKey := params.SessionKey != ""
|
||||
sessionKey := params.SessionKey
|
||||
if sessionKey == "" {
|
||||
sessionKey = sessions.BuildWSSessionKey(params.AgentID, uuid.NewString())
|
||||
}
|
||||
params.SessionKey = sessionKey
|
||||
|
||||
// Ownership check: when resuming an existing session, verify the caller owns it.
|
||||
// Skip for new sessions (Get returns nil) so first-message creation is not blocked.
|
||||
if params.SessionKey != "" && !canSeeAll(client.Role(), m.cfg.Gateway.OwnerIDs, userID) {
|
||||
if providedSessionKey && !canSeeAll(client.Role(), m.cfg.Gateway.OwnerIDs, userID) {
|
||||
if sess := m.sessions.Get(ctx, sessionKey); sess != nil && sess.UserID != userID {
|
||||
client.SendResponse(protocol.NewErrorResponse(req.ID, protocol.ErrUnauthorized, i18n.T(locale, i18n.MsgPermissionDenied, "session")))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// Detach from HTTP request context so agent runs survive page navigation/reconnect.
|
||||
// WithoutCancel preserves all context values (locale, user ID, etc.)
|
||||
// but HTTP request cancellation no longer propagates.
|
||||
// Explicit abort via chat.abort still works through the per-run cancel().
|
||||
runCtxBase := context.WithoutCancel(ctx)
|
||||
if userID != "" {
|
||||
runCtxBase = store.WithUserID(runCtxBase, userID)
|
||||
item := chatSendRequest{
|
||||
ctx: ctx,
|
||||
client: client,
|
||||
requestID: req.ID,
|
||||
params: params,
|
||||
loop: loop,
|
||||
userID: userID,
|
||||
sessionKey: sessionKey,
|
||||
}
|
||||
debounceKey := chatDebounceKey(userID, sessionKey)
|
||||
if m.debouncer == nil {
|
||||
m.debouncer = newChatDebouncer(m.dispatchChatSends)
|
||||
}
|
||||
if m.agents.IsSessionBusy(sessionKey) && agent.IsExactCancelKeyword(params.Message) {
|
||||
m.debouncer.Discard(debounceKey)
|
||||
m.abortChatSession(req.ID, client, sessionKey)
|
||||
return
|
||||
}
|
||||
if len(params.parseMedia()) > 0 {
|
||||
pending := m.debouncer.Take(debounceKey)
|
||||
m.dispatchChatSends(append(pending, item))
|
||||
return
|
||||
}
|
||||
if delay := chatDebounceDelay(m.cfg); delay > 0 {
|
||||
m.debouncer.Push(debounceKey, delay, item)
|
||||
return
|
||||
}
|
||||
m.dispatchChatSends([]chatSendRequest{item})
|
||||
}
|
||||
|
||||
// Mid-run injection: if session already has an active run, inject the message
|
||||
// into the running loop instead of starting a new concurrent run.
|
||||
if m.agents.IsSessionBusy(sessionKey) {
|
||||
// Exact cancel keyword detection: auto-abort when user sends "stop", "cancel", etc.
|
||||
if agent.IsExactCancelKeyword(params.Message) {
|
||||
results := m.agents.AbortRunsForSession(sessionKey)
|
||||
aborted := false
|
||||
for _, r := range results {
|
||||
if r.Stopped || r.Forced {
|
||||
aborted = true
|
||||
break
|
||||
}
|
||||
}
|
||||
client.SendResponse(protocol.NewOKResponse(req.ID, map[string]any{
|
||||
"cancelled": true,
|
||||
"aborted": aborted,
|
||||
}))
|
||||
return
|
||||
func (m *ChatMethods) abortChatSession(reqID string, client *gateway.Client, sessionKey string) {
|
||||
results := m.agents.AbortRunsForSession(sessionKey)
|
||||
aborted := false
|
||||
for _, r := range results {
|
||||
if r.Stopped || r.Forced {
|
||||
aborted = true
|
||||
break
|
||||
}
|
||||
}
|
||||
client.SendResponse(protocol.NewOKResponse(reqID, map[string]any{
|
||||
"cancelled": true,
|
||||
"aborted": aborted,
|
||||
}))
|
||||
}
|
||||
|
||||
func (m *ChatMethods) dispatchChatSends(requests []chatSendRequest) {
|
||||
if len(requests) == 0 {
|
||||
return
|
||||
}
|
||||
primary := requests[len(requests)-1]
|
||||
params := mergeChatSendRequests(requests)
|
||||
sessionKey := primary.sessionKey
|
||||
userID := primary.userID
|
||||
loop := primary.loop
|
||||
hasMedia := len(params.parseMedia()) > 0
|
||||
|
||||
// Mid-run injection: debounce rapid follow-ups into a single injected message.
|
||||
if !hasMedia && m.agents.IsSessionBusy(sessionKey) {
|
||||
injected := m.agents.InjectMessage(sessionKey, agent.InjectedMessage{
|
||||
Content: params.Message,
|
||||
UserID: userID,
|
||||
})
|
||||
if injected {
|
||||
client.SendResponse(protocol.NewOKResponse(req.ID, map[string]any{
|
||||
"injected": true,
|
||||
}))
|
||||
sendChatOK(requests, map[string]any{"injected": true})
|
||||
return
|
||||
}
|
||||
// Fallback: injection failed (channel full), proceed with new run
|
||||
// Fallback: injection failed (channel full), proceed with new run.
|
||||
}
|
||||
|
||||
// Detach from HTTP request context so agent runs survive page navigation/reconnect.
|
||||
// WithoutCancel preserves all context values (locale, user ID, etc.)
|
||||
// but HTTP request cancellation no longer propagates.
|
||||
// Explicit abort via chat.abort still works through the per-run cancel().
|
||||
runCtxBase := context.WithoutCancel(primary.ctx)
|
||||
if userID != "" {
|
||||
runCtxBase = store.WithUserID(runCtxBase, userID)
|
||||
}
|
||||
// Inject team dispatch tracker: gates team_tasks create (must search/list first)
|
||||
// and defers task dispatch to post-turn.
|
||||
runCtxBase, drainTeamDispatch := tools.InjectTeamDispatch(runCtxBase, m.postTurn)
|
||||
|
||||
// Create cancellable context for abort support (matching TS AbortController pattern).
|
||||
runCtx, cancel := context.WithCancel(runCtxBase)
|
||||
runID := uuid.NewString()
|
||||
injectCh := m.agents.RegisterRun(runCtxBase, runID, sessionKey, params.AgentID, cancel)
|
||||
|
||||
// Run agent asynchronously - events are broadcast via the event system
|
||||
@@ -284,8 +324,8 @@ func (m *ChatMethods) handleSend(ctx context.Context, client *gateway.Client, re
|
||||
WorkspaceChatID: userID, // mirror ChatID so vault chat_id isolation activates for WS direct flow
|
||||
RunID: runID,
|
||||
UserID: userID,
|
||||
Stream: params.Stream,
|
||||
InjectCh: injectCh,
|
||||
Stream: params.Stream,
|
||||
InjectCh: injectCh,
|
||||
// Wire trace ID back to the active run so force-abort can mark the
|
||||
// correct trace as cancelled if the goroutine does not exit within 3s.
|
||||
OnTraceCreated: func(traceID uuid.UUID) {
|
||||
@@ -297,17 +337,15 @@ func (m *ChatMethods) handleSend(ctx context.Context, client *gateway.Client, re
|
||||
// Send cancelled response so the frontend's chat.send promise resolves
|
||||
// instead of hanging until the 600s timeout.
|
||||
if runCtx.Err() != nil {
|
||||
client.SendResponse(protocol.NewOKResponse(req.ID, map[string]any{
|
||||
"cancelled": true,
|
||||
}))
|
||||
sendChatOK(requests, map[string]any{"cancelled": true})
|
||||
return
|
||||
}
|
||||
client.SendResponse(protocol.NewErrorResponse(req.ID, protocol.ErrInternal, err.Error()))
|
||||
sendChatError(requests, protocol.ErrInternal, err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
// Auto-generate conversation title on first message (label empty = never titled).
|
||||
if label := m.sessions.GetLabel(ctx, sessionKey); label == "" {
|
||||
if label := m.sessions.GetLabel(primary.ctx, sessionKey); label == "" {
|
||||
agentProvider := loop.Provider()
|
||||
agentModel := loop.Model()
|
||||
userMsg := params.Message
|
||||
@@ -324,7 +362,7 @@ func (m *ChatMethods) handleSend(ctx context.Context, client *gateway.Client, re
|
||||
return
|
||||
}
|
||||
bus.BroadcastForTenant(m.eventBus, protocol.EventSessionUpdated,
|
||||
client.TenantID(),
|
||||
primary.client.TenantID(),
|
||||
map[string]string{"sessionKey": sessionKey, "label": title, "userId": userID})
|
||||
}()
|
||||
}
|
||||
@@ -364,10 +402,22 @@ func (m *ChatMethods) handleSend(ctx context.Context, client *gateway.Client, re
|
||||
if len(mediaResults) > 0 {
|
||||
resp["media"] = mediaResults
|
||||
}
|
||||
client.SendResponse(protocol.NewOKResponse(req.ID, resp))
|
||||
sendChatOK(requests, resp)
|
||||
}()
|
||||
}
|
||||
|
||||
func sendChatOK(requests []chatSendRequest, payload map[string]any) {
|
||||
for _, request := range requests {
|
||||
request.client.SendResponse(protocol.NewOKResponse(request.requestID, payload))
|
||||
}
|
||||
}
|
||||
|
||||
func sendChatError(requests []chatSendRequest, code, message string) {
|
||||
for _, request := range requests {
|
||||
request.client.SendResponse(protocol.NewErrorResponse(request.requestID, code, message))
|
||||
}
|
||||
}
|
||||
|
||||
type chatHistoryParams struct {
|
||||
AgentID string `json:"agentId"`
|
||||
SessionKey string `json:"sessionKey"`
|
||||
|
||||
@@ -0,0 +1,143 @@
|
||||
package methods
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/nextlevelbuilder/goclaw/internal/agent"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/config"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/gateway"
|
||||
)
|
||||
|
||||
type chatSendRequest struct {
|
||||
ctx context.Context
|
||||
client *gateway.Client
|
||||
requestID string
|
||||
params chatSendParams
|
||||
loop agent.Agent
|
||||
userID string
|
||||
sessionKey string
|
||||
}
|
||||
|
||||
type chatDebouncer struct {
|
||||
mu sync.Mutex
|
||||
buffers map[string]*chatDebounceBuffer
|
||||
flushFn func([]chatSendRequest)
|
||||
}
|
||||
|
||||
type chatDebounceBuffer struct {
|
||||
items []chatSendRequest
|
||||
timer *time.Timer
|
||||
}
|
||||
|
||||
func newChatDebouncer(flushFn func([]chatSendRequest)) *chatDebouncer {
|
||||
return &chatDebouncer{
|
||||
buffers: make(map[string]*chatDebounceBuffer),
|
||||
flushFn: flushFn,
|
||||
}
|
||||
}
|
||||
|
||||
func (d *chatDebouncer) Push(key string, delay time.Duration, item chatSendRequest) {
|
||||
if delay <= 0 {
|
||||
d.flushFn([]chatSendRequest{item})
|
||||
return
|
||||
}
|
||||
|
||||
d.mu.Lock()
|
||||
buf, exists := d.buffers[key]
|
||||
if !exists {
|
||||
buf = &chatDebounceBuffer{}
|
||||
d.buffers[key] = buf
|
||||
}
|
||||
buf.items = append(buf.items, item)
|
||||
if buf.timer != nil {
|
||||
buf.timer.Stop()
|
||||
}
|
||||
buf.timer = time.AfterFunc(delay, func() {
|
||||
d.Flush(key)
|
||||
})
|
||||
d.mu.Unlock()
|
||||
}
|
||||
|
||||
func (d *chatDebouncer) Flush(key string) {
|
||||
items := d.Take(key)
|
||||
if len(items) == 0 {
|
||||
return
|
||||
}
|
||||
d.flushFn(items)
|
||||
}
|
||||
|
||||
func (d *chatDebouncer) Take(key string) []chatSendRequest {
|
||||
d.mu.Lock()
|
||||
buf, ok := d.buffers[key]
|
||||
if !ok || len(buf.items) == 0 {
|
||||
d.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
if buf.timer != nil {
|
||||
buf.timer.Stop()
|
||||
}
|
||||
items := buf.items
|
||||
delete(d.buffers, key)
|
||||
d.mu.Unlock()
|
||||
|
||||
return items
|
||||
}
|
||||
|
||||
func (d *chatDebouncer) Discard(key string) {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
buf, ok := d.buffers[key]
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if buf.timer != nil {
|
||||
buf.timer.Stop()
|
||||
}
|
||||
delete(d.buffers, key)
|
||||
}
|
||||
|
||||
func (d *chatDebouncer) Stop() {
|
||||
d.mu.Lock()
|
||||
keys := make([]string, 0, len(d.buffers))
|
||||
for key := range d.buffers {
|
||||
keys = append(keys, key)
|
||||
}
|
||||
d.mu.Unlock()
|
||||
|
||||
for _, key := range keys {
|
||||
d.Flush(key)
|
||||
}
|
||||
}
|
||||
|
||||
func mergeChatSendRequests(items []chatSendRequest) chatSendParams {
|
||||
if len(items) == 0 {
|
||||
return chatSendParams{}
|
||||
}
|
||||
last := items[len(items)-1].params
|
||||
parts := make([]string, 0, len(items))
|
||||
for _, item := range items {
|
||||
if item.params.Message != "" {
|
||||
parts = append(parts, item.params.Message)
|
||||
}
|
||||
}
|
||||
last.Message = strings.Join(parts, "\n")
|
||||
return last
|
||||
}
|
||||
|
||||
func chatDebounceDelay(cfg *config.Config) time.Duration {
|
||||
debounceMs := 0
|
||||
if cfg != nil {
|
||||
debounceMs = cfg.Gateway.InboundDebounceMs
|
||||
}
|
||||
if debounceMs == 0 {
|
||||
debounceMs = 1000
|
||||
}
|
||||
return time.Duration(debounceMs) * time.Millisecond
|
||||
}
|
||||
|
||||
func chatDebounceKey(userID, sessionKey string) string {
|
||||
return userID + ":" + sessionKey
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
package methods
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nextlevelbuilder/goclaw/internal/config"
|
||||
)
|
||||
|
||||
func TestMergeChatSendRequestsJoinsContentAndUsesLatestParams(t *testing.T) {
|
||||
items := []chatSendRequest{
|
||||
{params: chatSendParams{Message: "first", AgentID: "agent-a", SessionKey: "session-a", Stream: false}},
|
||||
{params: chatSendParams{Message: "", AgentID: "agent-a", SessionKey: "session-a", Stream: true}},
|
||||
{params: chatSendParams{Message: "second", AgentID: "agent-a", SessionKey: "session-a", Stream: true}},
|
||||
}
|
||||
|
||||
got := mergeChatSendRequests(items)
|
||||
if got.Message != "first\nsecond" {
|
||||
t.Fatalf("merged message = %q, want %q", got.Message, "first\nsecond")
|
||||
}
|
||||
if !got.Stream {
|
||||
t.Fatal("latest params should win for stream flag")
|
||||
}
|
||||
}
|
||||
|
||||
func TestChatDebouncerFlushesOnceAfterQuietWindow(t *testing.T) {
|
||||
out := make(chan []chatSendRequest, 1)
|
||||
d := newChatDebouncer(func(items []chatSendRequest) {
|
||||
out <- items
|
||||
})
|
||||
defer d.Stop()
|
||||
|
||||
d.Push("u1:s1", 20*time.Millisecond, chatSendRequest{params: chatSendParams{Message: "one"}})
|
||||
d.Push("u1:s1", 20*time.Millisecond, chatSendRequest{params: chatSendParams{Message: "two"}})
|
||||
|
||||
items := waitChatDebounce(t, out)
|
||||
if len(items) != 2 {
|
||||
t.Fatalf("flushed items = %d, want 2", len(items))
|
||||
}
|
||||
if got := mergeChatSendRequests(items).Message; got != "one\ntwo" {
|
||||
t.Fatalf("merged message = %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestChatDebouncerTakeDrainsPendingBeforeBypass(t *testing.T) {
|
||||
out := make(chan []chatSendRequest, 1)
|
||||
d := newChatDebouncer(func(items []chatSendRequest) {
|
||||
out <- items
|
||||
})
|
||||
defer d.Stop()
|
||||
|
||||
d.Push("u1:s1", time.Minute, chatSendRequest{params: chatSendParams{Message: "pending"}})
|
||||
|
||||
items := d.Take("u1:s1")
|
||||
if len(items) != 1 || items[0].params.Message != "pending" {
|
||||
t.Fatalf("flushed items = %#v", items)
|
||||
}
|
||||
assertNoChatDebounceFlush(t, out)
|
||||
}
|
||||
|
||||
func TestChatDebouncerDiscardDropsPendingBeforeCancel(t *testing.T) {
|
||||
out := make(chan []chatSendRequest, 1)
|
||||
d := newChatDebouncer(func(items []chatSendRequest) {
|
||||
out <- items
|
||||
})
|
||||
defer d.Stop()
|
||||
|
||||
d.Push("u1:s1", 20*time.Millisecond, chatSendRequest{params: chatSendParams{Message: "pending"}})
|
||||
d.Discard("u1:s1")
|
||||
|
||||
assertNoChatDebounceFlush(t, out)
|
||||
}
|
||||
|
||||
func TestChatDebounceDelayDefaultAndDisabled(t *testing.T) {
|
||||
if got := chatDebounceDelay(&config.Config{}); got != time.Second {
|
||||
t.Fatalf("default debounce = %s, want 1s", 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)
|
||||
}
|
||||
}
|
||||
|
||||
func waitChatDebounce(t *testing.T, ch <-chan []chatSendRequest) []chatSendRequest {
|
||||
t.Helper()
|
||||
select {
|
||||
case items := <-ch:
|
||||
return items
|
||||
case <-time.After(500 * time.Millisecond):
|
||||
t.Fatal("timed out waiting for chat debounce flush")
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
func assertNoChatDebounceFlush(t *testing.T, ch <-chan []chatSendRequest) {
|
||||
t.Helper()
|
||||
select {
|
||||
case items := <-ch:
|
||||
t.Fatalf("unexpected flush: %#v", items)
|
||||
case <-time.After(50 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
@@ -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 a new message, to coalesce rapid follow-ups. Set to -1 to disable.",
|
||||
"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.injectionAction": "Injection Detection Action",
|
||||
|
||||
"agents.title": "Agent Defaults",
|
||||
|
||||
@@ -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ý tin nhắn mới, để gộp các tin nhắn nhanh liên tiếp. Đặt -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 để dùng mặc định 1000ms, hoặc -1 để tắt.",
|
||||
"gateway.injectionAction": "Hành động phát hiện injection",
|
||||
|
||||
"agents.title": "Mặc định agent",
|
||||
|
||||
@@ -42,7 +42,7 @@
|
||||
"gateway.rateLimitRpm": "速率限制(RPM)",
|
||||
"gateway.rateLimitRpmTip": "每用户每分钟最大请求数。设为 0 禁用速率限制。",
|
||||
"gateway.inboundDebounceMs": "入站防抖(毫秒)",
|
||||
"gateway.inboundDebounceMsTip": "处理新消息前的延迟毫秒数,用于合并快速连续消息。设为 -1 禁用。",
|
||||
"gateway.inboundDebounceMsTip": "处理来自频道或 Web Chat 的快速连续消息前的延迟毫秒数。设为 0 使用默认 1000ms,设为 -1 禁用。",
|
||||
"gateway.injectionAction": "注入检测动作",
|
||||
|
||||
"agents.title": "Agent 默认值",
|
||||
|
||||
Reference in new issue
Block a user