mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-11 12:18:59 +00:00
feat(channels): pending message compaction — fix provider, wire auto-compact, add global config & UI
- Fix compact endpoint using random provider instead of agent's configured provider+model - Wire auto-compaction for all 5 channel types (telegram, discord, slack, feishu, zalo_personal) via PendingCompactable interface and InstanceLoader - Add global PendingCompactionConfig (threshold, keep_recent) to ChannelsConfig - Wire global config through InstanceLoader and PendingMessagesHandler - Increase compaction timeout from 45s to 180s for slow providers - Add pending compaction config card to Behavior tab in config page - Add HowItWorksCard (expanded by default) and toast notifications to pending messages page - Add i18n support for all new strings (en/vi/zh)
This commit is contained in:
1 parent
e2015835b4
commit
3b6bf645f3
22 files changed
+400
-41
No files matched your search
@@ -616,6 +616,9 @@ func runGateway() {
|
||||
server.SetBuiltinToolsHandler(builtinToolsH)
|
||||
}
|
||||
if pendingMessagesH != nil {
|
||||
if pc := cfg.Channels.PendingCompaction; pc != nil {
|
||||
pendingMessagesH.SetKeepRecent(pc.KeepRecent)
|
||||
}
|
||||
server.SetPendingMessagesHandler(pendingMessagesH)
|
||||
}
|
||||
|
||||
@@ -678,6 +681,8 @@ func runGateway() {
|
||||
var instanceLoader *channels.InstanceLoader
|
||||
if pgStores.ChannelInstances != nil {
|
||||
instanceLoader = channels.NewInstanceLoader(pgStores.ChannelInstances, pgStores.Agents, channelMgr, msgBus, pgStores.Pairing)
|
||||
instanceLoader.SetProviderRegistry(providerRegistry)
|
||||
instanceLoader.SetPendingCompactionConfig(cfg.Channels.PendingCompaction)
|
||||
instanceLoader.RegisterFactory(channels.TypeTelegram, telegram.FactoryWithStores(pgStores.Agents, pgStores.Teams, pgStores.PendingMessages))
|
||||
instanceLoader.RegisterFactory(channels.TypeDiscord, discord.FactoryWithPendingStore(pgStores.PendingMessages))
|
||||
instanceLoader.RegisterFactory(channels.TypeFeishu, feishu.FactoryWithPendingStore(pgStores.PendingMessages))
|
||||
|
||||
@@ -71,7 +71,7 @@ func wireHTTP(stores *store.Stores, token string, msgBus *bus.MessageBus, toolsR
|
||||
}
|
||||
|
||||
if stores != nil && stores.PendingMessages != nil {
|
||||
pendingMessagesH = httpapi.NewPendingMessagesHandler(stores.PendingMessages, token, providerReg)
|
||||
pendingMessagesH = httpapi.NewPendingMessagesHandler(stores.PendingMessages, stores.Agents, token, providerReg)
|
||||
}
|
||||
|
||||
return agentsH, skillsH, tracesH, mcpH, customToolsH, channelInstancesH, providersH, delegationsH, builtinToolsH, pendingMessagesH
|
||||
|
||||
@@ -308,6 +308,13 @@ func (c *BaseChannel) HandleMessage(senderID, chatID, content string, media []st
|
||||
c.bus.PublishInbound(msg)
|
||||
}
|
||||
|
||||
// PendingCompactable is optionally implemented by channels that have a PendingHistory
|
||||
// supporting LLM-based compaction. InstanceLoader uses this to wire compaction config
|
||||
// after channel creation.
|
||||
type PendingCompactable interface {
|
||||
SetPendingCompaction(cfg *CompactionConfig)
|
||||
}
|
||||
|
||||
// Truncate shortens a string to maxLen, appending "..." if truncated.
|
||||
func Truncate(s string, maxLen int) string {
|
||||
if len(s) <= maxLen {
|
||||
|
||||
@@ -98,6 +98,11 @@ func (c *Channel) Start(_ context.Context) error {
|
||||
// BlockReplyEnabled returns the per-channel block_reply override (nil = inherit gateway default).
|
||||
func (c *Channel) BlockReplyEnabled() *bool { return c.config.BlockReply }
|
||||
|
||||
// SetPendingCompaction configures LLM-based auto-compaction for pending messages.
|
||||
func (c *Channel) SetPendingCompaction(cfg *channels.CompactionConfig) {
|
||||
c.groupHistory.SetCompactionConfig(cfg)
|
||||
}
|
||||
|
||||
// Stop closes the Discord gateway connection.
|
||||
func (c *Channel) Stop(_ context.Context) error {
|
||||
c.groupHistory.StopFlusher()
|
||||
|
||||
@@ -121,6 +121,11 @@ func (c *Channel) Start(ctx context.Context) error {
|
||||
// BlockReplyEnabled returns the per-channel block_reply override (nil = inherit gateway default).
|
||||
func (c *Channel) BlockReplyEnabled() *bool { return c.cfg.BlockReply }
|
||||
|
||||
// SetPendingCompaction configures LLM-based auto-compaction for pending messages.
|
||||
func (c *Channel) SetPendingCompaction(cfg *channels.CompactionConfig) {
|
||||
c.groupHistory.SetCompactionConfig(cfg)
|
||||
}
|
||||
|
||||
// Stop shuts down the Feishu channel.
|
||||
func (c *Channel) Stop(_ context.Context) error {
|
||||
c.groupHistory.StopFlusher()
|
||||
|
||||
@@ -123,7 +123,7 @@ func (ph *PendingHistory) runCompaction(historyKey string, cfg *CompactionConfig
|
||||
// Force-flush buffer to ensure DB is consistent
|
||||
ph.flushNow()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 45*time.Second)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 180*time.Second)
|
||||
defer cancel()
|
||||
|
||||
// Check threshold from DB (may have been cleared since trigger)
|
||||
|
||||
@@ -10,6 +10,8 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/nextlevelbuilder/goclaw/internal/bus"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/config"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/providers"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/store"
|
||||
)
|
||||
|
||||
@@ -23,14 +25,16 @@ type ChannelFactory func(name string, creds json.RawMessage, cfg json.RawMessage
|
||||
// InstanceLoader loads channel instances from the database and registers them with the Manager.
|
||||
// Follows the DynamicToolLoader pattern: LoadAll at startup, Reload on cache invalidation.
|
||||
type InstanceLoader struct {
|
||||
store store.ChannelInstanceStore
|
||||
agentStore store.AgentStore
|
||||
factories map[string]ChannelFactory
|
||||
manager *Manager
|
||||
msgBus *bus.MessageBus
|
||||
pairingSvc store.PairingStore
|
||||
mu sync.Mutex
|
||||
loaded map[string]struct{} // channel names managed by this loader
|
||||
store store.ChannelInstanceStore
|
||||
agentStore store.AgentStore
|
||||
providerReg *providers.Registry
|
||||
pendingCompactCfg *config.PendingCompactionConfig
|
||||
factories map[string]ChannelFactory
|
||||
manager *Manager
|
||||
msgBus *bus.MessageBus
|
||||
pairingSvc store.PairingStore
|
||||
mu sync.Mutex
|
||||
loaded map[string]struct{} // channel names managed by this loader
|
||||
}
|
||||
|
||||
// NewInstanceLoader creates a new InstanceLoader.
|
||||
@@ -52,6 +56,18 @@ func NewInstanceLoader(
|
||||
}
|
||||
}
|
||||
|
||||
// SetProviderRegistry sets the provider registry for pending message compaction.
|
||||
// Must be called before LoadAll/Reload.
|
||||
func (l *InstanceLoader) SetProviderRegistry(reg *providers.Registry) {
|
||||
l.providerReg = reg
|
||||
}
|
||||
|
||||
// SetPendingCompactionConfig sets the global pending message compaction thresholds.
|
||||
// Must be called before LoadAll/Reload.
|
||||
func (l *InstanceLoader) SetPendingCompactionConfig(cfg *config.PendingCompactionConfig) {
|
||||
l.pendingCompactCfg = cfg
|
||||
}
|
||||
|
||||
// RegisterFactory registers a factory for a channel type (e.g., "telegram", "discord").
|
||||
func (l *InstanceLoader) RegisterFactory(channelType string, factory ChannelFactory) {
|
||||
l.factories[channelType] = factory
|
||||
@@ -205,8 +221,10 @@ func (l *InstanceLoader) loadInstance(ctx context.Context, inst store.ChannelIns
|
||||
}
|
||||
|
||||
// Resolve agent_key from UUID — the routing system (Router, session keys) uses agent_key, not UUID.
|
||||
var ag *store.AgentData
|
||||
if base, ok := ch.(interface{ SetAgentID(string) }); ok {
|
||||
ag, err := l.agentStore.GetByID(ctx, inst.AgentID)
|
||||
var err error
|
||||
ag, err = l.agentStore.GetByID(ctx, inst.AgentID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("agent %s not found for channel %s: %w", inst.AgentID, inst.Name, err)
|
||||
}
|
||||
@@ -216,6 +234,30 @@ func (l *InstanceLoader) loadInstance(ctx context.Context, inst store.ChannelIns
|
||||
if base, ok := ch.(interface{ SetType(string) }); ok {
|
||||
base.SetType(inst.ChannelType)
|
||||
}
|
||||
|
||||
// Wire pending message auto-compaction using the agent's configured provider+model.
|
||||
if pc, ok := ch.(PendingCompactable); ok && ag != nil && l.providerReg != nil && ag.Provider != "" {
|
||||
if p, err := l.providerReg.Get(ag.Provider); err == nil {
|
||||
model := ag.Model
|
||||
if model == "" {
|
||||
model = p.DefaultModel()
|
||||
}
|
||||
if model != "" {
|
||||
cc := &CompactionConfig{
|
||||
Provider: p,
|
||||
Model: model,
|
||||
}
|
||||
// Apply global threshold/keepRecent from config if set.
|
||||
if l.pendingCompactCfg != nil {
|
||||
cc.Threshold = l.pendingCompactCfg.Threshold
|
||||
cc.KeepRecent = l.pendingCompactCfg.KeepRecent
|
||||
}
|
||||
pc.SetPendingCompaction(cc)
|
||||
slog.Debug("pending compaction configured", "channel", inst.Name, "provider", ag.Provider, "model", model,
|
||||
"threshold", cc.Threshold, "keep_recent", cc.KeepRecent)
|
||||
}
|
||||
}
|
||||
}
|
||||
l.manager.RegisterChannel(inst.Name, ch)
|
||||
l.loaded[inst.Name] = struct{}{}
|
||||
|
||||
|
||||
@@ -280,6 +280,11 @@ func (c *Channel) handleEvent(evt socketmode.Event) {
|
||||
}
|
||||
}
|
||||
|
||||
// SetPendingCompaction configures LLM-based auto-compaction for pending messages.
|
||||
func (c *Channel) SetPendingCompaction(cfg *channels.CompactionConfig) {
|
||||
c.groupHistory.SetCompactionConfig(cfg)
|
||||
}
|
||||
|
||||
// Stop gracefully shuts down the Slack channel.
|
||||
func (c *Channel) Stop(_ context.Context) error {
|
||||
c.groupHistory.StopFlusher()
|
||||
|
||||
@@ -206,6 +206,11 @@ func (c *Channel) StreamEnabled(isGroup bool) bool {
|
||||
// BlockReplyEnabled returns the per-channel block_reply override (nil = inherit gateway default).
|
||||
func (c *Channel) BlockReplyEnabled() *bool { return c.config.BlockReply }
|
||||
|
||||
// SetPendingCompaction configures LLM-based auto-compaction for pending messages.
|
||||
func (c *Channel) SetPendingCompaction(cfg *channels.CompactionConfig) {
|
||||
c.groupHistory.SetCompactionConfig(cfg)
|
||||
}
|
||||
|
||||
// Stop shuts down the Telegram bot by cancelling the long polling context
|
||||
// and waiting for the polling goroutine to exit.
|
||||
func (c *Channel) Stop(_ context.Context) error {
|
||||
|
||||
@@ -123,6 +123,11 @@ func (c *Channel) Start(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetPendingCompaction configures LLM-based auto-compaction for pending messages.
|
||||
func (c *Channel) SetPendingCompaction(cfg *channels.CompactionConfig) {
|
||||
c.groupHistory.SetCompactionConfig(cfg)
|
||||
}
|
||||
|
||||
// Stop gracefully shuts down the Zalo Personal channel.
|
||||
func (c *Channel) Stop(_ context.Context) error {
|
||||
c.groupHistory.StopFlusher()
|
||||
|
||||
@@ -1,14 +1,23 @@
|
||||
package config
|
||||
|
||||
// PendingCompactionConfig configures LLM-based compaction of pending group messages.
|
||||
// When a group accumulates more than Threshold pending messages, older messages are
|
||||
// summarized by an LLM and replaced with a compact summary, keeping KeepRecent raw messages.
|
||||
type PendingCompactionConfig struct {
|
||||
Threshold int `json:"threshold,omitempty"` // trigger compaction when entries exceed this (default 50)
|
||||
KeepRecent int `json:"keep_recent,omitempty"` // keep this many recent raw messages after compaction (default 15)
|
||||
}
|
||||
|
||||
// ChannelsConfig contains per-channel configuration.
|
||||
type ChannelsConfig struct {
|
||||
Telegram TelegramConfig `json:"telegram"`
|
||||
Discord DiscordConfig `json:"discord"`
|
||||
Slack SlackConfig `json:"slack"`
|
||||
WhatsApp WhatsAppConfig `json:"whatsapp"`
|
||||
Zalo ZaloConfig `json:"zalo"`
|
||||
ZaloPersonal ZaloPersonalConfig `json:"zalo_personal"`
|
||||
Feishu FeishuConfig `json:"feishu"`
|
||||
Telegram TelegramConfig `json:"telegram"`
|
||||
Discord DiscordConfig `json:"discord"`
|
||||
Slack SlackConfig `json:"slack"`
|
||||
WhatsApp WhatsAppConfig `json:"whatsapp"`
|
||||
Zalo ZaloConfig `json:"zalo"`
|
||||
ZaloPersonal ZaloPersonalConfig `json:"zalo_personal"`
|
||||
Feishu FeishuConfig `json:"feishu"`
|
||||
PendingCompaction *PendingCompactionConfig `json:"pending_compaction,omitempty"` // global pending message compaction settings
|
||||
}
|
||||
|
||||
type TelegramConfig struct {
|
||||
|
||||
@@ -16,14 +16,19 @@ import (
|
||||
// PendingMessagesHandler handles pending message HTTP endpoints.
|
||||
type PendingMessagesHandler struct {
|
||||
store store.PendingMessageStore
|
||||
agentStore store.AgentStore
|
||||
token string
|
||||
providerReg *providers.Registry
|
||||
keepRecent int // global keepRecent from config (0 = use default 15)
|
||||
}
|
||||
|
||||
func NewPendingMessagesHandler(s store.PendingMessageStore, token string, providerReg *providers.Registry) *PendingMessagesHandler {
|
||||
return &PendingMessagesHandler{store: s, token: token, providerReg: providerReg}
|
||||
func NewPendingMessagesHandler(s store.PendingMessageStore, agentStore store.AgentStore, token string, providerReg *providers.Registry) *PendingMessagesHandler {
|
||||
return &PendingMessagesHandler{store: s, agentStore: agentStore, token: token, providerReg: providerReg}
|
||||
}
|
||||
|
||||
// SetKeepRecent sets the global keepRecent value from config.
|
||||
func (h *PendingMessagesHandler) SetKeepRecent(n int) { h.keepRecent = n }
|
||||
|
||||
func (h *PendingMessagesHandler) RegisterRoutes(mux *http.ServeMux) {
|
||||
mux.HandleFunc("GET /v1/pending-messages", h.authMiddleware(h.handleListGroups))
|
||||
mux.HandleFunc("GET /v1/pending-messages/messages", h.authMiddleware(h.handleListMessages))
|
||||
@@ -121,8 +126,8 @@ func (h *PendingMessagesHandler) handleCompact(w http.ResponseWriter, r *http.Re
|
||||
return
|
||||
}
|
||||
|
||||
// Resolve an LLM provider for summarization
|
||||
provider := h.resolveProvider()
|
||||
// Resolve an LLM provider for summarization using the default agent's config
|
||||
provider, model := h.resolveProviderAndModel()
|
||||
if provider == nil {
|
||||
// Fallback: hard delete if no provider available
|
||||
slog.Warn("compact.no_provider", "channel", req.ChannelName, "key", req.HistoryKey)
|
||||
@@ -134,10 +139,14 @@ func (h *PendingMessagesHandler) handleCompact(w http.ResponseWriter, r *http.Re
|
||||
return
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(r.Context(), 45*time.Second)
|
||||
ctx, cancel := context.WithTimeout(r.Context(), 180*time.Second)
|
||||
defer cancel()
|
||||
|
||||
remaining, err := channels.CompactGroup(ctx, h.store, req.ChannelName, req.HistoryKey, provider, provider.DefaultModel(), 15)
|
||||
keepRecent := h.keepRecent
|
||||
if keepRecent <= 0 {
|
||||
keepRecent = 15
|
||||
}
|
||||
remaining, err := channels.CompactGroup(ctx, h.store, req.ChannelName, req.HistoryKey, provider, model, keepRecent)
|
||||
if err != nil {
|
||||
slog.Warn("compact.failed", "channel", req.ChannelName, "key", req.HistoryKey, "error", err)
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
@@ -146,18 +155,38 @@ func (h *PendingMessagesHandler) handleCompact(w http.ResponseWriter, r *http.Re
|
||||
writeJSON(w, http.StatusOK, map[string]interface{}{"status": "ok", "method": "summarized", "remaining": remaining})
|
||||
}
|
||||
|
||||
// resolveProvider returns the first available LLM provider, or nil.
|
||||
func (h *PendingMessagesHandler) resolveProvider() providers.Provider {
|
||||
// resolveProviderAndModel resolves the LLM provider+model for pending message compaction.
|
||||
// Uses the default agent's configured provider and model so compaction uses the same
|
||||
// LLM as the agent that processes these messages.
|
||||
func (h *PendingMessagesHandler) resolveProviderAndModel() (providers.Provider, string) {
|
||||
if h.providerReg == nil {
|
||||
return nil
|
||||
return nil, ""
|
||||
}
|
||||
names := h.providerReg.List()
|
||||
if len(names) == 0 {
|
||||
return nil
|
||||
|
||||
// Use the default agent's provider+model
|
||||
if h.agentStore != nil {
|
||||
if ag, err := h.agentStore.GetDefault(context.Background()); err == nil && ag.Provider != "" {
|
||||
if p, err := h.providerReg.Get(ag.Provider); err == nil {
|
||||
model := ag.Model
|
||||
if model == "" {
|
||||
model = p.DefaultModel()
|
||||
}
|
||||
if model != "" {
|
||||
return p, model
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
p, err := h.providerReg.Get(names[0])
|
||||
if err != nil {
|
||||
return nil
|
||||
|
||||
// Fallback: first provider with a valid default model
|
||||
for _, name := range h.providerReg.List() {
|
||||
p, err := h.providerReg.Get(name)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
if p.DefaultModel() != "" {
|
||||
return p, p.DefaultModel()
|
||||
}
|
||||
}
|
||||
return p
|
||||
return nil, ""
|
||||
}
|
||||
@@ -124,6 +124,13 @@
|
||||
"behavior.blockReplyInfo": "Only the final response is sent. Intermediate tool status messages are suppressed.",
|
||||
"behavior.intentClassifyHint": "Classify user intent before routing to reduce unnecessary agent invocations.",
|
||||
"behavior.intentClassifyInfo": "Agent will only be invoked when the classifier detects actionable intent.",
|
||||
"behavior.pendingCompactionTitle": "Pending Message Compaction",
|
||||
"behavior.pendingCompactionDescription": "Automatically summarize old group messages using LLM when the buffer exceeds a threshold",
|
||||
"behavior.pendingCompactionThreshold": "Threshold",
|
||||
"behavior.pendingCompactionThresholdTip": "Trigger compaction when a group accumulates more than this many pending messages. Set to 0 to use the default (50).",
|
||||
"behavior.pendingCompactionKeepRecent": "Keep Recent",
|
||||
"behavior.pendingCompactionKeepRecentTip": "Number of recent raw messages to keep after compaction. Older messages are replaced with an LLM summary. Default is 15.",
|
||||
"behavior.pendingCompactionInfo": "When a group chat exceeds the threshold, the LLM automatically summarizes older messages into a single compact entry, keeping the most recent ones intact for context.",
|
||||
"behavior.rateLimitTitle": "Rate Limiting",
|
||||
"behavior.rateLimitDescription": "Inbound message size limits and per-user request throttling",
|
||||
"behavior.sessionsTitle": "Session Scoping",
|
||||
|
||||
@@ -3,6 +3,14 @@
|
||||
"description": "Buffered channel messages awaiting agent processing",
|
||||
"emptyTitle": "No pending messages",
|
||||
"emptyDescription": "No buffered message groups found. Messages appear here when channels buffer incoming messages before agent processing.",
|
||||
"howItWorks": {
|
||||
"title": "How it works",
|
||||
"step1": "Messages from channels (Telegram, Discord, etc.) are buffered here while the agent is busy or offline.",
|
||||
"step2": "When the buffer reaches the threshold (default: 50 messages), auto-compaction kicks in — an LLM summarizes older messages, keeping only the 15 most recent.",
|
||||
"step3": "On the agent's next turn, the compacted summary + recent messages are injected as context, so the agent catches up without being overwhelmed.",
|
||||
"compactAction": "Compact — Manually trigger LLM summarization now. Useful when you want the agent to process a clean summary instead of raw message flood.",
|
||||
"clearAction": "Clear — Permanently delete all buffered messages for a group. Use when messages are no longer relevant."
|
||||
},
|
||||
"columns": {
|
||||
"channel": "Channel",
|
||||
"group": "Group",
|
||||
@@ -28,5 +36,13 @@
|
||||
"title": "Messages — {{name}}",
|
||||
"noMessages": "No messages found.",
|
||||
"summary": "Summary"
|
||||
},
|
||||
"toast": {
|
||||
"compacted": "Compaction complete",
|
||||
"compactedSummarized": "Summarized to {{count}} messages",
|
||||
"compactedDeleted": "All messages deleted (no LLM provider available)",
|
||||
"failedCompact": "Compaction failed",
|
||||
"cleared": "Messages cleared",
|
||||
"failedClear": "Failed to clear messages"
|
||||
}
|
||||
}
|
||||
@@ -124,6 +124,13 @@
|
||||
"behavior.blockReplyInfo": "Chỉ phản hồi cuối được gửi. Thông báo trạng thái công cụ trung gian bị ẩn.",
|
||||
"behavior.intentClassifyHint": "Phân loại ý định người dùng trước khi định tuyến để giảm các lần gọi agent không cần thiết.",
|
||||
"behavior.intentClassifyInfo": "Agent chỉ được gọi khi bộ phân loại phát hiện ý định có thể thực hiện.",
|
||||
"behavior.pendingCompactionTitle": "Nén tin nhắn chờ",
|
||||
"behavior.pendingCompactionDescription": "Tự động tóm tắt tin nhắn nhóm cũ bằng LLM khi bộ đệm vượt ngưỡng",
|
||||
"behavior.pendingCompactionThreshold": "Ngưỡng",
|
||||
"behavior.pendingCompactionThresholdTip": "Kích hoạt nén khi nhóm tích lũy nhiều hơn số tin nhắn chờ này. Đặt 0 để dùng mặc định (50).",
|
||||
"behavior.pendingCompactionKeepRecent": "Giữ gần đây",
|
||||
"behavior.pendingCompactionKeepRecentTip": "Số tin nhắn gần đây giữ lại sau khi nén. Tin cũ hơn được thay bằng bản tóm tắt LLM. Mặc định là 15.",
|
||||
"behavior.pendingCompactionInfo": "Khi nhóm chat vượt ngưỡng, LLM tự động tóm tắt các tin nhắn cũ thành một mục rút gọn, giữ nguyên các tin gần đây nhất làm ngữ cảnh.",
|
||||
"behavior.rateLimitTitle": "Giới hạn tốc độ",
|
||||
"behavior.rateLimitDescription": "Giới hạn kích thước tin nhắn đến và kiểm soát tốc độ yêu cầu mỗi người dùng",
|
||||
"behavior.sessionsTitle": "Phạm vi phiên",
|
||||
|
||||
@@ -3,6 +3,14 @@
|
||||
"description": "Tin nhắn channel đã được lưu đệm đang chờ xử lý từ agent",
|
||||
"emptyTitle": "Không có tin nhắn đang chờ",
|
||||
"emptyDescription": "Không tìm thấy nhóm tin nhắn đệm. Tin nhắn xuất hiện ở đây khi channel lưu đệm tin nhắn đến trước khi agent xử lý.",
|
||||
"howItWorks": {
|
||||
"title": "Cơ chế hoạt động",
|
||||
"step1": "Tin nhắn từ các channel (Telegram, Discord, v.v.) được lưu đệm tại đây khi agent đang bận hoặc offline.",
|
||||
"step2": "Khi bộ đệm đạt ngưỡng (mặc định: 50 tin nhắn), hệ thống tự động tóm tắt — LLM tóm gọn các tin nhắn cũ, chỉ giữ lại 15 tin mới nhất.",
|
||||
"step3": "Ở lượt tiếp theo của agent, bản tóm tắt + tin nhắn gần đây được inject làm ngữ cảnh, giúp agent nắm bắt nhanh mà không bị quá tải.",
|
||||
"compactAction": "Tóm tắt — Kích hoạt LLM tóm tắt thủ công ngay. Hữu ích khi bạn muốn agent xử lý bản tóm tắt sạch thay vì hàng loạt tin nhắn thô.",
|
||||
"clearAction": "Xóa — Xóa vĩnh viễn tất cả tin nhắn đệm của một nhóm. Dùng khi tin nhắn không còn liên quan."
|
||||
},
|
||||
"columns": {
|
||||
"channel": "Channel",
|
||||
"group": "Nhóm",
|
||||
@@ -28,5 +36,13 @@
|
||||
"title": "Tin nhắn — {{name}}",
|
||||
"noMessages": "Không tìm thấy tin nhắn.",
|
||||
"summary": "Tóm tắt"
|
||||
},
|
||||
"toast": {
|
||||
"compacted": "Tóm tắt hoàn tất",
|
||||
"compactedSummarized": "Đã tóm tắt còn {{count}} tin nhắn",
|
||||
"compactedDeleted": "Đã xóa tất cả (không có LLM provider)",
|
||||
"failedCompact": "Tóm tắt thất bại",
|
||||
"cleared": "Đã xóa tin nhắn",
|
||||
"failedClear": "Không thể xóa tin nhắn"
|
||||
}
|
||||
}
|
||||
@@ -124,6 +124,13 @@
|
||||
"behavior.blockReplyInfo": "仅发送最终响应,中间工具状态消息被抑制。",
|
||||
"behavior.intentClassifyHint": "在路由之前对用户意图进行分类,以减少不必要的 Agent 调用。",
|
||||
"behavior.intentClassifyInfo": "仅当分类器检测到可执行意图时才调用 Agent。",
|
||||
"behavior.pendingCompactionTitle": "待处理消息压缩",
|
||||
"behavior.pendingCompactionDescription": "当缓冲区超过阈值时,自动使用 LLM 总结旧群组消息",
|
||||
"behavior.pendingCompactionThreshold": "阈值",
|
||||
"behavior.pendingCompactionThresholdTip": "当群组累积超过此数量的待处理消息时触发压缩。设为 0 使用默认值(50)。",
|
||||
"behavior.pendingCompactionKeepRecent": "保留最近",
|
||||
"behavior.pendingCompactionKeepRecentTip": "压缩后保留的最近原始消息数量。较旧的消息将替换为 LLM 摘要。默认为 15。",
|
||||
"behavior.pendingCompactionInfo": "当群聊超过阈值时,LLM 自动将旧消息总结为一条精简条目,保留最近的消息作为上下文。",
|
||||
"behavior.rateLimitTitle": "速率限制",
|
||||
"behavior.rateLimitDescription": "入站消息大小限制和每用户请求限流",
|
||||
"behavior.sessionsTitle": "会话范围",
|
||||
|
||||
@@ -3,6 +3,14 @@
|
||||
"description": "已缓冲的Channel消息,等待Agent处理",
|
||||
"emptyTitle": "暂无待处理消息",
|
||||
"emptyDescription": "未找到缓冲消息组。当Channel在Agent处理前缓冲传入消息时,消息将显示在此处。",
|
||||
"howItWorks": {
|
||||
"title": "工作机制",
|
||||
"step1": "来自各Channel(Telegram、Discord等)的消息在Agent忙碌或离线时缓冲在此。",
|
||||
"step2": "当缓冲区达到阈值(默认:50条消息)时,系统自动压缩 — LLM将旧消息总结为摘要,仅保留最近15条。",
|
||||
"step3": "在Agent的下一轮回复中,压缩摘要 + 最近消息作为上下文注入,让Agent快速了解情况而不会被淹没。",
|
||||
"compactAction": "压缩 — 立即手动触发LLM摘要。当您希望Agent处理简洁摘要而非大量原始消息时使用。",
|
||||
"clearAction": "清除 — 永久删除某组的所有缓冲消息。当消息不再相关时使用。"
|
||||
},
|
||||
"columns": {
|
||||
"channel": "Channel",
|
||||
"group": "组",
|
||||
@@ -28,5 +36,13 @@
|
||||
"title": "消息 — {{name}}",
|
||||
"noMessages": "未找到消息。",
|
||||
"summary": "摘要"
|
||||
},
|
||||
"toast": {
|
||||
"compacted": "压缩完成",
|
||||
"compactedSummarized": "已压缩为 {{count}} 条消息",
|
||||
"compactedDeleted": "已删除全部(无可用LLM Provider)",
|
||||
"failedCompact": "压缩失败",
|
||||
"cleared": "消息已清除",
|
||||
"failedClear": "清除消息失败"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,99 @@
|
||||
import { Archive, Clock, Info } from "lucide-react";
|
||||
import { useTranslation } from "react-i18next";
|
||||
import { Input } from "@/components/ui/input";
|
||||
import { Label } from "@/components/ui/label";
|
||||
import {
|
||||
Card,
|
||||
CardContent,
|
||||
CardDescription,
|
||||
CardHeader,
|
||||
CardTitle,
|
||||
} from "@/components/ui/card";
|
||||
|
||||
export interface PendingCompactionValues {
|
||||
threshold?: number;
|
||||
keep_recent?: number;
|
||||
}
|
||||
|
||||
interface Props {
|
||||
value: PendingCompactionValues;
|
||||
onChange: (v: PendingCompactionValues) => void;
|
||||
}
|
||||
|
||||
/** Global pending message compaction thresholds with visual emphasis. */
|
||||
export function BehaviorPendingCompactionCard({ value, onChange }: Props) {
|
||||
const { t } = useTranslation("config");
|
||||
|
||||
const update = (patch: Partial<PendingCompactionValues>) =>
|
||||
onChange({ ...value, ...patch });
|
||||
|
||||
return (
|
||||
<Card>
|
||||
<CardHeader className="pb-3">
|
||||
<CardTitle className="text-base">
|
||||
{t("behavior.pendingCompactionTitle")}
|
||||
</CardTitle>
|
||||
<CardDescription>
|
||||
{t("behavior.pendingCompactionDescription")}
|
||||
</CardDescription>
|
||||
</CardHeader>
|
||||
<CardContent className="space-y-0">
|
||||
{/* Threshold */}
|
||||
<div className="border-b py-4">
|
||||
<div className="flex items-start justify-between gap-4">
|
||||
<div className="flex items-start gap-3">
|
||||
<Archive className="mt-0.5 h-4 w-4 shrink-0 text-orange-500" />
|
||||
<div className="space-y-1">
|
||||
<Label className="text-sm font-medium">
|
||||
{t("behavior.pendingCompactionThreshold")}
|
||||
</Label>
|
||||
<p className="text-xs text-muted-foreground">
|
||||
{t("behavior.pendingCompactionThresholdTip")}
|
||||
</p>
|
||||
</div>
|
||||
</div>
|
||||
<Input
|
||||
type="number"
|
||||
value={value.threshold ?? ""}
|
||||
onChange={(e) => update({ threshold: Number(e.target.value) })}
|
||||
placeholder="50"
|
||||
min={0}
|
||||
className="w-24 shrink-0"
|
||||
/>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
{/* Keep Recent */}
|
||||
<div className="py-4">
|
||||
<div className="flex items-start justify-between gap-4">
|
||||
<div className="flex items-start gap-3">
|
||||
<Clock className="mt-0.5 h-4 w-4 shrink-0 text-blue-500" />
|
||||
<div className="space-y-1">
|
||||
<Label className="text-sm font-medium">
|
||||
{t("behavior.pendingCompactionKeepRecent")}
|
||||
</Label>
|
||||
<p className="text-xs text-muted-foreground">
|
||||
{t("behavior.pendingCompactionKeepRecentTip")}
|
||||
</p>
|
||||
</div>
|
||||
</div>
|
||||
<Input
|
||||
type="number"
|
||||
value={value.keep_recent ?? ""}
|
||||
onChange={(e) => update({ keep_recent: Number(e.target.value) })}
|
||||
placeholder="15"
|
||||
min={1}
|
||||
className="w-24 shrink-0"
|
||||
/>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
{/* Info banner */}
|
||||
<div className="flex items-start gap-2 rounded-md border border-orange-200 bg-orange-50 px-3 py-2 text-xs text-orange-700 dark:border-orange-800 dark:bg-orange-950/30 dark:text-orange-300">
|
||||
<Info className="mt-0.5 h-3.5 w-3.5 shrink-0" />
|
||||
<span>{t("behavior.pendingCompactionInfo")}</span>
|
||||
</div>
|
||||
</CardContent>
|
||||
</Card>
|
||||
);
|
||||
}
|
||||
@@ -6,6 +6,7 @@ import { BehaviorUxCard } from "./behavior-ux-card";
|
||||
import { BehaviorRateCard } from "./behavior-rate-card";
|
||||
import { BehaviorSessionsCard } from "./behavior-sessions-card";
|
||||
import { BehaviorSecurityCard } from "./behavior-security-card";
|
||||
import { BehaviorPendingCompactionCard, type PendingCompactionValues } from "./behavior-pending-compaction-card";
|
||||
|
||||
/* eslint-disable @typescript-eslint/no-explicit-any */
|
||||
|
||||
@@ -22,6 +23,7 @@ export function BehaviorSection({ config, onPatch, saving }: Props) {
|
||||
const ag = config.agents?.defaults ?? {};
|
||||
const tl = config.tools ?? {};
|
||||
const ss = config.sessions ?? {};
|
||||
const ch = config.channels ?? {};
|
||||
|
||||
// UX toggles (from gateway + agents.defaults)
|
||||
const [ux, setUx] = useState({
|
||||
@@ -49,6 +51,11 @@ export function BehaviorSection({ config, onPatch, saving }: Props) {
|
||||
scrub_credentials: tl.scrub_credentials,
|
||||
});
|
||||
|
||||
// Pending compaction (from channels.pending_compaction)
|
||||
const [pendingCompaction, setPendingCompaction] = useState<PendingCompactionValues>(
|
||||
ch.pending_compaction ?? {},
|
||||
);
|
||||
|
||||
const [dirty, setDirty] = useState(false);
|
||||
|
||||
// Reset drafts when external config changes
|
||||
@@ -68,6 +75,7 @@ export function BehaviorSection({ config, onPatch, saving }: Props) {
|
||||
injection_action: gw.injection_action,
|
||||
scrub_credentials: tl.scrub_credentials,
|
||||
});
|
||||
setPendingCompaction(ch.pending_compaction ?? {});
|
||||
setDirty(false);
|
||||
}, [config]); // eslint-disable-line react-hooks/exhaustive-deps
|
||||
|
||||
@@ -91,6 +99,7 @@ export function BehaviorSection({ config, onPatch, saving }: Props) {
|
||||
},
|
||||
tools: { ...tl, scrub_credentials: security.scrub_credentials },
|
||||
sessions: { ...ss, ...sessions },
|
||||
channels: { ...ch, pending_compaction: pendingCompaction },
|
||||
});
|
||||
};
|
||||
|
||||
@@ -100,6 +109,7 @@ export function BehaviorSection({ config, onPatch, saving }: Props) {
|
||||
<BehaviorRateCard value={rate} onChange={markDirty(setRate)} />
|
||||
<BehaviorSessionsCard value={sessions} onChange={markDirty(setSessions)} />
|
||||
<BehaviorSecurityCard value={security} onChange={markDirty(setSecurity)} />
|
||||
<BehaviorPendingCompactionCard value={pendingCompaction} onChange={markDirty(setPendingCompaction)} />
|
||||
|
||||
{dirty && (
|
||||
<div className="flex justify-end pt-2">
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
import { useState, useCallback } from "react";
|
||||
import i18next from "i18next";
|
||||
import { useHttp } from "@/hooks/use-ws";
|
||||
import { toast } from "@/stores/use-toast-store";
|
||||
import type { PendingMessageGroup, PendingMessage } from "../types";
|
||||
|
||||
export function usePendingMessages() {
|
||||
@@ -42,12 +44,23 @@ export function usePendingMessages() {
|
||||
const compactGroup = useCallback(
|
||||
async (channel: string, key: string) => {
|
||||
try {
|
||||
await http.post("/v1/pending-messages/compact", {
|
||||
channel_name: channel,
|
||||
history_key: key,
|
||||
});
|
||||
const res = await http.post<{ status: string; method?: string; remaining?: number }>(
|
||||
"/v1/pending-messages/compact",
|
||||
{ channel_name: channel, history_key: key },
|
||||
);
|
||||
const method = res?.method ?? "summarized";
|
||||
toast.success(
|
||||
i18next.t("pending-messages:toast.compacted"),
|
||||
method === "deleted"
|
||||
? i18next.t("pending-messages:toast.compactedDeleted")
|
||||
: i18next.t("pending-messages:toast.compactedSummarized", { count: res?.remaining ?? 0 }),
|
||||
);
|
||||
return true;
|
||||
} catch {
|
||||
} catch (err) {
|
||||
toast.error(
|
||||
i18next.t("pending-messages:toast.failedCompact"),
|
||||
err instanceof Error ? err.message : "",
|
||||
);
|
||||
return false;
|
||||
}
|
||||
},
|
||||
@@ -58,8 +71,13 @@ export function usePendingMessages() {
|
||||
async (channel: string, key: string) => {
|
||||
try {
|
||||
await http.delete(`/v1/pending-messages?channel=${encodeURIComponent(channel)}&key=${encodeURIComponent(key)}`);
|
||||
toast.success(i18next.t("pending-messages:toast.cleared"));
|
||||
return true;
|
||||
} catch {
|
||||
} catch (err) {
|
||||
toast.error(
|
||||
i18next.t("pending-messages:toast.failedClear"),
|
||||
err instanceof Error ? err.message : "",
|
||||
);
|
||||
return false;
|
||||
}
|
||||
},
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { useState, useEffect } from "react";
|
||||
import { useTranslation } from "react-i18next";
|
||||
import { Inbox, RefreshCw, Trash2, Archive, Loader2 } from "lucide-react";
|
||||
import { Inbox, RefreshCw, Trash2, Archive, Loader2, Info, ChevronDown, ChevronUp } from "lucide-react";
|
||||
import { Button } from "@/components/ui/button";
|
||||
import { Badge } from "@/components/ui/badge";
|
||||
import { PageHeader } from "@/components/shared/page-header";
|
||||
@@ -72,6 +72,8 @@ export function PendingMessagesPage() {
|
||||
}
|
||||
/>
|
||||
|
||||
<HowItWorksCard />
|
||||
|
||||
<div className="mt-4">
|
||||
{showSkeleton ? (
|
||||
<TableSkeleton rows={6} />
|
||||
@@ -181,6 +183,50 @@ export function PendingMessagesPage() {
|
||||
);
|
||||
}
|
||||
|
||||
function HowItWorksCard() {
|
||||
const { t } = useTranslation("pending-messages");
|
||||
const [open, setOpen] = useState(true);
|
||||
|
||||
return (
|
||||
<div className="mt-4">
|
||||
<button
|
||||
type="button"
|
||||
onClick={() => setOpen(!open)}
|
||||
className="flex w-full items-center gap-2 rounded-lg border bg-muted/30 px-4 py-2.5 text-left text-sm transition-colors hover:bg-muted/50"
|
||||
>
|
||||
<Info className="h-4 w-4 shrink-0 text-blue-500" />
|
||||
<span className="font-medium">{t("howItWorks.title")}</span>
|
||||
{open ? (
|
||||
<ChevronUp className="ml-auto h-4 w-4 text-muted-foreground" />
|
||||
) : (
|
||||
<ChevronDown className="ml-auto h-4 w-4 text-muted-foreground" />
|
||||
)}
|
||||
</button>
|
||||
{open && (
|
||||
<div className="rounded-b-lg border border-t-0 bg-muted/10 px-4 py-3 space-y-2.5 text-sm text-muted-foreground">
|
||||
<div className="flex gap-2.5">
|
||||
<span className="shrink-0 mt-0.5 flex h-5 w-5 items-center justify-center rounded-full bg-blue-500/10 text-xs font-semibold text-blue-600">1</span>
|
||||
<p>{t("howItWorks.step1")}</p>
|
||||
</div>
|
||||
<div className="flex gap-2.5">
|
||||
<span className="shrink-0 mt-0.5 flex h-5 w-5 items-center justify-center rounded-full bg-blue-500/10 text-xs font-semibold text-blue-600">2</span>
|
||||
<p>{t("howItWorks.step2")}</p>
|
||||
</div>
|
||||
<div className="flex gap-2.5">
|
||||
<span className="shrink-0 mt-0.5 flex h-5 w-5 items-center justify-center rounded-full bg-blue-500/10 text-xs font-semibold text-blue-600">3</span>
|
||||
<p>{t("howItWorks.step3")}</p>
|
||||
</div>
|
||||
<hr className="border-border/50" />
|
||||
<div className="space-y-1.5 text-xs">
|
||||
<p><Archive className="mr-1 inline h-3 w-3" /><strong>{t("howItWorks.compactAction")}</strong></p>
|
||||
<p><Trash2 className="mr-1 inline h-3 w-3" /><strong>{t("howItWorks.clearAction")}</strong></p>
|
||||
</div>
|
||||
</div>
|
||||
)}
|
||||
</div>
|
||||
);
|
||||
}
|
||||
|
||||
function ConfirmClearDialog({
|
||||
group,
|
||||
onConfirm,
|
||||
|
||||
Reference in new issue
Block a user