mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-03 05:20:01 +00:00
feat(channels): real LLM compaction for web UI compact endpoint
Extract CompactGroup() from runCompaction() for shared use by both auto-compact and HTTP compact endpoint. Wire provider registry into PendingMessagesHandler so /v1/pending-messages/compact performs actual LLM summarization (keep 15 recent + summary) instead of hard delete. Add rows-affected guard in Compact() to prevent duplicate summaries when concurrent compactions race on the same key.
This commit is contained in:
1 parent
30080f1acf
commit
4df60649e5
4 files changed
+121
-60
No files matched your search
@@ -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)
|
||||
pendingMessagesH = httpapi.NewPendingMessagesHandler(stores.PendingMessages, token, providerReg)
|
||||
}
|
||||
|
||||
return agentsH, skillsH, tracesH, mcpH, customToolsH, channelInstancesH, providersH, delegationsH, builtinToolsH, pendingMessagesH
|
||||
|
||||
@@ -41,49 +41,32 @@ func (ph *PendingHistory) MaybeCompact(historyKey string, currentCount int, cfg
|
||||
go ph.runCompaction(historyKey, cfg)
|
||||
}
|
||||
|
||||
// runCompaction performs LLM-based summarization of old messages.
|
||||
// Follows pattern from internal/agent/loop_compact.go.
|
||||
func (ph *PendingHistory) runCompaction(historyKey string, cfg *CompactionConfig) {
|
||||
defer ph.compacting.Delete(historyKey)
|
||||
|
||||
// Step 1: Force-flush buffer to ensure DB is consistent
|
||||
ph.flushNow()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 45*time.Second)
|
||||
defer cancel()
|
||||
|
||||
// Step 2: Read entries from DB
|
||||
entries, err := ph.store.ListByKey(ctx, ph.channelName, historyKey)
|
||||
if err != nil {
|
||||
slog.Warn("compaction.list_failed", "channel", ph.channelName, "key", historyKey, "error", err)
|
||||
return
|
||||
}
|
||||
|
||||
threshold := cfg.Threshold
|
||||
if threshold <= 0 {
|
||||
threshold = DefaultGroupHistoryLimit
|
||||
}
|
||||
if len(entries) <= threshold {
|
||||
return // cleared or below threshold
|
||||
}
|
||||
|
||||
// Step 3: Split entries
|
||||
keepRecent := cfg.KeepRecent
|
||||
// CompactGroup performs LLM-based compaction on a pending message group.
|
||||
// Reused by both auto-compact (channel) and HTTP compact endpoint.
|
||||
// Returns the number of entries remaining after compaction.
|
||||
func CompactGroup(ctx context.Context, s store.PendingMessageStore, channelName, historyKey string, provider providers.Provider, model string, keepRecent int) (int, error) {
|
||||
if keepRecent <= 0 {
|
||||
keepRecent = 15
|
||||
}
|
||||
if keepRecent >= len(entries) {
|
||||
return
|
||||
}
|
||||
splitIdx := len(entries) - keepRecent
|
||||
|
||||
// Step 1: Read entries from DB
|
||||
entries, err := s.ListByKey(ctx, channelName, historyKey)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("list entries: %w", err)
|
||||
}
|
||||
if keepRecent >= len(entries) {
|
||||
return len(entries), nil // nothing to summarize
|
||||
}
|
||||
|
||||
// Step 2: Split into old (to summarize) and recent (to keep)
|
||||
splitIdx := len(entries) - keepRecent
|
||||
toSummarize := entries[:splitIdx]
|
||||
deleteIDs := make([]uuid.UUID, len(toSummarize))
|
||||
for i, e := range toSummarize {
|
||||
deleteIDs[i] = e.ID
|
||||
}
|
||||
|
||||
// Step 4: Build text and call LLM
|
||||
// Step 3: Build text and call LLM for summarization
|
||||
var sb strings.Builder
|
||||
for _, e := range toSummarize {
|
||||
prefix := e.Sender
|
||||
@@ -97,33 +80,73 @@ func (ph *PendingHistory) runCompaction(historyKey string, cfg *CompactionConfig
|
||||
fmt.Fprintf(&sb, "%s%s: %s\n", prefix, ts, e.Body)
|
||||
}
|
||||
|
||||
resp, err := cfg.Provider.Chat(ctx, providers.ChatRequest{
|
||||
resp, err := provider.Chat(ctx, providers.ChatRequest{
|
||||
Messages: []providers.Message{{
|
||||
Role: "user",
|
||||
Content: "Summarize these group chat messages concisely, preserving key topics, decisions, names, and important context:\n\n" + sb.String(),
|
||||
}},
|
||||
Model: cfg.Model,
|
||||
Model: model,
|
||||
Options: map[string]any{"max_tokens": 512, "temperature": 0.3},
|
||||
})
|
||||
if err != nil {
|
||||
slog.Warn("compaction.llm_failed", "channel", ph.channelName, "key", historyKey, "error", err)
|
||||
return
|
||||
return 0, fmt.Errorf("llm summarize: %w", err)
|
||||
}
|
||||
|
||||
// Step 5: Compact in DB (atomic tx: delete old + insert summary)
|
||||
// Step 4: Compact in DB (atomic tx: delete old + insert summary)
|
||||
summary := &store.PendingMessage{
|
||||
ChannelName: ph.channelName,
|
||||
ChannelName: channelName,
|
||||
HistoryKey: historyKey,
|
||||
Sender: "[summary]",
|
||||
Body: resp.Content,
|
||||
IsSummary: true,
|
||||
}
|
||||
if err := ph.store.Compact(ctx, deleteIDs, summary); err != nil {
|
||||
slog.Warn("compaction.db_failed", "channel", ph.channelName, "key", historyKey, "error", err)
|
||||
if err := s.Compact(ctx, deleteIDs, summary); err != nil {
|
||||
return 0, fmt.Errorf("db compact: %w", err)
|
||||
}
|
||||
|
||||
remaining := keepRecent + 1 // kept messages + new summary
|
||||
slog.Info("compaction.done",
|
||||
"channel", channelName,
|
||||
"key", historyKey,
|
||||
"summarized", len(toSummarize),
|
||||
"kept", keepRecent,
|
||||
"total_after", remaining,
|
||||
)
|
||||
return remaining, nil
|
||||
}
|
||||
|
||||
// runCompaction performs LLM-based summarization of old messages.
|
||||
// Wraps CompactGroup with flush + RAM update + threshold check.
|
||||
func (ph *PendingHistory) runCompaction(historyKey string, cfg *CompactionConfig) {
|
||||
defer ph.compacting.Delete(historyKey)
|
||||
|
||||
// Force-flush buffer to ensure DB is consistent
|
||||
ph.flushNow()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 45*time.Second)
|
||||
defer cancel()
|
||||
|
||||
// Check threshold from DB (may have been cleared since trigger)
|
||||
entries, err := ph.store.ListByKey(ctx, ph.channelName, historyKey)
|
||||
if err != nil {
|
||||
slog.Warn("compaction.list_failed", "channel", ph.channelName, "key", historyKey, "error", err)
|
||||
return
|
||||
}
|
||||
threshold := cfg.Threshold
|
||||
if threshold <= 0 {
|
||||
threshold = DefaultGroupHistoryLimit
|
||||
}
|
||||
if len(entries) <= threshold {
|
||||
return
|
||||
}
|
||||
|
||||
// Step 6: Update RAM from DB
|
||||
_, err = CompactGroup(ctx, ph.store, ph.channelName, historyKey, cfg.Provider, cfg.Model, cfg.KeepRecent)
|
||||
if err != nil {
|
||||
slog.Warn("compaction.failed", "channel", ph.channelName, "key", historyKey, "error", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Update RAM from DB
|
||||
ph.mu.Lock()
|
||||
if _, exists := ph.entries[historyKey]; !exists {
|
||||
// Key was Clear()ed during compaction — remove stale summary
|
||||
@@ -132,7 +155,6 @@ func (ph *PendingHistory) runCompaction(historyKey string, cfg *CompactionConfig
|
||||
slog.Info("compaction.cleared_stale", "channel", ph.channelName, "key", historyKey)
|
||||
return
|
||||
}
|
||||
// Re-read from DB to get complete current state (summary + kept + new entries)
|
||||
fresh, err := ph.store.ListByKey(ctx, ph.channelName, historyKey)
|
||||
if err == nil {
|
||||
rebuilt := make([]HistoryEntry, 0, len(fresh))
|
||||
@@ -147,12 +169,4 @@ func (ph *PendingHistory) runCompaction(historyKey string, cfg *CompactionConfig
|
||||
ph.entries[historyKey] = rebuilt
|
||||
}
|
||||
ph.mu.Unlock()
|
||||
|
||||
slog.Info("compaction.done",
|
||||
"channel", ph.channelName,
|
||||
"key", historyKey,
|
||||
"summarized", len(toSummarize),
|
||||
"kept", keepRecent,
|
||||
"total_after", len(fresh),
|
||||
)
|
||||
}
|
||||
@@ -1,21 +1,27 @@
|
||||
package http
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/nextlevelbuilder/goclaw/internal/channels"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/i18n"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/providers"
|
||||
"github.com/nextlevelbuilder/goclaw/internal/store"
|
||||
)
|
||||
|
||||
// PendingMessagesHandler handles pending message HTTP endpoints.
|
||||
type PendingMessagesHandler struct {
|
||||
store store.PendingMessageStore
|
||||
token string
|
||||
store store.PendingMessageStore
|
||||
token string
|
||||
providerReg *providers.Registry
|
||||
}
|
||||
|
||||
func NewPendingMessagesHandler(s store.PendingMessageStore, token string) *PendingMessagesHandler {
|
||||
return &PendingMessagesHandler{store: s, token: token}
|
||||
func NewPendingMessagesHandler(s store.PendingMessageStore, token string, providerReg *providers.Registry) *PendingMessagesHandler {
|
||||
return &PendingMessagesHandler{store: s, token: token, providerReg: providerReg}
|
||||
}
|
||||
|
||||
func (h *PendingMessagesHandler) RegisterRoutes(mux *http.ServeMux) {
|
||||
@@ -101,7 +107,8 @@ type compactRequest struct {
|
||||
HistoryKey string `json:"history_key"`
|
||||
}
|
||||
|
||||
// POST /v1/pending-messages/compact — MVP: clear group and return success
|
||||
// POST /v1/pending-messages/compact — LLM-based summarization of old messages, keeping recent ones.
|
||||
// Falls back to hard delete if no LLM provider is available.
|
||||
func (h *PendingMessagesHandler) handleCompact(w http.ResponseWriter, r *http.Request) {
|
||||
locale := store.LocaleFromContext(r.Context())
|
||||
var req compactRequest
|
||||
@@ -114,9 +121,43 @@ func (h *PendingMessagesHandler) handleCompact(w http.ResponseWriter, r *http.Re
|
||||
return
|
||||
}
|
||||
|
||||
if err := h.store.DeleteByKey(r.Context(), req.ChannelName, req.HistoryKey); err != nil {
|
||||
// Resolve an LLM provider for summarization
|
||||
provider := h.resolveProvider()
|
||||
if provider == nil {
|
||||
// Fallback: hard delete if no provider available
|
||||
slog.Warn("compact.no_provider", "channel", req.ChannelName, "key", req.HistoryKey)
|
||||
if err := h.store.DeleteByKey(r.Context(), req.ChannelName, req.HistoryKey); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]string{"status": "ok", "method": "deleted"})
|
||||
return
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(r.Context(), 45*time.Second)
|
||||
defer cancel()
|
||||
|
||||
remaining, err := channels.CompactGroup(ctx, h.store, req.ChannelName, req.HistoryKey, provider, provider.DefaultModel(), 15)
|
||||
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()})
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]string{"status": "ok"})
|
||||
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 {
|
||||
if h.providerReg == nil {
|
||||
return nil
|
||||
}
|
||||
names := h.providerReg.List()
|
||||
if len(names) == 0 {
|
||||
return nil
|
||||
}
|
||||
p, err := h.providerReg.Get(names[0])
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
return p
|
||||
}
|
||||
@@ -103,7 +103,7 @@ func (s *PGPendingMessageStore) Compact(ctx context.Context, deleteIDs []uuid.UU
|
||||
args[i] = id
|
||||
}
|
||||
|
||||
_, err = tx.ExecContext(ctx,
|
||||
res, err := tx.ExecContext(ctx,
|
||||
fmt.Sprintf("DELETE FROM channel_pending_messages WHERE id IN (%s)", strings.Join(placeholders, ",")),
|
||||
args...,
|
||||
)
|
||||
@@ -111,6 +111,12 @@ func (s *PGPendingMessageStore) Compact(ctx context.Context, deleteIDs []uuid.UU
|
||||
return fmt.Errorf("compact delete: %w", err)
|
||||
}
|
||||
|
||||
// Guard: if another compaction already deleted these rows, skip summary insertion
|
||||
affected, _ := res.RowsAffected()
|
||||
if affected == 0 {
|
||||
return nil // already compacted by concurrent caller
|
||||
}
|
||||
|
||||
// Insert summary row
|
||||
if summary.ID == uuid.Nil {
|
||||
summary.ID = uuid.Must(uuid.NewV7())
|
||||
|
||||
Reference in new issue
Block a user