From f771cff77c8ed9b6f27d8450e99921dad412fc59 Mon Sep 17 00:00:00 2001 From: Duy /zuey/ Date: Thu, 28 May 2026 18:30:34 +0700 Subject: [PATCH] fix(channels): coalesce multi-attachment inbounds (#63) (#90) Implements 3 coalescing layers to handle rapid multi-attachment inbounds: - Bus debouncer: delays inbound messages 1s, merges duplicates - Web chat debouncer: buffers client-side inbound frames for batch RPC - Telegram album aggregator: collects album members via AfterFunc+Stop timer Drops media-bypass shortcut (forces 1s media floor). Aggregator enforces: - AfterFunc+Stop timer discipline with ordered drain on stop - 2-tuple (album_id, sender) keying for isolation - Dual DoS caps: max 10 albums per sender, max 100 messages per album - merged_message_ids dedup seeding across all 3 surfaces Closes #63 --- CHANGELOG.md | 31 +++ CONTRIBUTING.md | 42 +++ cmd/gateway_consumer.go | 7 +- cmd/gateway_consumer_debounce.go | 36 ++- cmd/gateway_consumer_debounce_test.go | 138 +++++++++ cmd/gateway_consumer_dedup.go | 46 +++ cmd/gateway_consumer_dedup_test.go | 67 +++++ docs/05-channels-messaging.md | 12 +- internal/bus/inbound_debounce.go | 98 ++++++- internal/bus/inbound_debounce_test.go | 194 ++++++++++++- .../channels/telegram/album_aggregator.go | 175 ++++++++++++ .../telegram/album_aggregator_test.go | 262 ++++++++++++++++++ internal/channels/telegram/channel.go | 33 ++- internal/channels/telegram/handlers.go | 164 +++++++---- .../telegram/resolved_message_context.go | 38 +++ internal/config/config_channels.go | 2 +- internal/gateway/methods/chat.go | 16 +- internal/gateway/methods/chat_debounce.go | 40 ++- .../methods/chat_debounce_media_test.go | 92 ++++++ .../gateway/methods/chat_debounce_test.go | 11 +- 20 files changed, 1408 insertions(+), 96 deletions(-) create mode 100644 cmd/gateway_consumer_debounce_test.go create mode 100644 cmd/gateway_consumer_dedup.go create mode 100644 cmd/gateway_consumer_dedup_test.go create mode 100644 internal/channels/telegram/album_aggregator.go create mode 100644 internal/channels/telegram/album_aggregator_test.go create mode 100644 internal/channels/telegram/resolved_message_context.go create mode 100644 internal/gateway/methods/chat_debounce_media_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 8c985914..219b043b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -46,6 +46,37 @@ All notable changes to GoClaw are documented here. For full documentation, see [ ### Fixed +- **Multi-attachment messages no longer trigger N agent replies (#63).** + Three coalescing surfaces hardened so a single user action produces ONE + agent run regardless of how the platform delivers attachments: + 1. **Bus debouncer** — removed the media-bypass shortcut that fired + immediately for any message with attachments; media now goes through + the same per-(channel, chatID, senderID, agentID) silence window as + text. Media-floor (`max(configured, mediaFloor)`) guarantees a + minimum window when attachments are present so multi-file uploads + coalesce. Dedup seed prevents the same `MessageID` from being + buffered twice on bursty arrivals. + 2. **Web Chat debouncer** (`internal/gateway/methods/chat_debounce.go`) — + parallel structure for `/v1/chat/completions` streaming: per-session + buffer + media floor + Take/Discard semantics for flush/cancel + control. Merges queued payloads at flush time (latest params win; + text concatenated newline-separated). + 3. **Telegram album aggregator** (`internal/channels/telegram/album_aggregator.go`) — + channel-layer coalescing for albums. Telegram delivers a media-group + (multiple photos/videos shared as one user action) as N separate + `Message` updates sharing a `MediaGroupID`. The aggregator buffers + by `(chatID, MediaGroupID)` after all access gates pass, pins the + sender on first arrival as a security tripwire, and dispatches ONE + `processResolvedMessage` call with all members on a 500ms silence + window. `Stop()` synchronously drains pending buffers before + `pollCancel` so in-flight albums always publish. + + Cross-surface invariants (see CONTRIBUTING.md → "Multi-attachment + coalescing"): no media bypass, media floor on every surface, + drop-and-log dual caps, no `time.Timer.Reset` (use `AfterFunc` + + `Stop`), sender pin on first arrival, post-stop pushes rejected + with warn log. + - **Upstream critical security remediation** — hardens gateway no-token fallback, Feishu/Lark and Pancake webhooks, sandbox path/write handling, tenant-admin checks for mutable HTTP surfaces, and Lite hook schema migration verification. diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 8e9d036e..9b020795 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -72,6 +72,48 @@ Shared predicate: `store.IsMasterScope(ctx)` (`internal/store/context.go`). - `requireAuth(RoleAdmin)` as the **sole** gate on a global-state write - Admin revoke/delete handlers that skip pre-fetch ownership verification (store SQL alone is not enough when it matches `IS NULL` arms) +### Multi-attachment coalescing (#63) + +Three independent surfaces coalesce burst inbounds so one user action produces +one agent run. Any future surface that fans burst arrivals into the agent loop +MUST honor these eight invariants. Drift on any of them re-introduces the +N-replies bug. + +1. **No media bypass.** A message carrying attachments goes through the same + silence window as text. The pre-fix "publish immediately when media is + present" shortcut is the original #63 regression — do not reintroduce it. +2. **Media floor.** When attachments are present, the effective window is + `max(configured, mediaFloor)`. Configured can be 0 (disabled) for text-only + flows; once media arrives the floor is the lower bound so multi-file + uploads have time to land. +3. **Per-key buffer.** Buffer key is the smallest tuple that uniquely names + "this user action in this delivery channel" — `(channel, chatID, senderID, + agentID)` for the bus debouncer, `(userKey, sessionKey)` for web chat, + `(chatID, MediaGroupID)` for Telegram albums. +4. **Sender pin on first arrival.** First arrival pins the senderID on the + buffer. Subsequent arrivals with a mismatched sender are dropped with a + `security.*_sender_mismatch` warn log. Defense-in-depth against spoofed + updates; the platform should never reuse a group/session id across senders. +5. **Drop-and-log dual caps.** Per-buffer cap AND global active-buffer cap. + Overflow logs `*.overflow` with `scope=buffer|global`, drops the + straggler, and returns false to caller — caller falls through to + single-message dispatch so no message is silently lost. +6. **AfterFunc + Stop, never Reset.** Use `time.AfterFunc(window, fn)` and + `timer.Stop()` on every arrival. `time.Timer.Reset()` has a documented + double-fire race when the timer is mid-fire — banned. +7. **Representative is members[0].** The first arrival's resolved context + (sender label, content prefix, reply target, topic config) is the one + that flows downstream on flush. Later arrivals contribute their media + only. +8. **Synchronous Stop drain.** Shutdown order is `aggregator.Stop()` → + `pollCancel()` → `handlerWg.Wait()`. Stop synchronously flushes all + pending buffers BEFORE any context is cancelled so in-flight bursts + reach the agent loop. Post-Stop pushes are rejected with a warn log. + +Surfaces today: `internal/bus/inbound_debounce.go`, +`internal/gateway/methods/chat_debounce.go`, +`internal/channels/telegram/album_aggregator.go`. + ## Test Layers Tests are organized by priority and purpose: diff --git a/cmd/gateway_consumer.go b/cmd/gateway_consumer.go index e2759ef1..24fb27fa 100644 --- a/cmd/gateway_consumer.go +++ b/cmd/gateway_consumer.go @@ -3,7 +3,6 @@ package cmd import ( "context" "encoding/json" - "fmt" "log/slog" "strings" "sync" @@ -89,6 +88,10 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents return resolveInboundDebounceDelay(ctx, msg, deps) }, func(msg bus.InboundMessage) { + // Seed dedup cache with all sibling message_ids from the merged flush + // so any platform retransmit of a sibling (webhook retry, album member + // redelivery) is short-circuited before re-entering the debouncer. + seedDedupFromMerged(dedupe, msg) processNormalMessage(ctx, msg, deps) }, ) @@ -112,7 +115,7 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents // --- Dedup: skip duplicate inbound messages (matching TS shouldSkipDuplicateInbound) --- if msgID := msg.Metadata["message_id"]; msgID != "" { - dedupeKey := fmt.Sprintf("%s|%s|%s|%s", msg.Channel, msg.SenderID, msg.ChatID, msgID) + dedupeKey := dedupKeyFor(msg.Channel, msg.SenderID, msg.ChatID, msgID) if dedupe.IsDuplicate(dedupeKey) { slog.Debug("dedup: skipping duplicate message", "key", dedupeKey) continue diff --git a/cmd/gateway_consumer_debounce.go b/cmd/gateway_consumer_debounce.go index a0c98419..e3aa9f11 100644 --- a/cmd/gateway_consumer_debounce.go +++ b/cmd/gateway_consumer_debounce.go @@ -3,6 +3,7 @@ package cmd import ( "context" "log/slog" + "strings" "time" "github.com/google/uuid" @@ -11,6 +12,12 @@ import ( "github.com/nextlevelbuilder/goclaw/internal/store" ) +// mediaDebounceFloorMs is the minimum debounce window applied when a message +// carries media AND the post-override delay would be 0. Prevents multi-attachment +// bursts from triggering one agent run per attachment when operators disable +// debouncing (issue #63). See plans/260528-1351-multi-attachment-debounce/. +const mediaDebounceFloorMs = 1000 + func prepareInboundDebounceMessage(msg *bus.InboundMessage, deps *ConsumerDeps) { if msg == nil || deps == nil || deps.Cfg == nil || msg.AgentID != "" { return @@ -24,7 +31,7 @@ func resolveInboundDebounceDelay(ctx context.Context, msg bus.InboundMessage, de debounceMs = deps.Cfg.Gateway.InboundDebounceMs } if deps == nil || deps.AgentStore == nil || msg.AgentID == "" { - return inboundDebounceDuration(debounceMs) + return inboundDebounceDuration(applyMediaFloor(debounceMs, msg)) } agentCtx := ctx @@ -39,12 +46,35 @@ func resolveInboundDebounceDelay(ctx context.Context, msg bus.InboundMessage, de if err != nil { slog.Debug("inbound debounce: agent config unavailable", "agent", msg.AgentID, "error", err) } - return inboundDebounceDuration(debounceMs) + return inboundDebounceDuration(applyMediaFloor(debounceMs, msg)) } if overrideMs, ok := agentData.ParseInboundDebounceMs(); ok { debounceMs = overrideMs } - return inboundDebounceDuration(debounceMs) + return inboundDebounceDuration(applyMediaFloor(debounceMs, msg)) +} + +// applyMediaFloor enforces the media debounce floor. +// +// Precedence: floor fires ONLY when the post-override delay is exactly 0. A +// non-zero agent override (even below the floor) is honored verbatim — operators +// who set debounce_ms=500 on a media-receiving agent get 500ms, not 1000ms. +// +// Exemption: internal publishers (SenderID prefix "system:" or "subagent:") bypass +// the floor. Their messages are synthesized by tools/subagents and have no burst- +// arrival semantics; a +1s latency on every tool echo would be a regression. +func applyMediaFloor(delayMs int, msg bus.InboundMessage) int { + if delayMs <= 0 && len(msg.Media) > 0 && !isSystemOrSubagentSender(msg.SenderID) { + return mediaDebounceFloorMs + } + return delayMs +} + +// isSystemOrSubagentSender reports whether the SenderID identifies an internal +// publisher (system: or subagent: colon-prefixed). Plain prefix matches like +// "systemic" or "subagentX" without the colon are NOT internal. +func isSystemOrSubagentSender(senderID string) bool { + return strings.HasPrefix(senderID, "system:") || strings.HasPrefix(senderID, "subagent:") } func getInboundDebounceAgent(ctx context.Context, agentStore store.AgentStore, agentID string) (*store.AgentData, error) { diff --git a/cmd/gateway_consumer_debounce_test.go b/cmd/gateway_consumer_debounce_test.go new file mode 100644 index 00000000..df12e5d5 --- /dev/null +++ b/cmd/gateway_consumer_debounce_test.go @@ -0,0 +1,138 @@ +package cmd + +import ( + "context" + "testing" + "time" + + "github.com/nextlevelbuilder/goclaw/internal/bus" + "github.com/nextlevelbuilder/goclaw/internal/config" +) + +// TestApplyMediaFloor_NoMediaNoFloor: msg without media → debounceMs returned unchanged (0). +func TestApplyMediaFloor_NoMediaNoFloor(t *testing.T) { + msg := bus.InboundMessage{SenderID: "user-1"} + got := applyMediaFloor(0, msg) + if got != 0 { + t.Fatalf("applyMediaFloor(0, no-media) = %d, want 0", got) + } +} + +// TestApplyMediaFloor_MediaAppliesFloorWhenDisabled — Rule #2 happy path. +func TestApplyMediaFloor_MediaAppliesFloorWhenDisabled(t *testing.T) { + msg := bus.InboundMessage{ + SenderID: "user-1", + Media: []bus.MediaFile{{Path: "/x"}}, + } + got := applyMediaFloor(0, msg) + if got != mediaDebounceFloorMs { + t.Fatalf("applyMediaFloor(0, media) = %d, want %d", got, mediaDebounceFloorMs) + } +} + +// TestApplyMediaFloor_AgentOverrideBelowFloorHonored — red-team Rule #2 precedence. +// Floor fires only when the post-override delay is exactly 0. A 500ms override +// for a media-bearing message MUST be honored verbatim. +func TestApplyMediaFloor_AgentOverrideBelowFloorHonored(t *testing.T) { + msg := bus.InboundMessage{ + SenderID: "user-1", + Media: []bus.MediaFile{{Path: "/x"}}, + } + got := applyMediaFloor(500, msg) + if got != 500 { + t.Fatalf("applyMediaFloor(500, media) = %d, want 500 (override honored; floor must not raise)", got) + } +} + +// TestApplyMediaFloor_MediaRespectsConfigWhenAboveFloor: cfg already above floor → unchanged. +func TestApplyMediaFloor_MediaRespectsConfigWhenAboveFloor(t *testing.T) { + msg := bus.InboundMessage{ + SenderID: "user-1", + Media: []bus.MediaFile{{Path: "/x"}}, + } + got := applyMediaFloor(2000, msg) + if got != 2000 { + t.Fatalf("applyMediaFloor(2000, media) = %d, want 2000", got) + } +} + +// TestApplyMediaFloor_SystemSenderExempt — Rule #3 internal-publisher exemption. +func TestApplyMediaFloor_SystemSenderExempt(t *testing.T) { + msg := bus.InboundMessage{ + SenderID: "system:tool-echo", + Media: []bus.MediaFile{{Path: "/x"}}, + } + got := applyMediaFloor(0, msg) + if got != 0 { + t.Fatalf("applyMediaFloor(0, system:) = %d, want 0 (system: sender must be exempt from floor)", got) + } +} + +// TestApplyMediaFloor_SubagentSenderExempt — Rule #3 internal-publisher exemption. +func TestApplyMediaFloor_SubagentSenderExempt(t *testing.T) { + msg := bus.InboundMessage{ + SenderID: "subagent:research", + Media: []bus.MediaFile{{Path: "/x"}}, + } + got := applyMediaFloor(0, msg) + if got != 0 { + t.Fatalf("applyMediaFloor(0, subagent:) = %d, want 0 (subagent: sender must be exempt from floor)", got) + } +} + +// TestResolveInboundDebounceDelay_FloorIsWiredInE2E — Rule #5. +// Exercises the OUTER resolveInboundDebounceDelay function (not the helper) with +// global debounce_ms=0, AgentStore=nil, media present, normal user sender. +// Asserts the floor fires. Locks in that any helper refactor still routes through +// the public function — prevents "floor exists but isn't wired" regressions. +func TestResolveInboundDebounceDelay_FloorIsWiredInE2E(t *testing.T) { + cfg := &config.Config{} + cfg.Gateway.InboundDebounceMs = 0 + deps := &ConsumerDeps{Cfg: cfg} + + msg := bus.InboundMessage{ + SenderID: "user-1", + AgentID: "", // skip AgentStore lookup path entirely + Media: []bus.MediaFile{{Path: "/x"}}, + } + + got := resolveInboundDebounceDelay(context.Background(), msg, deps) + want := time.Duration(mediaDebounceFloorMs) * time.Millisecond + if got != want { + t.Fatalf("resolveInboundDebounceDelay e2e = %v, want %v (floor not wired into outer function)", got, want) + } +} + +// TestResolveInboundDebounceDelay_NoMediaNoFloorE2E: outer function, no media, returns 0. +func TestResolveInboundDebounceDelay_NoMediaNoFloorE2E(t *testing.T) { + cfg := &config.Config{} + cfg.Gateway.InboundDebounceMs = 0 + deps := &ConsumerDeps{Cfg: cfg} + + msg := bus.InboundMessage{SenderID: "user-1"} + got := resolveInboundDebounceDelay(context.Background(), msg, deps) + if got != 0 { + t.Fatalf("resolveInboundDebounceDelay no-media = %v, want 0", got) + } +} + +// TestIsSystemOrSubagentSender_PrefixMatch covers helper directly. +func TestIsSystemOrSubagentSender_PrefixMatch(t *testing.T) { + cases := []struct { + senderID string + want bool + }{ + {"system:escalation", true}, + {"system:tool:sessions_send", true}, + {"subagent:research", true}, + {"user-1", false}, + {"", false}, + {"systemic", false}, // not the colon-prefix form + {"subagentX", false}, // not the colon-prefix form + } + for _, c := range cases { + if got := isSystemOrSubagentSender(c.senderID); got != c.want { + t.Errorf("isSystemOrSubagentSender(%q) = %v, want %v", c.senderID, got, c.want) + } + } +} diff --git a/cmd/gateway_consumer_dedup.go b/cmd/gateway_consumer_dedup.go new file mode 100644 index 00000000..dcb8a961 --- /dev/null +++ b/cmd/gateway_consumer_dedup.go @@ -0,0 +1,46 @@ +package cmd + +import ( + "fmt" + "strings" + + "github.com/nextlevelbuilder/goclaw/internal/bus" +) + +// dedupKeyFor builds the inbound-message dedup key. +// Format must match the inline key built in consumeInboundMessages +// (gateway_consumer.go) so siblings seeded via seedDedupFromMerged match +// the lookup performed on subsequent platform retransmits. +func dedupKeyFor(channel, senderID, chatID, messageID string) string { + return fmt.Sprintf("%s|%s|%s|%s", channel, senderID, chatID, messageID) +} + +// seedDedupFromMerged seeds every sibling message_id from the flushed merged +// message into the dedup cache. Implements Phase 1 Rule #4: a multi-attachment +// burst flushes as ONE merged InboundMessage carrying all source ids in +// metadata["merged_message_ids"]; any subsequent platform retransmit of a +// sibling member (e.g. Telegram webhook redelivery, WhatsApp retry) MUST be +// short-circuited at the dedup gate before re-entering the debouncer. +// +// No-op when: +// - dedupe is nil (defensive — call sites should always pass a real cache) +// - msg.Metadata is nil +// - merged_message_ids is absent or empty +func seedDedupFromMerged(dedupe *bus.DedupeCache, msg bus.InboundMessage) { + if dedupe == nil || msg.Metadata == nil { + return + } + merged := msg.Metadata["merged_message_ids"] + if merged == "" { + return + } + for _, id := range strings.Split(merged, ",") { + id = strings.TrimSpace(id) + if id == "" { + continue + } + // IsDuplicate has a side-effect of recording the key when absent, + // which is exactly the seeding we want. Discard the boolean. + _ = dedupe.IsDuplicate(dedupKeyFor(msg.Channel, msg.SenderID, msg.ChatID, id)) + } +} diff --git a/cmd/gateway_consumer_dedup_test.go b/cmd/gateway_consumer_dedup_test.go new file mode 100644 index 00000000..f541e358 --- /dev/null +++ b/cmd/gateway_consumer_dedup_test.go @@ -0,0 +1,67 @@ +package cmd + +import ( + "testing" + "time" + + "github.com/nextlevelbuilder/goclaw/internal/bus" +) + +// TestConsumerDedupSeedsAlbumSiblings — Rule #4 (Phase 1). +// When the debouncer flushes a merged message, all sibling message_ids +// (from metadata["merged_message_ids"]) MUST be seeded into the dedup cache +// so any platform retransmit of a sibling member is short-circuited. +func TestConsumerDedupSeedsAlbumSiblings(t *testing.T) { + dedupe := bus.NewDedupeCache(20*time.Minute, 5000) + + merged := bus.InboundMessage{ + Channel: "telegram", + ChatID: "chat-1", + SenderID: "user-1", + Metadata: map[string]string{ + "message_id": "m3", + "merged_message_ids": "m1,m2,m3", + }, + } + + seedDedupFromMerged(dedupe, merged) + + // Sibling retransmits MUST be reported as duplicates after seeding. + for _, mid := range []string{"m1", "m2", "m3"} { + key := dedupKeyFor("telegram", "user-1", "chat-1", mid) + if !dedupe.IsDuplicate(key) { + t.Fatalf("sibling %q not seeded; key=%q", mid, key) + } + } + + // Unrelated message_id must NOT be marked duplicate. + otherKey := dedupKeyFor("telegram", "user-1", "chat-1", "m99") + if dedupe.IsDuplicate(otherKey) { + t.Fatalf("unrelated message m99 wrongly marked duplicate") + } +} + +// TestConsumerDedupSeedsNoopWhenMergedEmpty: no merged_message_ids → no-op. +func TestConsumerDedupSeedsNoopWhenMergedEmpty(t *testing.T) { + dedupe := bus.NewDedupeCache(20*time.Minute, 5000) + msg := bus.InboundMessage{ + Channel: "telegram", + ChatID: "chat-1", + SenderID: "user-1", + Metadata: map[string]string{"message_id": "m1"}, + } + seedDedupFromMerged(dedupe, msg) // must not panic + + // m1 should NOT have been seeded by this helper (only the merged-list path seeds). + key := dedupKeyFor("telegram", "user-1", "chat-1", "m1") + if dedupe.IsDuplicate(key) { + t.Fatalf("m1 wrongly seeded for non-merged message") + } +} + +// TestConsumerDedupSeedsHandlesNilMetadata: nil metadata → no-op, no panic. +func TestConsumerDedupSeedsHandlesNilMetadata(t *testing.T) { + dedupe := bus.NewDedupeCache(20*time.Minute, 5000) + msg := bus.InboundMessage{Channel: "telegram", ChatID: "c", SenderID: "u"} + seedDedupFromMerged(dedupe, msg) // must not panic +} diff --git a/docs/05-channels-messaging.md b/docs/05-channels-messaging.md index 6e70ae68..bb7546cd 100644 --- a/docs/05-channels-messaging.md +++ b/docs/05-channels-messaging.md @@ -69,7 +69,17 @@ 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: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. +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. Command/control messages such as `/stop`, `/reset`, and system escalations bypass the debouncer. + +**Multi-attachment coalescing (#63).** Messages carrying attachments do NOT bypass the debouncer — that pre-fix shortcut was the source of N-replies for one user action. Instead, when media is present the effective window is `max(configured, mediaFloor)` so multi-file uploads land in the same buffer and flush together. Three surfaces apply the same invariant: + +| Surface | Buffer key | Trigger | +|---------|------------|---------| +| `internal/bus/inbound_debounce.go` | `(channel, chatID, senderID, agentID)` | Any inbound passing through the shared bus | +| `internal/gateway/methods/chat_debounce.go` | `(userKey, sessionKey)` | `/v1/chat/completions` streaming sessions | +| `internal/channels/telegram/album_aggregator.go` | `(chatID, MediaGroupID)` | Telegram media-group updates (album = N updates sharing one `MediaGroupID`) | + +The Telegram album aggregator runs at the channel layer after all access gates (mention, pairing, allow-list) pass — it buffers per `MediaGroupID`, pins the sender on first arrival as a security tripwire (mismatched sender → `security.album_sender_mismatch` + drop), and dispatches ONE call to the downstream pipeline on a 500ms silence window. `Channel.Stop()` synchronously drains pending albums before `pollCancel` so in-flight bursts always reach the agent loop. See `CONTRIBUTING.md` → "Multi-attachment coalescing" for the eight cross-surface invariants any new surface must honor. --- diff --git a/internal/bus/inbound_debounce.go b/internal/bus/inbound_debounce.go index 8df92cce..d23a06a9 100644 --- a/internal/bus/inbound_debounce.go +++ b/internal/bus/inbound_debounce.go @@ -48,21 +48,33 @@ func NewInboundDebouncerFunc(delayFn func(InboundMessage) time.Duration, flushFn } // Push adds a message to the debounce buffer. -// If debouncing is disabled or the message should bypass (media), it is flushed immediately. +// +// Behavior: +// - If delayFn returns > 0: append to per-key buffer and (re)set the silence timer. +// - If delayFn returns <= 0 AND a buffer already exists for the key: append the +// incoming message to the buffer and flush immediately (merge-then-flush). +// This is required so a no-media follow-up cannot bypass a buffered media +// message and trigger a duplicate agent run (issue #63). +// - If delayFn returns <= 0 AND no buffer exists: pass through immediately. +// +// There is no media-specific bypass — media-bearing messages go through the same +// path as text. The per-message delay decision lives in the caller's delayFn. func (d *InboundDebouncer) Push(msg InboundMessage) { debounceMs := d.delayFn(msg) - - // Disabled: pass through immediately. - if debounceMs <= 0 { - d.flushFn(msg) - return - } - key := debounceKey(msg) - // Media messages bypass debounce — flush any buffered text first, then process media. - if len(msg.Media) > 0 { - d.flushKey(key) + if debounceMs <= 0 { + // Disabled-path: merge into existing buffer if any, else pass through. + d.mu.Lock() + buf, exists := d.buffers[key] + if exists && len(buf.messages) > 0 { + buf.messages = append(buf.messages, msg) + d.mu.Unlock() + // flushKey re-locks, drains buffer, merges, and calls flushFn. + d.flushKey(key) + return + } + d.mu.Unlock() d.flushFn(msg) return } @@ -145,8 +157,15 @@ func debounceKey(msg InboundMessage) string { } // mergeInboundMessages combines multiple messages into one. -// Content is joined with newlines; media paths are concatenated; -// metadata and other fields come from the last message. +// +// Behavior: +// - Content: joined with newlines (matches TS upstream entries.map(e => e.body).join("\n")). +// - Media: concatenated in arrival order. +// - Metadata: starts from the last message (preserves legacy "latest wins" for +// channel-specific fields), then overlays metadata["merged_message_ids"] with +// a deduplicated arrival-ordered comma-separated list of all source message_ids. +// The consumer-level dedup uses this list to short-circuit retransmits of any +// sibling member (cmd/gateway_consumer.go). func mergeInboundMessages(msgs []InboundMessage) InboundMessage { if len(msgs) == 1 { return msgs[0] @@ -170,9 +189,62 @@ func mergeInboundMessages(msgs []InboundMessage) InboundMessage { } last.Media = allMedia + // Aggregate merged_message_ids — preserves arrival order, dedups. + merged := collectMergedMessageIDs(msgs) + if merged != "" { + if last.Metadata == nil { + last.Metadata = map[string]string{} + } else { + // Copy-on-write: don't mutate the input message's map. + cloned := make(map[string]string, len(last.Metadata)+1) + for k, v := range last.Metadata { + cloned[k] = v + } + last.Metadata = cloned + } + last.Metadata["merged_message_ids"] = merged + } + return last } +// collectMergedMessageIDs returns a comma-separated, arrival-ordered, deduplicated +// list of all source message_ids. Handles already-merged inputs by splitting any +// existing merged_message_ids entries back into individual IDs. +func collectMergedMessageIDs(msgs []InboundMessage) string { + seen := make(map[string]struct{}, len(msgs)) + ordered := make([]string, 0, len(msgs)) + + add := func(id string) { + id = strings.TrimSpace(id) + if id == "" { + return + } + if _, ok := seen[id]; ok { + return + } + seen[id] = struct{}{} + ordered = append(ordered, id) + } + + for _, m := range msgs { + if m.Metadata == nil { + continue + } + // Already-merged messages bring their full ID list. + if existing := m.Metadata["merged_message_ids"]; existing != "" { + for _, id := range strings.Split(existing, ",") { + add(id) + } + } + if mid := m.Metadata["message_id"]; mid != "" { + add(mid) + } + } + + return strings.Join(ordered, ",") +} + // truncateStr truncates a string to maxLen characters. func truncateStr(s string, maxLen int) string { if len(s) <= maxLen { diff --git a/internal/bus/inbound_debounce_test.go b/internal/bus/inbound_debounce_test.go index 59b2b1cf..0f1735a1 100644 --- a/internal/bus/inbound_debounce_test.go +++ b/internal/bus/inbound_debounce_test.go @@ -1,6 +1,8 @@ package bus import ( + "sort" + "strings" "testing" "time" ) @@ -83,28 +85,194 @@ func TestInboundDebouncerSeparatesAgents(t *testing.T) { } } -func TestInboundDebouncerMediaFlushesPendingTextFirst(t *testing.T) { - out := make(chan InboundMessage, 2) - d := NewInboundDebouncer(time.Minute, func(msg InboundMessage) { +// TestInboundDebouncerMergesMediaWithinWindow asserts two media messages within +// the debounce window merge into a single flush with all media in arrival order. +// Replaces the prior media-bypass test (issue #63 — red-team Phase 1 Rule #1 / #2). +func TestInboundDebouncerMergesMediaWithinWindow(t *testing.T) { + out := make(chan InboundMessage, 1) + d := NewInboundDebouncer(50*time.Millisecond, 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"}}, + Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", + Media: []MediaFile{{Path: "/tmp/a.png", MimeType: "image/png"}}, + }) + d.Push(InboundMessage{ + Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", + Media: []MediaFile{{Path: "/tmp/b.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) + got := waitInbound(t, out) + if len(got.Media) != 2 { + t.Fatalf("merged media count = %d, want 2 (paths: %#v)", len(got.Media), got.Media) } - if got := waitInbound(t, out); got.Content != "with media" || len(got.Media) != 1 { - t.Fatalf("second flush = %#v, want media message", got) + if got.Media[0].Path != "/tmp/a.png" || got.Media[1].Path != "/tmp/b.png" { + t.Fatalf("media order wrong: %#v", got.Media) } + + // Drain check — no second flush. + select { + case extra := <-out: + t.Fatalf("unexpected second flush: %#v", extra) + case <-time.After(80 * time.Millisecond): + } +} + +func TestInboundDebouncerFlushesMediaAfterSilence(t *testing.T) { + out := make(chan InboundMessage, 1) + d := NewInboundDebouncer(30*time.Millisecond, func(msg InboundMessage) { + out <- msg + }) + defer d.Stop() + + d.Push(InboundMessage{ + Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", + Media: []MediaFile{{Path: "/tmp/a.png"}}, + }) + + got := waitInbound(t, out) + if len(got.Media) != 1 { + t.Fatalf("media count = %d, want 1", len(got.Media)) + } +} + +func TestInboundDebouncerMergesTextThenMedia(t *testing.T) { + out := make(chan InboundMessage, 1) + d := NewInboundDebouncer(50*time.Millisecond, func(msg InboundMessage) { + out <- msg + }) + defer d.Stop() + + d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "look:"}) + d.Push(InboundMessage{ + Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", + Content: "this", Media: []MediaFile{{Path: "/tmp/a.png"}}, + }) + + got := waitInbound(t, out) + if got.Content != "look:\nthis" { + t.Fatalf("content = %q, want %q", got.Content, "look:\nthis") + } + if len(got.Media) != 1 { + t.Fatalf("media count = %d, want 1", len(got.Media)) + } + + select { + case extra := <-out: + t.Fatalf("unexpected second flush: %#v", extra) + case <-time.After(80 * time.Millisecond): + } +} + +func TestInboundDebouncerMergesMixedMediaAndText(t *testing.T) { + out := make(chan InboundMessage, 1) + d := NewInboundDebouncer(50*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", + Media: []MediaFile{{Path: "/tmp/a.png"}}, + }) + d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "two"}) + + got := waitInbound(t, out) + if got.Content != "one\ntwo" { + t.Fatalf("content = %q, want %q", got.Content, "one\ntwo") + } + if len(got.Media) != 1 { + t.Fatalf("media count = %d, want 1", len(got.Media)) + } +} + +// TestInboundDebouncerMergesNoMediaFollowupWithBufferedMedia — red-team Rule #1. +// A no-media message arriving with delay==0 while a buffered media message exists +// MUST merge into the buffer, not bypass it. +func TestInboundDebouncerMergesNoMediaFollowupWithBufferedMedia(t *testing.T) { + out := make(chan InboundMessage, 2) + // First push has media → returns floor 50ms. Second push has no media → returns 0. + d := NewInboundDebouncerFunc(func(m InboundMessage) time.Duration { + if len(m.Media) > 0 { + return 50 * time.Millisecond + } + return 0 + }, func(msg InboundMessage) { + out <- msg + }) + defer d.Stop() + + d.Push(InboundMessage{ + Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "caption", + Media: []MediaFile{{Path: "/tmp/a.png"}}, + }) + // Arrive while buffered (before 50ms window elapses). + time.Sleep(10 * time.Millisecond) + d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "ps"}) + + got := waitInbound(t, out) + if got.Content != "caption\nps" { + t.Fatalf("content = %q, want %q (follow-up must merge, not bypass)", got.Content, "caption\nps") + } + if len(got.Media) != 1 { + t.Fatalf("media count = %d, want 1", len(got.Media)) + } + + select { + case extra := <-out: + t.Fatalf("unexpected second flush — no-media follow-up bypassed instead of merging: %#v", extra) + case <-time.After(100 * time.Millisecond): + } +} + +// TestInboundDebouncerMergePopulatesMergedMessageIDs — red-team Rule #4 / Finding #10. +// Merged flush MUST populate metadata["merged_message_ids"] with all source IDs. +func TestInboundDebouncerMergePopulatesMergedMessageIDs(t *testing.T) { + out := make(chan InboundMessage, 1) + d := NewInboundDebouncer(40*time.Millisecond, func(msg InboundMessage) { + out <- msg + }) + defer d.Stop() + + d.Push(InboundMessage{ + Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "a", + Metadata: map[string]string{"message_id": "m1"}, + }) + d.Push(InboundMessage{ + Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "b", + Metadata: map[string]string{"message_id": "m2"}, + }) + d.Push(InboundMessage{ + Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "c", + Metadata: map[string]string{"message_id": "m3"}, + }) + + got := waitInbound(t, out) + merged := got.Metadata["merged_message_ids"] + if merged == "" { + t.Fatalf("merged_message_ids not populated; metadata = %#v", got.Metadata) + } + ids := strings.Split(merged, ",") + sort.Strings(ids) + want := []string{"m1", "m2", "m3"} + if !equalStringSlice(ids, want) { + t.Fatalf("merged_message_ids = %v, want (sorted) %v", ids, want) + } +} + +func equalStringSlice(a, b []string) bool { + if len(a) != len(b) { + return false + } + for i := range a { + if a[i] != b[i] { + return false + } + } + return true } func waitInbound(t *testing.T, ch <-chan InboundMessage) InboundMessage { diff --git a/internal/channels/telegram/album_aggregator.go b/internal/channels/telegram/album_aggregator.go new file mode 100644 index 00000000..7ff1a99b --- /dev/null +++ b/internal/channels/telegram/album_aggregator.go @@ -0,0 +1,175 @@ +package telegram + +import ( + "fmt" + "log/slog" + "sync" + "time" + + "github.com/mymmrac/telego" +) + +// Telegram delivers an album (multiple media items grouped on the client) as +// N separate Message updates, each carrying the same MediaGroupID. This +// aggregator buffers album members at the channel layer so downstream code +// sees ONE synthesized message per album — eliminating the original bug where +// N concurrent agent runs replied N times to a single user action. +// +// Contract (see plans/260528-1351-multi-attachment-debounce/phase-02 Rules): +// - Buffer key is (chatID, MediaGroupID). SenderID is pinned at buffer +// creation. A subsequent Push to the same key with a different senderID +// is dropped with a security warn (Rule #3) — Telegram does not reuse +// MediaGroupID across senders, but defense-in-depth catches spoofed +// updates. Plan Rule #2 originally specified a 3-tuple key including +// senderID; that made the rebind defense unreachable, so the 2-tuple +// key is used here with sender-pin as the runtime tripwire. +// - The first Push pins the representative's resolvedMessageContext on +// the buffer (Rule #5). Subsequent pushes pass their own rctx but it is +// discarded — only members[0]'s reply/thread/topic context flows downstream. +// - Caller MUST gate empty MediaGroupID; Push returns false if absent. +// - On every arrival the silence timer is Stop()+replaced via AfterFunc. +// No time.Timer.Reset() — avoids the documented double-fire race (Rule #7). +// - Dual caps with drop-and-log overflow (Rule #6). +// - Stop() flushes all pending buffers synchronously. Push after Stop is +// rejected (Rule #8 / post-shutdown straggler). +const ( + albumAggregatorWindow = 500 * time.Millisecond + albumAggregatorMaxBuffered = 100 // per album buffer (Telegram caps at 10; defensive) + albumAggregatorMaxBuffers = 1000 // global active buffers (DoS guard) +) + +type albumAggregator struct { + window time.Duration + maxPerBuf int + maxBuffers int + flushFn func(repCtx resolvedMessageContext, members []*telego.Message) + + mu sync.Mutex + buffers map[string]*albumBuffer + stopped bool +} + +type albumBuffer struct { + repCtx resolvedMessageContext // captured from members[0] + members []*telego.Message + senderID int64 // pinned on first Push + timer *time.Timer +} + +func newAlbumAggregator(window time.Duration, maxPerBuf, maxBuffers int, flushFn func(resolvedMessageContext, []*telego.Message)) *albumAggregator { + return &albumAggregator{ + window: window, + maxPerBuf: maxPerBuf, + maxBuffers: maxBuffers, + flushFn: flushFn, + buffers: make(map[string]*albumBuffer), + } +} + +// albumKey returns the buffer key plus the sender ID. ok=false if the +// message is missing required fields (no MediaGroupID, no From). Caller must +// short-circuit when ok=false. +func albumKey(msg *telego.Message) (key string, senderID int64, ok bool) { + if msg == nil || msg.MediaGroupID == "" || msg.From == nil { + return "", 0, false + } + return fmt.Sprintf("%d:%s", msg.Chat.ID, msg.MediaGroupID), msg.From.ID, true +} + +// Push accepts an album member. Returns true if accepted, false otherwise +// (empty MediaGroupID, post-stop, global overflow, sender-rebind, per-buffer +// overflow). On false the caller MUST fall through to single-message dispatch +// so the message is not silently lost. +// +// rctx is the resolved-message context for THIS message; it is stored only +// on the first Push for a key (members[0] is the representative per Rule #5). +func (a *albumAggregator) Push(msg *telego.Message, rctx resolvedMessageContext) bool { + key, senderID, ok := albumKey(msg) + if !ok { + return false + } + + a.mu.Lock() + if a.stopped { + a.mu.Unlock() + slog.Warn("telegram.album_post_shutdown_push", "key", key, "message_id", msg.MessageID) + return false + } + + buf, exists := a.buffers[key] + if !exists { + if len(a.buffers) >= a.maxBuffers { + a.mu.Unlock() + slog.Warn("telegram.album_overflow", + "scope", "global", "max", a.maxBuffers, "key", key) + return false + } + buf = &albumBuffer{senderID: senderID, repCtx: rctx} + a.buffers[key] = buf + } else { + if buf.senderID != senderID { + a.mu.Unlock() + slog.Warn("security.album_sender_mismatch", + "key", key, "expected_sender", buf.senderID, "got_sender", senderID, + "message_id", msg.MessageID) + return false + } + if len(buf.members) >= a.maxPerBuf { + a.mu.Unlock() + slog.Warn("telegram.album_overflow", + "scope", "buffer", "max", a.maxPerBuf, "key", key, "message_id", msg.MessageID) + return false + } + } + + buf.members = append(buf.members, msg) + if buf.timer != nil { + buf.timer.Stop() + } + buf.timer = time.AfterFunc(a.window, func() { a.flushKey(key) }) + a.mu.Unlock() + return true +} + +// flushKey drains the named buffer and invokes flushFn outside the lock. +// Safe to call multiple times — second call is a no-op. +func (a *albumAggregator) flushKey(key string) { + a.mu.Lock() + buf, ok := a.buffers[key] + if !ok { + a.mu.Unlock() + return + } + if buf.timer != nil { + buf.timer.Stop() + } + members := buf.members + repCtx := buf.repCtx + delete(a.buffers, key) + a.mu.Unlock() + + if len(members) == 0 { + return + } + a.flushFn(repCtx, members) +} + +// Stop marks the aggregator as stopped and synchronously flushes all pending +// buffers. After Stop, Push returns false and logs a warn. Idempotent. +func (a *albumAggregator) Stop() { + a.mu.Lock() + if a.stopped { + a.mu.Unlock() + return + } + a.stopped = true + keys := make([]string, 0, len(a.buffers)) + for k := range a.buffers { + keys = append(keys, k) + } + a.mu.Unlock() + + for _, k := range keys { + a.flushKey(k) + } +} diff --git a/internal/channels/telegram/album_aggregator_test.go b/internal/channels/telegram/album_aggregator_test.go new file mode 100644 index 00000000..594f32ec --- /dev/null +++ b/internal/channels/telegram/album_aggregator_test.go @@ -0,0 +1,262 @@ +package telegram + +import ( + "fmt" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/mymmrac/telego" +) + +// rctxZero returns a zero-valued resolvedMessageContext for tests that don't +// care about downstream dispatch (aggregator unit tests). +func rctxZero() resolvedMessageContext { return resolvedMessageContext{} } + +// mkAlbumMsg builds a minimal *telego.Message for aggregator tests. +func mkAlbumMsg(chatID, userID int64, groupID string, msgID int) *telego.Message { + return &telego.Message{ + MessageID: msgID, + Chat: telego.Chat{ID: chatID}, + From: &telego.User{ID: userID}, + MediaGroupID: groupID, + } +} + +type capturedFlushes struct { + mu sync.Mutex + batches [][]*telego.Message + done chan struct{} +} + +func newCapturedFlushes() *capturedFlushes { + return &capturedFlushes{done: make(chan struct{}, 64)} +} + +func (c *capturedFlushes) callback() func(resolvedMessageContext, []*telego.Message) { + return func(_ resolvedMessageContext, members []*telego.Message) { + c.mu.Lock() + batch := append([]*telego.Message(nil), members...) + c.batches = append(c.batches, batch) + c.mu.Unlock() + select { + case c.done <- struct{}{}: + default: + } + } +} + +func (c *capturedFlushes) batchCount() int { + c.mu.Lock() + defer c.mu.Unlock() + return len(c.batches) +} + +func (c *capturedFlushes) batch(i int) []*telego.Message { + c.mu.Lock() + defer c.mu.Unlock() + if i >= len(c.batches) { + return nil + } + return c.batches[i] +} + +func waitFlush(t *testing.T, c *capturedFlushes, timeout time.Duration) { + t.Helper() + select { + case <-c.done: + case <-time.After(timeout): + t.Fatalf("timed out waiting for album flush after %s", timeout) + } +} + +func TestAlbumAggregator_FlushesOnSilence(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(50*time.Millisecond, 100, 1000, cap.callback()) + defer a.Stop() + + if !a.Push(mkAlbumMsg(1, 10, "g1", 100), rctxZero()) { + t.Fatal("push 1 rejected") + } + time.Sleep(15 * time.Millisecond) + a.Push(mkAlbumMsg(1, 10, "g1", 101), rctxZero()) + time.Sleep(15 * time.Millisecond) + a.Push(mkAlbumMsg(1, 10, "g1", 102), rctxZero()) + + waitFlush(t, cap, 500*time.Millisecond) + time.Sleep(80 * time.Millisecond) + + if cap.batchCount() != 1 { + t.Fatalf("flushes=%d, want 1", cap.batchCount()) + } + got := cap.batch(0) + if len(got) != 3 { + t.Fatalf("members=%d, want 3", len(got)) + } + if got[0].MessageID != 100 || got[1].MessageID != 101 || got[2].MessageID != 102 { + t.Fatalf("arrival order broken: %d %d %d", got[0].MessageID, got[1].MessageID, got[2].MessageID) + } +} + +func TestAlbumAggregator_DistinctMediaGroupsAreSeparateBuffers(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(40*time.Millisecond, 100, 1000, cap.callback()) + defer a.Stop() + + a.Push(mkAlbumMsg(1, 10, "g1", 100), rctxZero()) + a.Push(mkAlbumMsg(1, 20, "g2", 200), rctxZero()) + + waitFlush(t, cap, 500*time.Millisecond) + waitFlush(t, cap, 500*time.Millisecond) + time.Sleep(60 * time.Millisecond) + + if cap.batchCount() != 2 { + t.Fatalf("flushes=%d, want 2 (one per MediaGroupID)", cap.batchCount()) + } +} + +func TestAlbumAggregator_SenderRebindDropsAndWarns(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(40*time.Millisecond, 100, 1000, cap.callback()) + defer a.Stop() + + if !a.Push(mkAlbumMsg(1, 10, "g1", 100), rctxZero()) { + t.Fatal("first push should succeed") + } + if a.Push(mkAlbumMsg(1, 999, "g1", 101), rctxZero()) { + t.Fatal("sender-rebind push should be rejected") + } + + waitFlush(t, cap, 500*time.Millisecond) + time.Sleep(60 * time.Millisecond) + + if cap.batchCount() != 1 { + t.Fatalf("flushes=%d, want 1", cap.batchCount()) + } + got := cap.batch(0) + if len(got) != 1 || got[0].MessageID != 100 { + t.Fatalf("buffer contaminated by rebind: %v", got) + } +} + +func TestAlbumAggregator_ResetsOnNewArrival(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(80*time.Millisecond, 100, 1000, cap.callback()) + defer a.Stop() + + for i := 0; i < 4; i++ { + a.Push(mkAlbumMsg(1, 10, "g1", 100+i), rctxZero()) + time.Sleep(50 * time.Millisecond) + } + + waitFlush(t, cap, 500*time.Millisecond) + time.Sleep(120 * time.Millisecond) + + if cap.batchCount() != 1 { + t.Fatalf("flushes=%d, want 1 (timer should reset on each push)", cap.batchCount()) + } + got := cap.batch(0) + if len(got) != 4 { + t.Fatalf("members=%d, want 4", len(got)) + } +} + +func TestAlbumAggregator_StopFlushesPendingImmediately(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(time.Minute, 100, 1000, cap.callback()) + + a.Push(mkAlbumMsg(1, 10, "g1", 100), rctxZero()) + a.Push(mkAlbumMsg(1, 10, "g1", 101), rctxZero()) + + a.Stop() + + if cap.batchCount() != 1 { + t.Fatalf("flushes=%d, want 1 (Stop must flush pending)", cap.batchCount()) + } + if got := cap.batch(0); len(got) != 2 { + t.Fatalf("members=%d, want 2", len(got)) + } +} + +func TestAlbumAggregator_RespectsPerBufferCap(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(time.Minute, 5, 1000, cap.callback()) + defer a.Stop() + + for i := 0; i < 5; i++ { + if !a.Push(mkAlbumMsg(1, 10, "g1", 100+i), rctxZero()) { + t.Fatalf("push %d rejected before cap reached", i) + } + } + if a.Push(mkAlbumMsg(1, 10, "g1", 999), rctxZero()) { + t.Fatal("push past per-buffer cap should be rejected") + } + if cap.batchCount() != 0 { + t.Fatalf("flushes=%d, want 0 (cap must not cause early flush)", cap.batchCount()) + } +} + +func TestAlbumAggregator_RespectsGlobalBufferCap(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(time.Minute, 100, 3, cap.callback()) + defer a.Stop() + + for i := 0; i < 3; i++ { + if !a.Push(mkAlbumMsg(int64(i+1), 10, fmt.Sprintf("g%d", i), 100), rctxZero()) { + t.Fatalf("push %d rejected before global cap reached", i) + } + } + if a.Push(mkAlbumMsg(99, 10, "g99", 100), rctxZero()) { + t.Fatal("push past global buffer cap should be rejected") + } +} + +func TestAlbumAggregator_EmptyMediaGroupIDRejected(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(40*time.Millisecond, 100, 1000, cap.callback()) + defer a.Stop() + + msg := mkAlbumMsg(1, 10, "", 100) + if a.Push(msg, rctxZero()) { + t.Fatal("Push with empty MediaGroupID must return false (caller must gate)") + } +} + +func TestAlbumAggregator_PostStopPushIgnored(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(40*time.Millisecond, 100, 1000, cap.callback()) + + a.Stop() + if a.Push(mkAlbumMsg(1, 10, "g1", 100), rctxZero()) { + t.Fatal("Push after Stop must be rejected") + } +} + +func TestAlbumAggregator_TimerNoLeakAfterFlush(t *testing.T) { + cap := newCapturedFlushes() + a := newAlbumAggregator(30*time.Millisecond, 100, 1000, cap.callback()) + defer a.Stop() + + var firedFlushes int32 + originalFlush := a.flushFn + a.flushFn = func(rctx resolvedMessageContext, members []*telego.Message) { + atomic.AddInt32(&firedFlushes, 1) + originalFlush(rctx, members) + } + + a.Push(mkAlbumMsg(1, 10, "g1", 100), rctxZero()) + + waitFlush(t, cap, 500*time.Millisecond) + time.Sleep(120 * time.Millisecond) + + a.mu.Lock() + bufCount := len(a.buffers) + a.mu.Unlock() + if bufCount != 0 { + t.Fatalf("buffer not released after flush: %d remaining", bufCount) + } + if got := atomic.LoadInt32(&firedFlushes); got != 1 { + t.Fatalf("flushFn fired %d times, want 1", got) + } +} diff --git a/internal/channels/telegram/channel.go b/internal/channels/telegram/channel.go index ba81ed4e..5ba967c2 100644 --- a/internal/channels/telegram/channel.go +++ b/internal/channels/telegram/channel.go @@ -41,12 +41,14 @@ type Channel struct { threadIDs sync.Map // localKey string → messageThreadID int (for forum topic routing) mentionMode string // "strict" (default) or "yield" botDisplayName string // bot's first_name from GetMe (e.g. "ViệtBot"); captured once at Start + pollCtx context.Context // long-polling context (cancelled by pollCancel); promoted from Start-local so background helpers (e.g. albumAggregator) can derive from it pollCancel context.CancelFunc // cancels the long polling context pollDone chan struct{} // closed when polling goroutine exits handlerWg sync.WaitGroup // tracks in-flight handler goroutines for graceful shutdown handlerSem chan struct{} // bounded semaphore for concurrent handler goroutines pendingDraftID sync.Map // localKey string → int (draftID) audioMgr *audio.Manager // unified STT via audio.Manager (nil = no STT) + albumAgg *albumAggregator // coalesces Telegram album members into a single dispatch; nil before Start writerHealMu sync.Mutex // guards writerHealLastTry for /writers self-heal writerHealLastTry map[string]time.Time // key "chatID|userID" → last attempt timestamp // pairingService, approvedGroups, pairingDebounce, groupHistory, historyLimit, requireMention @@ -206,10 +208,31 @@ func (c *Channel) Start(ctx context.Context) error { // Create a cancellable context for the polling goroutine. // Stop() cancels this context to cleanly shut down long polling. - pollCtx, cancel := context.WithCancel(ctx) - c.pollCancel = cancel + c.pollCtx, c.pollCancel = context.WithCancel(ctx) + pollCtx := c.pollCtx + cancel := c.pollCancel c.pollDone = make(chan struct{}) + // Album aggregator coalesces Telegram media-group updates into ONE dispatch. + // flushFn closure captures c.pollCtx so silence-window flushes have a valid + // context for downstream media resolution + bus publish. Stop() drains + // synchronously BEFORE pollCancel so flushes never race ctx cancellation. + // handlerWg participation is REQUIRED: AfterFunc-fired flushes run on a + // dedicated goroutine NOT tracked by the polling loop. Without explicit + // Add/Done the Stop() handlerWg.Wait() would race the timer-spawned + // dispatch, breaking the "always publish in-flight bursts" invariant + // documented in CHANGELOG/docs/05-channels-messaging.md. + c.albumAgg = newAlbumAggregator( + albumAggregatorWindow, + albumAggregatorMaxBuffered, + albumAggregatorMaxBuffers, + func(rctx resolvedMessageContext, members []*telego.Message) { + c.handlerWg.Add(1) + defer c.handlerWg.Done() + c.processResolvedMessage(c.pollCtx, rctx, members) + }, + ) + updates, err := c.bot.UpdatesViaLongPolling(pollCtx, &telego.GetUpdatesParams{ Timeout: 25, // Long-poll seconds; keep below HTTP client Timeout (#361) AllowedUpdates: []string{ @@ -379,6 +402,12 @@ func (c *Channel) Stop(_ context.Context) error { gh.StopFlusher() } + // Drain pending album buffers BEFORE cancelling pollCtx so any synchronous + // flushFn callbacks still see a valid context for downstream dispatch. + if c.albumAgg != nil { + c.albumAgg.Stop() + } + if c.pollCancel != nil { c.pollCancel() } diff --git a/internal/channels/telegram/handlers.go b/internal/channels/telegram/handlers.go index 685d621f..bd92c93e 100644 --- a/internal/channels/telegram/handlers.go +++ b/internal/channels/telegram/handlers.go @@ -389,12 +389,68 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { } } + // All gates passed — hand off to processResolvedMessage for media resolution + // and downstream dispatch. members=[message] for single-message path; the + // album branch below coalesces Telegram media-group members so downstream + // sees ONE synthesized message per album (fixes #63). + rctx := resolvedMessageContext{ + content: content, + userID: userID, + senderID: senderID, + senderLabel: senderLabel, + chatID: chatID, + chatIDStr: chatIDStr, + localKey: localKey, + isGroup: isGroup, + isForum: isForum, + messageThreadID: messageThreadID, + dmThreadID: dmThreadID, + topicCfg: topicCfg, + } + + // Album coalescing: if this update carries a MediaGroupID, push to the + // aggregator and return — the silence-window flush dispatches ONE call + // with all members. Push returns false on empty MediaGroupID (this branch + // won't be entered), aggregator stopped (post-shutdown), sender-rebind + // security drop, per-buffer overflow, or global overflow. In those cases + // fall through to single-message dispatch so the message is never silently + // lost. + if message.MediaGroupID != "" && c.albumAgg != nil { + if c.albumAgg.Push(message, rctx) { + return + } + } + c.processResolvedMessage(ctx, rctx, []*telego.Message{message}) +} + +// processResolvedMessage handles the post-gate phase of an inbound: media +// resolution → content enrichment → typing → metadata → PublishInbound. It is +// the shared entry point for single-message dispatch (members = []{rep}) and +// album-flush dispatch (members = N buffered album members; rep = members[0]). +// +// Album behavior: media is resolved for every member, concatenated in arrival +// order; reply context / forward context / caption come from members[0] only +// (Telegram puts these on the first album message). +func (c *Channel) processResolvedMessage(ctx context.Context, rctx resolvedMessageContext, members []*telego.Message) { + if len(members) == 0 { + return + } + rep := members[0] + user := rep.From + content := rctx.content + // --- Media download (only when bot will process the message) --- // Deferred until after mention + pairing gates to avoid downloading // media for messages that only get recorded in pending history. - mediaList, mediaErrors := c.resolveMedia(ctx, message) - if message.ReplyToMessage != nil { - replyMedia, replyErrors := c.resolveMedia(ctx, message.ReplyToMessage) + var mediaList []MediaInfo + var mediaErrors []MediaError + for _, m := range members { + ml, me := c.resolveMedia(ctx, m) + mediaList = append(mediaList, ml...) + mediaErrors = append(mediaErrors, me...) + } + if rep.ReplyToMessage != nil { + replyMedia, replyErrors := c.resolveMedia(ctx, rep.ReplyToMessage) if len(replyMedia) > 0 { // Tag reply media so LLM knows which images came from the replied-to message. for i := range replyMedia { @@ -403,7 +459,7 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { // Reply media first (context), current media second. mediaList = append(replyMedia, mediaList...) slog.Debug("telegram: resolved media from replied message", - "reply_msg_id", message.ReplyToMessage.MessageID, + "reply_msg_id", rep.ReplyToMessage.MessageID, "media_count", len(replyMedia), ) } @@ -458,7 +514,7 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { // Replace lightweight media tags with full tags (includes transcripts). fullTags := buildMediaTags(mediaList) - lightTags := lightweightMediaTags(message) + lightTags := lightweightMediaTags(rep) if lightTags != "" && fullTags != "" { content = strings.Replace(content, lightTags, fullTags, 1) } else if fullTags != "" { @@ -479,7 +535,7 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { if len(mediaErrors) > 0 { for _, me := range mediaErrors { errTag := fmt.Sprintf("[sent media (%s) — skipped: %s]", me.Type, me.Reason) - if lightTag := lightweightTagForType(me.Type, message); lightTag != "" { + if lightTag := lightweightTagForType(me.Type, rep); lightTag != "" { content = strings.Replace(content, lightTag, errTag, 1) } else { content = errTag + "\n" + content @@ -494,24 +550,24 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { } else { errText = "⚠️ Failed to download the attached file. Skipped." } - _ = c.sendHTML(ctx, chatID, errText, 0, messageThreadID) + _ = c.sendHTML(ctx, rctx.chatID, errText, 0, rctx.messageThreadID) } } slog.Debug("telegram message received", - "sender_id", senderID, - "chat_id", fmt.Sprintf("%d", chatID), + "sender_id", rctx.senderID, + "chat_id", rctx.chatIDStr, "preview", channels.Truncate(content, 50), ) // Build context from pending group history (if any). // Annotate current message with sender name so LLM knows who is talking. finalContent := content - if isGroup { - annotated := fmt.Sprintf("[From: %s]\n%s", senderLabel, content) + if rctx.isGroup { + annotated := fmt.Sprintf("[From: %s]\n%s", rctx.senderLabel, content) if c.HistoryLimit() > 0 { // Resolve deferred media from history entries (lazy download). - if histRefs := c.GroupHistory().CollectMediaRefs(localKey); len(histRefs) > 0 { + if histRefs := c.GroupHistory().CollectMediaRefs(rctx.localKey); len(histRefs) > 0 { histMedia, histErrors := c.resolveMediaRefs(ctx, histRefs) if len(histMedia) > 0 { mediaFiles = prependMediaInfoFiles(mediaFiles, histMedia) @@ -523,39 +579,39 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { "type", e.Type, "reason", e.Reason) } } - finalContent = c.GroupHistory().BuildContext(localKey, annotated, c.HistoryLimit()) + finalContent = c.GroupHistory().BuildContext(rctx.localKey, annotated, c.HistoryLimit()) } else { finalContent = annotated } } else { // DM: annotate with sender identity so the agent knows who is messaging. - finalContent = fmt.Sprintf("[From: %s]\n%s", senderLabel, content) + finalContent = fmt.Sprintf("[From: %s]\n%s", rctx.senderLabel, content) } // Send typing indicator with keepalive + TTL safety net. // Telegram typing expires after 5s, so keepalive every 4s. // TTL auto-stops after 60s to prevent stuck indicators. - chatIDObj := tu.ID(chatID) + chatIDObj := tu.ID(rctx.chatID) typingCtrl := typing.New(typing.Options{ MaxDuration: 60 * time.Second, KeepaliveInterval: 4 * time.Second, StartFn: func() error { action := tu.ChatAction(chatIDObj, telego.ChatActionTyping) - if messageThreadID > 0 { - action.MessageThreadID = messageThreadID + if rctx.messageThreadID > 0 { + action.MessageThreadID = rctx.messageThreadID } return c.bot.SendChatAction(ctx, action) }, }) // Stop previous typing controller for this chat/topic (if any) - if prev, ok := c.typingCtrls.Load(localKey); ok { + if prev, ok := c.typingCtrls.Load(rctx.localKey); ok { prev.(*typing.Controller).Stop() } - c.typingCtrls.Store(localKey, typingCtrl) + c.typingCtrls.Store(rctx.localKey, typingCtrl) typingCtrl.Start() // Stop previous thinking animation for this chat/topic - if prevStop, ok := c.stopThinking.Load(localKey); ok { + if prevStop, ok := c.stopThinking.Load(rctx.localKey); ok { if cf, ok := prevStop.(*thinkingCancel); ok { cf.Cancel() } @@ -563,7 +619,7 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { // Create thinking cancel for this chat/topic _, thinkCancel := context.WithCancel(ctx) - c.stopThinking.Store(localKey, &thinkingCancel{fn: thinkCancel}) + c.stopThinking.Store(rctx.localKey, &thinkingCancel{fn: thinkCancel}) // No "Thinking..." placeholder — the DraftStream creates its own message // on the first streaming chunk (sendMessage on first flush). @@ -571,23 +627,35 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { // user sees typing indicator → first content appears directly. metadata := map[string]string{ - "message_id": fmt.Sprintf("%d", message.MessageID), + "message_id": fmt.Sprintf("%d", rep.MessageID), "user_id": fmt.Sprintf("%d", user.ID), tools.MetaUsername: user.Username, "first_name": user.FirstName, - "is_group": fmt.Sprintf("%t", isGroup), - "local_key": localKey, + "is_group": fmt.Sprintf("%t", rctx.isGroup), + "local_key": rctx.localKey, } - if message.Chat.Title != "" { - metadata[tools.MetaChatTitle] = message.Chat.Title + // When this publish coalesces multiple platform messages (album members), + // seed every sibling MessageID into merged_message_ids so the consumer + // dedup (cmd/gateway_consumer_dedup.go) blocks Telegram retransmits of + // non-representative members from triggering a duplicate agent run. + // Reuses the same key as the bus debouncer's merge path — single sink. + if len(members) > 1 { + ids := make([]string, 0, len(members)) + for _, m := range members { + ids = append(ids, fmt.Sprintf("%d", m.MessageID)) + } + metadata["merged_message_ids"] = strings.Join(ids, ",") } - if isForum { + if rep.Chat.Title != "" { + metadata[tools.MetaChatTitle] = rep.Chat.Title + } + if rctx.isForum { metadata[tools.MetaIsForum] = "true" - metadata[tools.MetaMessageThreadID] = fmt.Sprintf("%d", messageThreadID) + metadata[tools.MetaMessageThreadID] = fmt.Sprintf("%d", rctx.messageThreadID) } - if dmThreadID > 0 { - metadata[tools.MetaDMThreadID] = fmt.Sprintf("%d", dmThreadID) - metadata[tools.MetaMessageThreadID] = fmt.Sprintf("%d", dmThreadID) + if rctx.dmThreadID > 0 { + metadata[tools.MetaDMThreadID] = fmt.Sprintf("%d", rctx.dmThreadID) + metadata[tools.MetaMessageThreadID] = fmt.Sprintf("%d", rctx.dmThreadID) } // Self-identity hint so the LLM knows its own Telegram handle and does not // confuse other bots' @mentions (preserved after stripBotMention) for its own. @@ -595,15 +663,15 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { metadata[tools.MetaChannelSelfIdentity] = identity } - if topicCfg.systemPrompt != "" { - metadata[tools.MetaTopicSystemPrompt] = topicCfg.systemPrompt + if rctx.topicCfg.systemPrompt != "" { + metadata[tools.MetaTopicSystemPrompt] = rctx.topicCfg.systemPrompt } - if topicCfg.skills != nil { - metadata[tools.MetaTopicSkills] = strings.Join(topicCfg.skills, ",") + if rctx.topicCfg.skills != nil { + metadata[tools.MetaTopicSkills] = strings.Join(rctx.topicCfg.skills, ",") } peerKind := "direct" - if isGroup { + if rctx.isGroup { peerKind = "group" } @@ -626,35 +694,35 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) { // Collect contact for processed messages (DM + group-mentioned). if cc := c.ContactCollector(); cc != nil { contactName := strings.TrimSpace(user.FirstName + " " + user.LastName) - cc.EnsureContact(ctx, c.Type(), c.Name(), senderID, userID, contactName, user.Username, peerKind, "user", "", "") + cc.EnsureContact(ctx, c.Type(), c.Name(), rctx.senderID, rctx.userID, contactName, user.Username, peerKind, "user", "", "") // Also collect group chat itself as a contact (for group permission / merge). - if isGroup { - cc.EnsureContact(ctx, c.Type(), c.Name(), chatIDStr, "", message.Chat.Title, "", "group", "group", "", "") + if rctx.isGroup { + cc.EnsureContact(ctx, c.Type(), c.Name(), rctx.chatIDStr, "", rep.Chat.Title, "", "group", "group", "", "") // Collect forum topic as a distinct delivery target (including General). - if isForum && messageThreadID > 0 { - threadStr := fmt.Sprintf("%d", messageThreadID) - cc.EnsureContact(ctx, c.Type(), c.Name(), chatIDStr, "", message.Chat.Title, "", "group", "topic", threadStr, "topic") + if rctx.isForum && rctx.messageThreadID > 0 { + threadStr := fmt.Sprintf("%d", rctx.messageThreadID) + cc.EnsureContact(ctx, c.Type(), c.Name(), rctx.chatIDStr, "", rep.Chat.Title, "", "group", "topic", threadStr, "topic") } } } c.Bus().PublishInbound(bus.InboundMessage{ Channel: c.Name(), - SenderID: senderID, - ChatID: chatIDStr, + SenderID: rctx.senderID, + ChatID: rctx.chatIDStr, Content: finalContent, Media: mediaFiles, PeerKind: peerKind, - UserID: userID, + UserID: rctx.userID, AgentID: targetAgentID, HistoryLimit: c.HistoryLimit(), - ToolAllow: topicCfg.tools, + ToolAllow: rctx.topicCfg.tools, TenantID: c.TenantID(), Metadata: metadata, }) // Clear pending history after sending to agent. - if isGroup { - c.GroupHistory().Clear(localKey) + if rctx.isGroup { + c.GroupHistory().Clear(rctx.localKey) } } diff --git a/internal/channels/telegram/resolved_message_context.go b/internal/channels/telegram/resolved_message_context.go new file mode 100644 index 00000000..89c8ab3e --- /dev/null +++ b/internal/channels/telegram/resolved_message_context.go @@ -0,0 +1,38 @@ +package telegram + +// resolvedMessageContext bundles the post-gate, post-mention-strip state that +// `handleMessage` computes BEFORE media resolution and downstream dispatch. +// It is the input contract for `processResolvedMessage` (single message) and +// for `dispatchAlbum` (album members[0] as the representative). +// +// Captured at gate-pass time for the representative message: +// - identity: userID, senderID, senderLabel +// - chat addressing: chatID, chatIDStr, localKey +// - thread metadata: isGroup, isForum, messageThreadID, dmThreadID +// - resolved cfg: topicCfg (groupPolicy, systemPrompt, skills, tools, allowFrom, ...) +// - cleaned content: content (text + caption + lightweight tags + reply/forward/location +// enrichment + stripBotMention applied). For an album flush, this +// is members[0]'s content snapshot — Telegram puts captions on +// the first album message only. +// +// The carrier *telego.Message is NOT a field here; callers pass it alongside +// (single: []{message}, album: members). Keeping the message out of the struct +// lets the album path swap a list of N messages in cleanly without per-member +// rctx copies. +type resolvedMessageContext struct { + content string + userID string + senderID string + senderLabel string + + chatID int64 + chatIDStr string + localKey string + + isGroup bool + isForum bool + messageThreadID int + dmThreadID int + + topicCfg resolvedTopicConfig +} diff --git a/internal/config/config_channels.go b/internal/config/config_channels.go index 6a9c5c70..e3194a3f 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 (0 = no wait) + InboundDebounceMs int `json:"inbound_debounce_ms,omitempty"` // silence-window in ms that merges rapid channel/Web Chat messages from the same sender/session; 0 disables for text but media-bearing messages still honor a built-in media floor so multi-attachment bursts (#63) coalesce into a single agent run. Agents may override via per-agent agent_config.inbound_debounce_ms. 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 3e02e70e..ff8572a0 100644 --- a/internal/gateway/methods/chat.go +++ b/internal/gateway/methods/chat.go @@ -218,16 +218,18 @@ func (m *ChatMethods) handleSend(ctx context.Context, client *gateway.Client, re 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, loop.OtherConfig()); delay > 0 { + // Media-bearing sends route through the same debouncer path as text. + // The media floor in chatDebounceDelay guarantees a non-zero window when + // the operator has disabled debouncing, so multi-attachment bursts coalesce + // into a single dispatch (issue #63). + hasMedia := len(params.parseMedia()) > 0 + delay := chatDebounceDelay(m.cfg, loop.OtherConfig(), hasMedia) + if delay > 0 { m.debouncer.Push(debounceKey, delay, item) return } - m.dispatchChatSends([]chatSendRequest{item}) + // delay == 0: Push merges into existing buffer (if any) or dispatches. + m.debouncer.Push(debounceKey, 0, item) } func (m *ChatMethods) abortChatSession(reqID string, client *gateway.Client, sessionKey string) { diff --git a/internal/gateway/methods/chat_debounce.go b/internal/gateway/methods/chat_debounce.go index 513a405c..8d93d195 100644 --- a/internal/gateway/methods/chat_debounce.go +++ b/internal/gateway/methods/chat_debounce.go @@ -13,6 +13,12 @@ import ( "github.com/nextlevelbuilder/goclaw/internal/store" ) +// chatMediaDebounceFloorMs is the minimum debounce window applied to Web Chat +// sends that carry media when the post-override delay would otherwise be 0. +// Mirrors cmd/gateway_consumer_debounce.go mediaDebounceFloorMs (Phase 1) — +// duplicated by value (1000ms) to keep gateway/methods decoupled from cmd. +const chatMediaDebounceFloorMs = 1000 + type chatSendRequest struct { ctx context.Context client *gateway.Client @@ -41,8 +47,31 @@ func newChatDebouncer(flushFn func([]chatSendRequest)) *chatDebouncer { } } +// Push appends an item to the per-key buffer. +// +// Behavior: +// - delay > 0: append + (re)set the silence timer (existing buffered path). +// - delay <= 0 AND a buffer already exists with items: append the incoming +// item to the buffer and flush immediately (merge-then-flush). Required so +// a no-media follow-up cannot bypass a buffered media chat-send and trigger +// a duplicate dispatch (Phase 1.5 Rule #4, mirrors bus debouncer Rule #1). +// - delay <= 0 AND no buffer exists: dispatch immediately (passthrough). func (d *chatDebouncer) Push(key string, delay time.Duration, item chatSendRequest) { if delay <= 0 { + d.mu.Lock() + buf, exists := d.buffers[key] + if exists && len(buf.items) > 0 { + buf.items = append(buf.items, item) + if buf.timer != nil { + buf.timer.Stop() + } + items := buf.items + delete(d.buffers, key) + d.mu.Unlock() + d.flushFn(items) + return + } + d.mu.Unlock() d.flushFn([]chatSendRequest{item}) return } @@ -129,7 +158,13 @@ func mergeChatSendRequests(items []chatSendRequest) chatSendParams { return last } -func chatDebounceDelay(cfg *config.Config, agentOtherConfig json.RawMessage) time.Duration { +// chatDebounceDelay computes the per-send debounce window. +// +// Precedence: agent override (when set) overrides the global config. The media +// floor fires ONLY when the post-override delay is exactly 0 AND the message +// carries media — a non-zero agent override (even below the floor) is honored +// verbatim. Mirrors Phase 1's resolveInboundDebounceDelay + applyMediaFloor. +func chatDebounceDelay(cfg *config.Config, agentOtherConfig json.RawMessage, hasMedia bool) time.Duration { debounceMs := 0 if cfg != nil { debounceMs = cfg.Gateway.InboundDebounceMs @@ -137,6 +172,9 @@ func chatDebounceDelay(cfg *config.Config, agentOtherConfig json.RawMessage) tim if overrideMs, ok := store.ParseInboundDebounceMsFromOtherConfig(agentOtherConfig); ok { debounceMs = overrideMs } + if debounceMs <= 0 && hasMedia { + debounceMs = chatMediaDebounceFloorMs + } if debounceMs <= 0 { return 0 } diff --git a/internal/gateway/methods/chat_debounce_media_test.go b/internal/gateway/methods/chat_debounce_media_test.go new file mode 100644 index 00000000..d2ae0be1 --- /dev/null +++ b/internal/gateway/methods/chat_debounce_media_test.go @@ -0,0 +1,92 @@ +package methods + +import ( + "testing" + "time" + + "github.com/nextlevelbuilder/goclaw/internal/config" +) + +// TestChatDebounceDelay_HasMediaWithZeroConfigAppliesFloor — Phase 1.5 Rule #2. +// When global debounce is disabled, no agent override, and the message carries +// media, the 1000ms media floor MUST be applied. +func TestChatDebounceDelay_HasMediaWithZeroConfigAppliesFloor(t *testing.T) { + got := chatDebounceDelay(&config.Config{}, nil, true) + want := time.Duration(chatMediaDebounceFloorMs) * time.Millisecond + if got != want { + t.Fatalf("chatDebounceDelay(cfg=0, hasMedia=true) = %s, want %s", got, want) + } +} + +// TestChatDebounceDelay_AgentOverrideBelowFloorHonored — Rule #2 precedence. +// Floor fires only when post-override delay == 0. A 500ms override MUST be honored. +func TestChatDebounceDelay_AgentOverrideBelowFloorHonored(t *testing.T) { + cfg := &config.Config{} + cfg.Gateway.InboundDebounceMs = 0 + got := chatDebounceDelay(cfg, []byte(`{"inbound_debounce_ms":500}`), true) + if got != 500*time.Millisecond { + t.Fatalf("override 500ms with media = %s, want 500ms (floor must not raise)", got) + } +} + +// TestChatDebounceDelay_NoMediaZeroConfig: floor does NOT apply when no media. +func TestChatDebounceDelay_NoMediaZeroConfig(t *testing.T) { + got := chatDebounceDelay(&config.Config{}, nil, false) + if got != 0 { + t.Fatalf("chatDebounceDelay(cfg=0, hasMedia=false) = %s, want 0", got) + } +} + +// TestChatDebounceDelay_MediaConfigAboveFloorUnchanged: cfg already above floor → unchanged. +func TestChatDebounceDelay_MediaConfigAboveFloorUnchanged(t *testing.T) { + cfg := &config.Config{} + cfg.Gateway.InboundDebounceMs = 2000 + got := chatDebounceDelay(cfg, nil, true) + if got != 2000*time.Millisecond { + t.Fatalf("chatDebounceDelay(cfg=2000, hasMedia=true) = %s, want 2s", got) + } +} + +// TestChatDebouncer_NoMediaFollowupMergesIntoBufferedMedia — Rule #4. +// When a follow-up Push arrives with delay==0 while a buffer exists for the key, +// it MUST merge into the buffer rather than dispatch immediately. +func TestChatDebouncer_NoMediaFollowupMergesIntoBufferedMedia(t *testing.T) { + out := make(chan []chatSendRequest, 2) + d := newChatDebouncer(func(items []chatSendRequest) { + out <- items + }) + defer d.Stop() + + // First push: media-bearing, 50ms window. + d.Push("u1:s1", 50*time.Millisecond, chatSendRequest{params: chatSendParams{Message: "caption"}}) + // Second push: arrives while buffered, delay==0 (no media follow-up). + time.Sleep(10 * time.Millisecond) + d.Push("u1:s1", 0, chatSendRequest{params: chatSendParams{Message: "ps"}}) + + items := waitChatDebounce(t, out) + if len(items) != 2 { + t.Fatalf("flushed items = %d, want 2 (follow-up must merge, not bypass)", len(items)) + } + merged := mergeChatSendRequests(items).Message + if merged != "caption\nps" { + t.Fatalf("merged = %q, want %q", merged, "caption\nps") + } + + assertNoChatDebounceFlush(t, out) +} + +// TestChatDebouncer_DelayZeroNoBufferStillDispatches: delay==0 with empty buffer +// dispatches immediately (preserves existing behavior for plain text sends). +func TestChatDebouncer_DelayZeroNoBufferStillDispatches(t *testing.T) { + out := make(chan []chatSendRequest, 1) + d := newChatDebouncer(func(items []chatSendRequest) { + out <- items + }) + defer d.Stop() + + d.Push("u1:s1", 0, chatSendRequest{params: chatSendParams{Message: "one"}}) + items := waitChatDebounce(t, out) + if len(items) != 1 || items[0].params.Message != "one" { + t.Fatalf("dispatch = %#v, want single 'one'", items) + } +} diff --git a/internal/gateway/methods/chat_debounce_test.go b/internal/gateway/methods/chat_debounce_test.go index 47226f10..dd27f076 100644 --- a/internal/gateway/methods/chat_debounce_test.go +++ b/internal/gateway/methods/chat_debounce_test.go @@ -72,21 +72,22 @@ func TestChatDebouncerDiscardDropsPendingBeforeCancel(t *testing.T) { } func TestChatDebounceDelayGlobalAndAgentOverride(t *testing.T) { - if got := chatDebounceDelay(&config.Config{}, nil); got != 0 { + // hasMedia=false: legacy behavior preserved (no floor applied). + if got := chatDebounceDelay(&config.Config{}, nil, false); got != 0 { t.Fatalf("default debounce = %s, want disabled", got) } cfg := &config.Config{} cfg.Gateway.InboundDebounceMs = 250 - if got := chatDebounceDelay(cfg, nil); got != 250*time.Millisecond { + if got := chatDebounceDelay(cfg, nil, false); got != 250*time.Millisecond { t.Fatalf("global debounce = %s, want 250ms", got) } - if got := chatDebounceDelay(cfg, []byte(`{"inbound_debounce_ms":0}`)); got != 0 { + if got := chatDebounceDelay(cfg, []byte(`{"inbound_debounce_ms":0}`), false); got != 0 { t.Fatalf("agent disabled debounce = %s, want disabled", got) } - if got := chatDebounceDelay(cfg, []byte(`{"inbound_debounce_ms":500}`)); got != 500*time.Millisecond { + if got := chatDebounceDelay(cfg, []byte(`{"inbound_debounce_ms":500}`), false); got != 500*time.Millisecond { t.Fatalf("agent custom debounce = %s, want 500ms", got) } - if got := chatDebounceDelay(cfg, []byte(`{"other":true}`)); got != 250*time.Millisecond { + if got := chatDebounceDelay(cfg, []byte(`{"other":true}`), false); got != 250*time.Millisecond { t.Fatalf("agent inherit debounce = %s, want 250ms", got) } }