mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-11 03:13:24 +00:00
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
This commit is contained in:
1 parent
a591473546
commit
f771cff77c
20 files changed
+1408
-96
No files matched your search
@@ -46,6 +46,37 @@ All notable changes to GoClaw are documented here. For full documentation, see [
|
|||||||
|
|
||||||
### Fixed
|
### 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,
|
- **Upstream critical security remediation** — hardens gateway no-token fallback,
|
||||||
Feishu/Lark and Pancake webhooks, sandbox path/write handling, tenant-admin
|
Feishu/Lark and Pancake webhooks, sandbox path/write handling, tenant-admin
|
||||||
checks for mutable HTTP surfaces, and Lite hook schema migration verification.
|
checks for mutable HTTP surfaces, and Lite hook schema migration verification.
|
||||||
|
|||||||
@@ -72,6 +72,48 @@ Shared predicate: `store.IsMasterScope(ctx)` (`internal/store/context.go`).
|
|||||||
- `requireAuth(RoleAdmin)` as the **sole** gate on a global-state write
|
- `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)
|
- 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
|
## Test Layers
|
||||||
|
|
||||||
Tests are organized by priority and purpose:
|
Tests are organized by priority and purpose:
|
||||||
|
|||||||
@@ -3,7 +3,6 @@ package cmd
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
@@ -89,6 +88,10 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents
|
|||||||
return resolveInboundDebounceDelay(ctx, msg, deps)
|
return resolveInboundDebounceDelay(ctx, msg, deps)
|
||||||
},
|
},
|
||||||
func(msg bus.InboundMessage) {
|
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)
|
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) ---
|
// --- Dedup: skip duplicate inbound messages (matching TS shouldSkipDuplicateInbound) ---
|
||||||
if msgID := msg.Metadata["message_id"]; msgID != "" {
|
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) {
|
if dedupe.IsDuplicate(dedupeKey) {
|
||||||
slog.Debug("dedup: skipping duplicate message", "key", dedupeKey)
|
slog.Debug("dedup: skipping duplicate message", "key", dedupeKey)
|
||||||
continue
|
continue
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package cmd
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
@@ -11,6 +12,12 @@ import (
|
|||||||
"github.com/nextlevelbuilder/goclaw/internal/store"
|
"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) {
|
func prepareInboundDebounceMessage(msg *bus.InboundMessage, deps *ConsumerDeps) {
|
||||||
if msg == nil || deps == nil || deps.Cfg == nil || msg.AgentID != "" {
|
if msg == nil || deps == nil || deps.Cfg == nil || msg.AgentID != "" {
|
||||||
return
|
return
|
||||||
@@ -24,7 +31,7 @@ func resolveInboundDebounceDelay(ctx context.Context, msg bus.InboundMessage, de
|
|||||||
debounceMs = deps.Cfg.Gateway.InboundDebounceMs
|
debounceMs = deps.Cfg.Gateway.InboundDebounceMs
|
||||||
}
|
}
|
||||||
if deps == nil || deps.AgentStore == nil || msg.AgentID == "" {
|
if deps == nil || deps.AgentStore == nil || msg.AgentID == "" {
|
||||||
return inboundDebounceDuration(debounceMs)
|
return inboundDebounceDuration(applyMediaFloor(debounceMs, msg))
|
||||||
}
|
}
|
||||||
|
|
||||||
agentCtx := ctx
|
agentCtx := ctx
|
||||||
@@ -39,12 +46,35 @@ func resolveInboundDebounceDelay(ctx context.Context, msg bus.InboundMessage, de
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
slog.Debug("inbound debounce: agent config unavailable", "agent", msg.AgentID, "error", err)
|
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 {
|
if overrideMs, ok := agentData.ParseInboundDebounceMs(); ok {
|
||||||
debounceMs = overrideMs
|
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) {
|
func getInboundDebounceAgent(ctx context.Context, agentStore store.AgentStore, agentID string) (*store.AgentData, error) {
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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))
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -69,7 +69,17 @@ The consumer routes system messages based on sender ID prefixes:
|
|||||||
|
|
||||||
### Inbound Debounce
|
### 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.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
|||||||
@@ -48,21 +48,33 @@ func NewInboundDebouncerFunc(delayFn func(InboundMessage) time.Duration, flushFn
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Push adds a message to the debounce buffer.
|
// 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) {
|
func (d *InboundDebouncer) Push(msg InboundMessage) {
|
||||||
debounceMs := d.delayFn(msg)
|
debounceMs := d.delayFn(msg)
|
||||||
|
|
||||||
// Disabled: pass through immediately.
|
|
||||||
if debounceMs <= 0 {
|
|
||||||
d.flushFn(msg)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
key := debounceKey(msg)
|
key := debounceKey(msg)
|
||||||
|
|
||||||
// Media messages bypass debounce — flush any buffered text first, then process media.
|
if debounceMs <= 0 {
|
||||||
if len(msg.Media) > 0 {
|
// Disabled-path: merge into existing buffer if any, else pass through.
|
||||||
d.flushKey(key)
|
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)
|
d.flushFn(msg)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -145,8 +157,15 @@ func debounceKey(msg InboundMessage) string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// mergeInboundMessages combines multiple messages into one.
|
// 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 {
|
func mergeInboundMessages(msgs []InboundMessage) InboundMessage {
|
||||||
if len(msgs) == 1 {
|
if len(msgs) == 1 {
|
||||||
return msgs[0]
|
return msgs[0]
|
||||||
@@ -170,9 +189,62 @@ func mergeInboundMessages(msgs []InboundMessage) InboundMessage {
|
|||||||
}
|
}
|
||||||
last.Media = allMedia
|
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
|
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.
|
// truncateStr truncates a string to maxLen characters.
|
||||||
func truncateStr(s string, maxLen int) string {
|
func truncateStr(s string, maxLen int) string {
|
||||||
if len(s) <= maxLen {
|
if len(s) <= maxLen {
|
||||||
|
|||||||
@@ -1,6 +1,8 @@
|
|||||||
package bus
|
package bus
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"sort"
|
||||||
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
@@ -83,28 +85,194 @@ func TestInboundDebouncerSeparatesAgents(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestInboundDebouncerMediaFlushesPendingTextFirst(t *testing.T) {
|
// TestInboundDebouncerMergesMediaWithinWindow asserts two media messages within
|
||||||
out := make(chan InboundMessage, 2)
|
// the debounce window merge into a single flush with all media in arrival order.
|
||||||
d := NewInboundDebouncer(time.Minute, func(msg InboundMessage) {
|
// 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
|
out <- msg
|
||||||
})
|
})
|
||||||
defer d.Stop()
|
defer d.Stop()
|
||||||
|
|
||||||
d.Push(InboundMessage{Channel: "telegram", ChatID: "chat-1", SenderID: "user-1", Content: "pending"})
|
|
||||||
d.Push(InboundMessage{
|
d.Push(InboundMessage{
|
||||||
Channel: "telegram",
|
Channel: "telegram", ChatID: "chat-1", SenderID: "user-1",
|
||||||
ChatID: "chat-1",
|
Media: []MediaFile{{Path: "/tmp/a.png", MimeType: "image/png"}},
|
||||||
SenderID: "user-1",
|
})
|
||||||
Content: "with media",
|
d.Push(InboundMessage{
|
||||||
Media: []MediaFile{{Path: "/tmp/a.png", MimeType: "image/png"}},
|
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 {
|
got := waitInbound(t, out)
|
||||||
t.Fatalf("first flush = %#v, want pending text without media", got)
|
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 {
|
if got.Media[0].Path != "/tmp/a.png" || got.Media[1].Path != "/tmp/b.png" {
|
||||||
t.Fatalf("second flush = %#v, want media message", got)
|
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 {
|
func waitInbound(t *testing.T, ch <-chan InboundMessage) InboundMessage {
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -41,12 +41,14 @@ type Channel struct {
|
|||||||
threadIDs sync.Map // localKey string → messageThreadID int (for forum topic routing)
|
threadIDs sync.Map // localKey string → messageThreadID int (for forum topic routing)
|
||||||
mentionMode string // "strict" (default) or "yield"
|
mentionMode string // "strict" (default) or "yield"
|
||||||
botDisplayName string // bot's first_name from GetMe (e.g. "ViệtBot"); captured once at Start
|
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
|
pollCancel context.CancelFunc // cancels the long polling context
|
||||||
pollDone chan struct{} // closed when polling goroutine exits
|
pollDone chan struct{} // closed when polling goroutine exits
|
||||||
handlerWg sync.WaitGroup // tracks in-flight handler goroutines for graceful shutdown
|
handlerWg sync.WaitGroup // tracks in-flight handler goroutines for graceful shutdown
|
||||||
handlerSem chan struct{} // bounded semaphore for concurrent handler goroutines
|
handlerSem chan struct{} // bounded semaphore for concurrent handler goroutines
|
||||||
pendingDraftID sync.Map // localKey string → int (draftID)
|
pendingDraftID sync.Map // localKey string → int (draftID)
|
||||||
audioMgr *audio.Manager // unified STT via audio.Manager (nil = no STT)
|
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
|
writerHealMu sync.Mutex // guards writerHealLastTry for /writers self-heal
|
||||||
writerHealLastTry map[string]time.Time // key "chatID|userID" → last attempt timestamp
|
writerHealLastTry map[string]time.Time // key "chatID|userID" → last attempt timestamp
|
||||||
// pairingService, approvedGroups, pairingDebounce, groupHistory, historyLimit, requireMention
|
// 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.
|
// Create a cancellable context for the polling goroutine.
|
||||||
// Stop() cancels this context to cleanly shut down long polling.
|
// Stop() cancels this context to cleanly shut down long polling.
|
||||||
pollCtx, cancel := context.WithCancel(ctx)
|
c.pollCtx, c.pollCancel = context.WithCancel(ctx)
|
||||||
c.pollCancel = cancel
|
pollCtx := c.pollCtx
|
||||||
|
cancel := c.pollCancel
|
||||||
c.pollDone = make(chan struct{})
|
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{
|
updates, err := c.bot.UpdatesViaLongPolling(pollCtx, &telego.GetUpdatesParams{
|
||||||
Timeout: 25, // Long-poll seconds; keep below HTTP client Timeout (#361)
|
Timeout: 25, // Long-poll seconds; keep below HTTP client Timeout (#361)
|
||||||
AllowedUpdates: []string{
|
AllowedUpdates: []string{
|
||||||
@@ -379,6 +402,12 @@ func (c *Channel) Stop(_ context.Context) error {
|
|||||||
gh.StopFlusher()
|
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 {
|
if c.pollCancel != nil {
|
||||||
c.pollCancel()
|
c.pollCancel()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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) ---
|
// --- Media download (only when bot will process the message) ---
|
||||||
// Deferred until after mention + pairing gates to avoid downloading
|
// Deferred until after mention + pairing gates to avoid downloading
|
||||||
// media for messages that only get recorded in pending history.
|
// media for messages that only get recorded in pending history.
|
||||||
mediaList, mediaErrors := c.resolveMedia(ctx, message)
|
var mediaList []MediaInfo
|
||||||
if message.ReplyToMessage != nil {
|
var mediaErrors []MediaError
|
||||||
replyMedia, replyErrors := c.resolveMedia(ctx, message.ReplyToMessage)
|
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 {
|
if len(replyMedia) > 0 {
|
||||||
// Tag reply media so LLM knows which images came from the replied-to message.
|
// Tag reply media so LLM knows which images came from the replied-to message.
|
||||||
for i := range replyMedia {
|
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.
|
// Reply media first (context), current media second.
|
||||||
mediaList = append(replyMedia, mediaList...)
|
mediaList = append(replyMedia, mediaList...)
|
||||||
slog.Debug("telegram: resolved media from replied message",
|
slog.Debug("telegram: resolved media from replied message",
|
||||||
"reply_msg_id", message.ReplyToMessage.MessageID,
|
"reply_msg_id", rep.ReplyToMessage.MessageID,
|
||||||
"media_count", len(replyMedia),
|
"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).
|
// Replace lightweight media tags with full tags (includes transcripts).
|
||||||
fullTags := buildMediaTags(mediaList)
|
fullTags := buildMediaTags(mediaList)
|
||||||
lightTags := lightweightMediaTags(message)
|
lightTags := lightweightMediaTags(rep)
|
||||||
if lightTags != "" && fullTags != "" {
|
if lightTags != "" && fullTags != "" {
|
||||||
content = strings.Replace(content, lightTags, fullTags, 1)
|
content = strings.Replace(content, lightTags, fullTags, 1)
|
||||||
} else if fullTags != "" {
|
} else if fullTags != "" {
|
||||||
@@ -479,7 +535,7 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) {
|
|||||||
if len(mediaErrors) > 0 {
|
if len(mediaErrors) > 0 {
|
||||||
for _, me := range mediaErrors {
|
for _, me := range mediaErrors {
|
||||||
errTag := fmt.Sprintf("[sent media (%s) — skipped: %s]", me.Type, me.Reason)
|
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)
|
content = strings.Replace(content, lightTag, errTag, 1)
|
||||||
} else {
|
} else {
|
||||||
content = errTag + "\n" + content
|
content = errTag + "\n" + content
|
||||||
@@ -494,24 +550,24 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) {
|
|||||||
} else {
|
} else {
|
||||||
errText = "⚠️ Failed to download the attached file. Skipped."
|
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",
|
slog.Debug("telegram message received",
|
||||||
"sender_id", senderID,
|
"sender_id", rctx.senderID,
|
||||||
"chat_id", fmt.Sprintf("%d", chatID),
|
"chat_id", rctx.chatIDStr,
|
||||||
"preview", channels.Truncate(content, 50),
|
"preview", channels.Truncate(content, 50),
|
||||||
)
|
)
|
||||||
|
|
||||||
// Build context from pending group history (if any).
|
// Build context from pending group history (if any).
|
||||||
// Annotate current message with sender name so LLM knows who is talking.
|
// Annotate current message with sender name so LLM knows who is talking.
|
||||||
finalContent := content
|
finalContent := content
|
||||||
if isGroup {
|
if rctx.isGroup {
|
||||||
annotated := fmt.Sprintf("[From: %s]\n%s", senderLabel, content)
|
annotated := fmt.Sprintf("[From: %s]\n%s", rctx.senderLabel, content)
|
||||||
if c.HistoryLimit() > 0 {
|
if c.HistoryLimit() > 0 {
|
||||||
// Resolve deferred media from history entries (lazy download).
|
// 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)
|
histMedia, histErrors := c.resolveMediaRefs(ctx, histRefs)
|
||||||
if len(histMedia) > 0 {
|
if len(histMedia) > 0 {
|
||||||
mediaFiles = prependMediaInfoFiles(mediaFiles, histMedia)
|
mediaFiles = prependMediaInfoFiles(mediaFiles, histMedia)
|
||||||
@@ -523,39 +579,39 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) {
|
|||||||
"type", e.Type, "reason", e.Reason)
|
"type", e.Type, "reason", e.Reason)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
finalContent = c.GroupHistory().BuildContext(localKey, annotated, c.HistoryLimit())
|
finalContent = c.GroupHistory().BuildContext(rctx.localKey, annotated, c.HistoryLimit())
|
||||||
} else {
|
} else {
|
||||||
finalContent = annotated
|
finalContent = annotated
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
// DM: annotate with sender identity so the agent knows who is messaging.
|
// 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.
|
// Send typing indicator with keepalive + TTL safety net.
|
||||||
// Telegram typing expires after 5s, so keepalive every 4s.
|
// Telegram typing expires after 5s, so keepalive every 4s.
|
||||||
// TTL auto-stops after 60s to prevent stuck indicators.
|
// TTL auto-stops after 60s to prevent stuck indicators.
|
||||||
chatIDObj := tu.ID(chatID)
|
chatIDObj := tu.ID(rctx.chatID)
|
||||||
typingCtrl := typing.New(typing.Options{
|
typingCtrl := typing.New(typing.Options{
|
||||||
MaxDuration: 60 * time.Second,
|
MaxDuration: 60 * time.Second,
|
||||||
KeepaliveInterval: 4 * time.Second,
|
KeepaliveInterval: 4 * time.Second,
|
||||||
StartFn: func() error {
|
StartFn: func() error {
|
||||||
action := tu.ChatAction(chatIDObj, telego.ChatActionTyping)
|
action := tu.ChatAction(chatIDObj, telego.ChatActionTyping)
|
||||||
if messageThreadID > 0 {
|
if rctx.messageThreadID > 0 {
|
||||||
action.MessageThreadID = messageThreadID
|
action.MessageThreadID = rctx.messageThreadID
|
||||||
}
|
}
|
||||||
return c.bot.SendChatAction(ctx, action)
|
return c.bot.SendChatAction(ctx, action)
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
// Stop previous typing controller for this chat/topic (if any)
|
// 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()
|
prev.(*typing.Controller).Stop()
|
||||||
}
|
}
|
||||||
c.typingCtrls.Store(localKey, typingCtrl)
|
c.typingCtrls.Store(rctx.localKey, typingCtrl)
|
||||||
typingCtrl.Start()
|
typingCtrl.Start()
|
||||||
|
|
||||||
// Stop previous thinking animation for this chat/topic
|
// 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 {
|
if cf, ok := prevStop.(*thinkingCancel); ok {
|
||||||
cf.Cancel()
|
cf.Cancel()
|
||||||
}
|
}
|
||||||
@@ -563,7 +619,7 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) {
|
|||||||
|
|
||||||
// Create thinking cancel for this chat/topic
|
// Create thinking cancel for this chat/topic
|
||||||
_, thinkCancel := context.WithCancel(ctx)
|
_, 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
|
// No "Thinking..." placeholder — the DraftStream creates its own message
|
||||||
// on the first streaming chunk (sendMessage on first flush).
|
// 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.
|
// user sees typing indicator → first content appears directly.
|
||||||
|
|
||||||
metadata := map[string]string{
|
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),
|
"user_id": fmt.Sprintf("%d", user.ID),
|
||||||
tools.MetaUsername: user.Username,
|
tools.MetaUsername: user.Username,
|
||||||
"first_name": user.FirstName,
|
"first_name": user.FirstName,
|
||||||
"is_group": fmt.Sprintf("%t", isGroup),
|
"is_group": fmt.Sprintf("%t", rctx.isGroup),
|
||||||
"local_key": localKey,
|
"local_key": rctx.localKey,
|
||||||
}
|
}
|
||||||
if message.Chat.Title != "" {
|
// When this publish coalesces multiple platform messages (album members),
|
||||||
metadata[tools.MetaChatTitle] = message.Chat.Title
|
// 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.MetaIsForum] = "true"
|
||||||
metadata[tools.MetaMessageThreadID] = fmt.Sprintf("%d", messageThreadID)
|
metadata[tools.MetaMessageThreadID] = fmt.Sprintf("%d", rctx.messageThreadID)
|
||||||
}
|
}
|
||||||
if dmThreadID > 0 {
|
if rctx.dmThreadID > 0 {
|
||||||
metadata[tools.MetaDMThreadID] = fmt.Sprintf("%d", dmThreadID)
|
metadata[tools.MetaDMThreadID] = fmt.Sprintf("%d", rctx.dmThreadID)
|
||||||
metadata[tools.MetaMessageThreadID] = fmt.Sprintf("%d", dmThreadID)
|
metadata[tools.MetaMessageThreadID] = fmt.Sprintf("%d", rctx.dmThreadID)
|
||||||
}
|
}
|
||||||
// Self-identity hint so the LLM knows its own Telegram handle and does not
|
// Self-identity hint so the LLM knows its own Telegram handle and does not
|
||||||
// confuse other bots' @mentions (preserved after stripBotMention) for its own.
|
// 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
|
metadata[tools.MetaChannelSelfIdentity] = identity
|
||||||
}
|
}
|
||||||
|
|
||||||
if topicCfg.systemPrompt != "" {
|
if rctx.topicCfg.systemPrompt != "" {
|
||||||
metadata[tools.MetaTopicSystemPrompt] = topicCfg.systemPrompt
|
metadata[tools.MetaTopicSystemPrompt] = rctx.topicCfg.systemPrompt
|
||||||
}
|
}
|
||||||
if topicCfg.skills != nil {
|
if rctx.topicCfg.skills != nil {
|
||||||
metadata[tools.MetaTopicSkills] = strings.Join(topicCfg.skills, ",")
|
metadata[tools.MetaTopicSkills] = strings.Join(rctx.topicCfg.skills, ",")
|
||||||
}
|
}
|
||||||
|
|
||||||
peerKind := "direct"
|
peerKind := "direct"
|
||||||
if isGroup {
|
if rctx.isGroup {
|
||||||
peerKind = "group"
|
peerKind = "group"
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -626,35 +694,35 @@ func (c *Channel) handleMessage(ctx context.Context, update telego.Update) {
|
|||||||
// Collect contact for processed messages (DM + group-mentioned).
|
// Collect contact for processed messages (DM + group-mentioned).
|
||||||
if cc := c.ContactCollector(); cc != nil {
|
if cc := c.ContactCollector(); cc != nil {
|
||||||
contactName := strings.TrimSpace(user.FirstName + " " + user.LastName)
|
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).
|
// Also collect group chat itself as a contact (for group permission / merge).
|
||||||
if isGroup {
|
if rctx.isGroup {
|
||||||
cc.EnsureContact(ctx, c.Type(), c.Name(), chatIDStr, "", message.Chat.Title, "", "group", "group", "", "")
|
cc.EnsureContact(ctx, c.Type(), c.Name(), rctx.chatIDStr, "", rep.Chat.Title, "", "group", "group", "", "")
|
||||||
// Collect forum topic as a distinct delivery target (including General).
|
// Collect forum topic as a distinct delivery target (including General).
|
||||||
if isForum && messageThreadID > 0 {
|
if rctx.isForum && rctx.messageThreadID > 0 {
|
||||||
threadStr := fmt.Sprintf("%d", messageThreadID)
|
threadStr := fmt.Sprintf("%d", rctx.messageThreadID)
|
||||||
cc.EnsureContact(ctx, c.Type(), c.Name(), chatIDStr, "", message.Chat.Title, "", "group", "topic", threadStr, "topic")
|
cc.EnsureContact(ctx, c.Type(), c.Name(), rctx.chatIDStr, "", rep.Chat.Title, "", "group", "topic", threadStr, "topic")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
c.Bus().PublishInbound(bus.InboundMessage{
|
c.Bus().PublishInbound(bus.InboundMessage{
|
||||||
Channel: c.Name(),
|
Channel: c.Name(),
|
||||||
SenderID: senderID,
|
SenderID: rctx.senderID,
|
||||||
ChatID: chatIDStr,
|
ChatID: rctx.chatIDStr,
|
||||||
Content: finalContent,
|
Content: finalContent,
|
||||||
Media: mediaFiles,
|
Media: mediaFiles,
|
||||||
PeerKind: peerKind,
|
PeerKind: peerKind,
|
||||||
UserID: userID,
|
UserID: rctx.userID,
|
||||||
AgentID: targetAgentID,
|
AgentID: targetAgentID,
|
||||||
HistoryLimit: c.HistoryLimit(),
|
HistoryLimit: c.HistoryLimit(),
|
||||||
ToolAllow: topicCfg.tools,
|
ToolAllow: rctx.topicCfg.tools,
|
||||||
TenantID: c.TenantID(),
|
TenantID: c.TenantID(),
|
||||||
Metadata: metadata,
|
Metadata: metadata,
|
||||||
})
|
})
|
||||||
|
|
||||||
// Clear pending history after sending to agent.
|
// Clear pending history after sending to agent.
|
||||||
if isGroup {
|
if rctx.isGroup {
|
||||||
c.GroupHistory().Clear(localKey)
|
c.GroupHistory().Clear(rctx.localKey)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -370,7 +370,7 @@ type GatewayConfig struct {
|
|||||||
MaxMessageChars int `json:"max_message_chars,omitempty"` // max user message characters (default 32000)
|
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)
|
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"
|
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
|
Quota *QuotaConfig `json:"quota,omitempty"` // per-user/group request quotas
|
||||||
BlockReply *bool `json:"block_reply,omitempty"` // deliver intermediate text during tool iterations (default false)
|
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)
|
ToolStatus *bool `json:"tool_status,omitempty"` // show tool name in streaming preview during tool execution (default true)
|
||||||
|
|||||||
@@ -218,16 +218,18 @@ func (m *ChatMethods) handleSend(ctx context.Context, client *gateway.Client, re
|
|||||||
m.abortChatSession(req.ID, client, sessionKey)
|
m.abortChatSession(req.ID, client, sessionKey)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if len(params.parseMedia()) > 0 {
|
// Media-bearing sends route through the same debouncer path as text.
|
||||||
pending := m.debouncer.Take(debounceKey)
|
// The media floor in chatDebounceDelay guarantees a non-zero window when
|
||||||
m.dispatchChatSends(append(pending, item))
|
// the operator has disabled debouncing, so multi-attachment bursts coalesce
|
||||||
return
|
// into a single dispatch (issue #63).
|
||||||
}
|
hasMedia := len(params.parseMedia()) > 0
|
||||||
if delay := chatDebounceDelay(m.cfg, loop.OtherConfig()); delay > 0 {
|
delay := chatDebounceDelay(m.cfg, loop.OtherConfig(), hasMedia)
|
||||||
|
if delay > 0 {
|
||||||
m.debouncer.Push(debounceKey, delay, item)
|
m.debouncer.Push(debounceKey, delay, item)
|
||||||
return
|
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) {
|
func (m *ChatMethods) abortChatSession(reqID string, client *gateway.Client, sessionKey string) {
|
||||||
|
|||||||
@@ -13,6 +13,12 @@ import (
|
|||||||
"github.com/nextlevelbuilder/goclaw/internal/store"
|
"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 {
|
type chatSendRequest struct {
|
||||||
ctx context.Context
|
ctx context.Context
|
||||||
client *gateway.Client
|
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) {
|
func (d *chatDebouncer) Push(key string, delay time.Duration, item chatSendRequest) {
|
||||||
if delay <= 0 {
|
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})
|
d.flushFn([]chatSendRequest{item})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -129,7 +158,13 @@ func mergeChatSendRequests(items []chatSendRequest) chatSendParams {
|
|||||||
return last
|
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
|
debounceMs := 0
|
||||||
if cfg != nil {
|
if cfg != nil {
|
||||||
debounceMs = cfg.Gateway.InboundDebounceMs
|
debounceMs = cfg.Gateway.InboundDebounceMs
|
||||||
@@ -137,6 +172,9 @@ func chatDebounceDelay(cfg *config.Config, agentOtherConfig json.RawMessage) tim
|
|||||||
if overrideMs, ok := store.ParseInboundDebounceMsFromOtherConfig(agentOtherConfig); ok {
|
if overrideMs, ok := store.ParseInboundDebounceMsFromOtherConfig(agentOtherConfig); ok {
|
||||||
debounceMs = overrideMs
|
debounceMs = overrideMs
|
||||||
}
|
}
|
||||||
|
if debounceMs <= 0 && hasMedia {
|
||||||
|
debounceMs = chatMediaDebounceFloorMs
|
||||||
|
}
|
||||||
if debounceMs <= 0 {
|
if debounceMs <= 0 {
|
||||||
return 0
|
return 0
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -72,21 +72,22 @@ func TestChatDebouncerDiscardDropsPendingBeforeCancel(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func TestChatDebounceDelayGlobalAndAgentOverride(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)
|
t.Fatalf("default debounce = %s, want disabled", got)
|
||||||
}
|
}
|
||||||
cfg := &config.Config{}
|
cfg := &config.Config{}
|
||||||
cfg.Gateway.InboundDebounceMs = 250
|
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)
|
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)
|
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)
|
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)
|
t.Fatalf("agent inherit debounce = %s, want 250ms", got)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in new issue
Block a user