mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-11 12:18:59 +00:00
fix(channelmemory): honor Discord parent channel excludes (#1390)
* fix(channelmemory): harden passive extraction * fix(channelmemory): honor discord parent channel excludes --------- Co-authored-by: Collective Developer <man@collective.dev>
This commit is contained in:
1 parent
8367a9bd0c
commit
7263f771c8
21 files changed
+272
-104
No files matched your search
@@ -47,6 +47,7 @@ type ProcessAllEvent struct {
|
||||
Type string `json:"type"`
|
||||
ChannelName string `json:"channel_name,omitempty"`
|
||||
HistoryKey string `json:"history_key,omitempty"`
|
||||
GroupMessageCount int `json:"group_message_count,omitempty"`
|
||||
Run *store.ChannelMemoryExtractionRun `json:"run,omitempty"`
|
||||
Error string `json:"error,omitempty"`
|
||||
RunCount int `json:"run_count"`
|
||||
@@ -57,12 +58,13 @@ type ProcessAllEvent struct {
|
||||
}
|
||||
|
||||
type GroupOption struct {
|
||||
ChannelName string `json:"channel_name"`
|
||||
HistoryKey string `json:"history_key"`
|
||||
GroupTitle string `json:"group_title,omitempty"`
|
||||
MessageCount int `json:"message_count"`
|
||||
LastActivity time.Time `json:"last_activity"`
|
||||
Excluded bool `json:"excluded"`
|
||||
ChannelName string `json:"channel_name"`
|
||||
HistoryKey string `json:"history_key"`
|
||||
ParentHistoryKey string `json:"parent_history_key,omitempty"`
|
||||
GroupTitle string `json:"group_title,omitempty"`
|
||||
MessageCount int `json:"message_count"`
|
||||
LastActivity time.Time `json:"last_activity"`
|
||||
Excluded bool `json:"excluded"`
|
||||
}
|
||||
|
||||
func (s *Service) Status(ctx context.Context, inst *store.ChannelInstanceData) (*Status, error) {
|
||||
@@ -109,12 +111,13 @@ func (s *Service) GroupOptions(ctx context.Context, inst *store.ChannelInstanceD
|
||||
continue
|
||||
}
|
||||
out = append(out, GroupOption{
|
||||
ChannelName: group.ChannelName,
|
||||
HistoryKey: group.HistoryKey,
|
||||
GroupTitle: titles[group.ChannelName+":"+group.HistoryKey],
|
||||
MessageCount: group.MessageCount,
|
||||
LastActivity: group.LastActivity,
|
||||
Excluded: contains(cfg.ExcludeHistoryKeys, group.HistoryKey),
|
||||
ChannelName: group.ChannelName,
|
||||
HistoryKey: group.HistoryKey,
|
||||
ParentHistoryKey: group.ParentHistoryKey,
|
||||
GroupTitle: titles[group.ChannelName+":"+group.HistoryKey],
|
||||
MessageCount: group.MessageCount,
|
||||
LastActivity: group.LastActivity,
|
||||
Excluded: contains(cfg.ExcludeHistoryKeys, group.HistoryKey) || contains(cfg.ExcludeHistoryKeys, group.ParentHistoryKey),
|
||||
})
|
||||
}
|
||||
return out, nil
|
||||
@@ -130,7 +133,7 @@ func (s *Service) RunNow(ctx context.Context, inst *store.ChannelInstanceData, t
|
||||
return nil, err
|
||||
}
|
||||
for _, group := range groups {
|
||||
if group.ChannelName != inst.Name || !eligibleHistoryKey(group.HistoryKey, cfg) {
|
||||
if group.ChannelName != inst.Name || !eligibleHistoryGroup(group, cfg) {
|
||||
continue
|
||||
}
|
||||
messages, err := s.unprocessedMessages(ctx, inst.ID, group)
|
||||
@@ -163,7 +166,7 @@ func (s *Service) RunAllWithProgress(ctx context.Context, inst *store.ChannelIns
|
||||
}
|
||||
result := &ProcessAllResult{}
|
||||
for _, group := range groups {
|
||||
if group.ChannelName != inst.Name || !eligibleHistoryKey(group.HistoryKey, cfg) {
|
||||
if group.ChannelName != inst.Name || !eligibleHistoryGroup(group, cfg) {
|
||||
continue
|
||||
}
|
||||
messages, err := s.unprocessedMessages(ctx, inst.ID, group)
|
||||
@@ -178,7 +181,7 @@ func (s *Service) RunAllWithProgress(ctx context.Context, inst *store.ChannelIns
|
||||
}
|
||||
if len(messages) < cfg.MinMessages {
|
||||
result.SkippedGroupCount++
|
||||
if err := emitProcessAllEvent(emit, "group_skipped", group, nil, "", result); err != nil {
|
||||
if err := emitProcessAllEvent(emit, "group_skipped", group, len(messages), nil, "", result); err != nil {
|
||||
return result, err
|
||||
}
|
||||
continue
|
||||
@@ -186,7 +189,7 @@ func (s *Service) RunAllWithProgress(ctx context.Context, inst *store.ChannelIns
|
||||
run, err := s.runMessages(ctx, inst, cfg, group, messages, trigger)
|
||||
if err != nil {
|
||||
result.ErrorCount++
|
||||
if emitErr := emitProcessAllEvent(emit, "group_failed", group, nil, err.Error(), result); emitErr != nil {
|
||||
if emitErr := emitProcessAllEvent(emit, "group_failed", group, len(messages), nil, err.Error(), result); emitErr != nil {
|
||||
return result, emitErr
|
||||
}
|
||||
continue
|
||||
@@ -195,17 +198,17 @@ func (s *Service) RunAllWithProgress(ctx context.Context, inst *store.ChannelIns
|
||||
result.RunCount++
|
||||
result.MessageCount += run.MessageCount
|
||||
result.ItemCount += run.ItemCount
|
||||
if err := emitProcessAllEvent(emit, "group_completed", group, run, "", result); err != nil {
|
||||
if err := emitProcessAllEvent(emit, "group_completed", group, run.MessageCount, run, "", result); err != nil {
|
||||
return result, err
|
||||
}
|
||||
}
|
||||
if err := emitProcessAllEvent(emit, "final", store.PendingMessageGroup{}, nil, "", result); err != nil {
|
||||
if err := emitProcessAllEvent(emit, "final", store.PendingMessageGroup{}, 0, nil, "", result); err != nil {
|
||||
return result, err
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func emitProcessAllEvent(emit func(ProcessAllEvent) error, typ string, group store.PendingMessageGroup, run *store.ChannelMemoryExtractionRun, errMsg string, result *ProcessAllResult) error {
|
||||
func emitProcessAllEvent(emit func(ProcessAllEvent) error, typ string, group store.PendingMessageGroup, groupMessageCount int, run *store.ChannelMemoryExtractionRun, errMsg string, result *ProcessAllResult) error {
|
||||
if emit == nil {
|
||||
return nil
|
||||
}
|
||||
@@ -213,6 +216,7 @@ func emitProcessAllEvent(emit func(ProcessAllEvent) error, typ string, group sto
|
||||
Type: typ,
|
||||
ChannelName: group.ChannelName,
|
||||
HistoryKey: group.HistoryKey,
|
||||
GroupMessageCount: groupMessageCount,
|
||||
Run: run,
|
||||
Error: errMsg,
|
||||
RunCount: result.RunCount,
|
||||
@@ -231,7 +235,7 @@ func (s *Service) UnprocessedMessageCount(ctx context.Context, inst *store.Chann
|
||||
}
|
||||
total := 0
|
||||
for _, group := range groups {
|
||||
if group.ChannelName != inst.Name || !eligibleHistoryKey(group.HistoryKey, cfg) {
|
||||
if group.ChannelName != inst.Name || !eligibleHistoryGroup(group, cfg) {
|
||||
continue
|
||||
}
|
||||
messages, err := s.unprocessedMessages(ctx, inst.ID, group)
|
||||
|
||||
@@ -258,7 +258,11 @@ func TestRunAllSkipsLowVolumeGroupsAndContinues(t *testing.T) {
|
||||
}
|
||||
svc := &Service{Pending: pending, Extractions: &fakeExtractionStore{}}
|
||||
|
||||
result, err := svc.RunAll(context.Background(), inst, "scheduled")
|
||||
var events []ProcessAllEvent
|
||||
result, err := svc.RunAllWithProgress(context.Background(), inst, "scheduled", func(event ProcessAllEvent) error {
|
||||
events = append(events, event)
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("RunAll returned error: %v", err)
|
||||
}
|
||||
@@ -268,6 +272,15 @@ func TestRunAllSkipsLowVolumeGroupsAndContinues(t *testing.T) {
|
||||
if result.RunCount != 0 {
|
||||
t.Fatalf("expected no runs for low-volume groups, got %d", result.RunCount)
|
||||
}
|
||||
if len(events) < 2 {
|
||||
t.Fatalf("expected skipped events, got %d", len(events))
|
||||
}
|
||||
if events[0].Type != "group_skipped" || events[0].GroupMessageCount != 1 {
|
||||
t.Fatalf("first skipped event = %+v, want group_message_count 1", events[0])
|
||||
}
|
||||
if events[1].Type != "group_skipped" || events[1].GroupMessageCount != 2 {
|
||||
t.Fatalf("second skipped event = %+v, want group_message_count 2", events[1])
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunAllSkipsExcludedHistoryKeys(t *testing.T) {
|
||||
@@ -299,6 +312,49 @@ func TestRunAllSkipsExcludedHistoryKeys(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunAllSkipsThreadWhenParentHistoryKeyExcluded(t *testing.T) {
|
||||
inst := &store.ChannelInstanceData{
|
||||
BaseModel: store.BaseModel{ID: uuid.New()},
|
||||
Name: "discord",
|
||||
Config: MergeIntoInstanceConfig(nil, Config{Enabled: true, MinMessages: 2, ExcludeHistoryKeys: []string{"parent-channel"}}),
|
||||
}
|
||||
pending := &fakePendingStore{
|
||||
groups: []store.PendingMessageGroup{
|
||||
{ChannelName: "discord", HistoryKey: "thread-1", ParentHistoryKey: "parent-channel", MessageCount: 3},
|
||||
},
|
||||
messages: map[string][]store.PendingMessage{
|
||||
"discord:thread-1": {
|
||||
{ID: uuid.New(), ChannelName: "discord", HistoryKey: "thread-1", ParentHistoryKey: "parent-channel", Body: "one", CreatedAt: time.Now().UTC()},
|
||||
{ID: uuid.New(), ChannelName: "discord", HistoryKey: "thread-1", ParentHistoryKey: "parent-channel", Body: "two", CreatedAt: time.Now().UTC()},
|
||||
{ID: uuid.New(), ChannelName: "discord", HistoryKey: "thread-1", ParentHistoryKey: "parent-channel", Body: "three", CreatedAt: time.Now().UTC()},
|
||||
},
|
||||
},
|
||||
}
|
||||
svc := &Service{Pending: pending, Extractions: &fakeExtractionStore{}}
|
||||
|
||||
result, err := svc.RunAll(context.Background(), inst, "scheduled")
|
||||
if err != nil {
|
||||
t.Fatalf("RunAll returned error: %v", err)
|
||||
}
|
||||
if result.RunCount != 0 || result.SkippedGroupCount != 0 {
|
||||
t.Fatalf("thread under excluded parent should not be processed or counted: %+v", result)
|
||||
}
|
||||
count, err := svc.UnprocessedMessageCount(context.Background(), inst)
|
||||
if err != nil {
|
||||
t.Fatalf("UnprocessedMessageCount returned error: %v", err)
|
||||
}
|
||||
if count != 0 {
|
||||
t.Fatalf("unprocessed count = %d, want 0 for thread under excluded parent", count)
|
||||
}
|
||||
options, err := svc.GroupOptions(context.Background(), inst)
|
||||
if err != nil {
|
||||
t.Fatalf("GroupOptions returned error: %v", err)
|
||||
}
|
||||
if len(options) != 1 || !options[0].Excluded || options[0].ParentHistoryKey != "parent-channel" {
|
||||
t.Fatalf("group option = %+v, want excluded thread with parent key", options)
|
||||
}
|
||||
}
|
||||
|
||||
func TestItemHashIsStableAcrossRuns(t *testing.T) {
|
||||
runA := &store.ChannelMemoryExtractionRun{ID: uuid.New(), ChannelInstanceID: uuid.New(), HistoryKey: "group"}
|
||||
runB := *runA
|
||||
|
||||
@@ -42,6 +42,13 @@ func eligibleHistoryKey(key string, cfg Config) bool {
|
||||
return key != "" && !strings.Contains(k, "dm") && !strings.Contains(k, "private")
|
||||
}
|
||||
|
||||
func eligibleHistoryGroup(group store.PendingMessageGroup, cfg Config) bool {
|
||||
if group.ParentHistoryKey != "" && slices.Contains(cfg.ExcludeHistoryKeys, group.ParentHistoryKey) {
|
||||
return false
|
||||
}
|
||||
return eligibleHistoryKey(group.HistoryKey, cfg)
|
||||
}
|
||||
|
||||
func messageSourceID(msg store.PendingMessage) string {
|
||||
if msg.PlatformMsgID != "" {
|
||||
return msg.PlatformMsgID
|
||||
|
||||
@@ -196,13 +196,15 @@ func (c *Channel) handleMessage(_ *discordgo.Session, m *discordgo.MessageCreate
|
||||
mediaPaths = append(mediaPaths, mf.Path)
|
||||
}
|
||||
}
|
||||
parentHistoryKey := c.parentHistoryKeyForChannel(ctx, channelID)
|
||||
c.GroupHistory().Record(channelID, channels.HistoryEntry{
|
||||
Sender: senderName,
|
||||
SenderID: senderID,
|
||||
Body: content,
|
||||
Media: mediaPaths,
|
||||
Timestamp: m.Timestamp,
|
||||
MessageID: m.ID,
|
||||
Sender: senderName,
|
||||
SenderID: senderID,
|
||||
Body: content,
|
||||
ParentHistoryKey: parentHistoryKey,
|
||||
Media: mediaPaths,
|
||||
Timestamp: m.Timestamp,
|
||||
MessageID: m.ID,
|
||||
}, c.HistoryLimit())
|
||||
|
||||
// Collect contact even when bot is not mentioned (cache prevents DB spam).
|
||||
|
||||
@@ -157,6 +157,32 @@ func TestDiscordThreadBackfillDoesNotDuplicatePendingThreadHistory(t *testing.T)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDiscordThreadPendingHistoryStoresParentChannel(t *testing.T) {
|
||||
server := newDiscordThreadBackfillServer(t, discordThreadBackfillFixture{})
|
||||
defer server.Close()
|
||||
ch, _ := newThreadBackfillTestChannel(t, server)
|
||||
|
||||
ch.handleMessage(ch.session, &discordgo.MessageCreate{Message: &discordgo.Message{
|
||||
ID: "prior-1",
|
||||
ChannelID: "thread-1",
|
||||
GuildID: "guild-1",
|
||||
Content: "unmentioned thread context",
|
||||
Author: &discordgo.User{ID: "user-1", Username: "Alice"},
|
||||
Timestamp: time.Now(),
|
||||
}})
|
||||
|
||||
entries := ch.GroupHistory().GetEntries("thread-1")
|
||||
if len(entries) != 1 {
|
||||
t.Fatalf("pending entries = %d, want 1: %#v", len(entries), entries)
|
||||
}
|
||||
if entries[0].ParentHistoryKey != "parent-1" {
|
||||
t.Fatalf("parent history key = %q, want parent-1", entries[0].ParentHistoryKey)
|
||||
}
|
||||
if entries[0].Body != "unmentioned thread context" {
|
||||
t.Fatalf("entry body = %q, want recorded content", entries[0].Body)
|
||||
}
|
||||
}
|
||||
|
||||
type discordThreadBackfillFixture struct {
|
||||
channelJSON string
|
||||
historyStatus int
|
||||
|
||||
@@ -90,20 +90,33 @@ func (c *Channel) backfillThreadHistory(ctx context.Context, m *discordgo.Messag
|
||||
}
|
||||
|
||||
func (c *Channel) isThreadChannel(ctx context.Context, channelID string) bool {
|
||||
_, ok := c.threadChannel(ctx, channelID)
|
||||
return ok
|
||||
}
|
||||
|
||||
func (c *Channel) parentHistoryKeyForChannel(ctx context.Context, channelID string) string {
|
||||
ch, ok := c.threadChannel(ctx, channelID)
|
||||
if !ok || ch.ParentID == "" {
|
||||
return ""
|
||||
}
|
||||
return ch.ParentID
|
||||
}
|
||||
|
||||
func (c *Channel) threadChannel(ctx context.Context, channelID string) (*discordgo.Channel, bool) {
|
||||
if c.session == nil {
|
||||
return false
|
||||
return nil, false
|
||||
}
|
||||
if c.session.State != nil {
|
||||
if ch, err := c.session.State.Channel(channelID); err == nil && ch != nil {
|
||||
return ch.IsThread()
|
||||
return ch, ch.IsThread()
|
||||
}
|
||||
}
|
||||
ch, err := c.session.Channel(channelID, discordgo.WithContext(ctx))
|
||||
if err != nil {
|
||||
slog.Warn("discord: thread channel lookup failed", "channel_id", channelID, "error", err)
|
||||
return false
|
||||
return nil, false
|
||||
}
|
||||
return ch != nil && ch.IsThread()
|
||||
return ch, ch != nil && ch.IsThread()
|
||||
}
|
||||
|
||||
func (c *Channel) resolveThreadHistoryAttachments(ctx context.Context, attachments []*discordgo.MessageAttachment, maxBytes int64, remaining int) []threadHistoryAttachment {
|
||||
|
||||
@@ -45,13 +45,14 @@ type MediaRef struct {
|
||||
|
||||
// HistoryEntry represents a single tracked group message.
|
||||
type HistoryEntry struct {
|
||||
Sender string
|
||||
SenderID string
|
||||
Body string
|
||||
Media []string // temp file paths for images/attachments (RAM-only, not persisted to DB)
|
||||
MediaRefs []MediaRef // deferred media refs for lazy download (RAM-only, not persisted)
|
||||
Timestamp time.Time
|
||||
MessageID string
|
||||
Sender string
|
||||
SenderID string
|
||||
Body string
|
||||
ParentHistoryKey string
|
||||
Media []string // temp file paths for images/attachments (RAM-only, not persisted to DB)
|
||||
MediaRefs []MediaRef // deferred media refs for lazy download (RAM-only, not persisted)
|
||||
Timestamp time.Time
|
||||
MessageID string
|
||||
}
|
||||
|
||||
// PendingHistory tracks group messages across multiple groups.
|
||||
@@ -181,13 +182,14 @@ func (ph *PendingHistory) Record(historyKey string, entry HistoryEntry, limit in
|
||||
// Queue for DB persistence (batched flush)
|
||||
if ph.store != nil {
|
||||
ph.enqueueFlush(store.PendingMessage{
|
||||
ChannelName: ph.channelName,
|
||||
HistoryKey: historyKey,
|
||||
Sender: entry.Sender,
|
||||
SenderID: entry.SenderID,
|
||||
Body: entry.Body,
|
||||
PlatformMsgID: entry.MessageID,
|
||||
CreatedAt: entry.Timestamp,
|
||||
ChannelName: ph.channelName,
|
||||
HistoryKey: historyKey,
|
||||
ParentHistoryKey: entry.ParentHistoryKey,
|
||||
Sender: entry.Sender,
|
||||
SenderID: entry.SenderID,
|
||||
Body: entry.Body,
|
||||
PlatformMsgID: entry.MessageID,
|
||||
CreatedAt: entry.Timestamp,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -215,11 +217,12 @@ func (ph *PendingHistory) loadFromDB(historyKey string) []HistoryEntry {
|
||||
entries := make([]HistoryEntry, 0, len(msgs))
|
||||
for _, m := range msgs {
|
||||
entries = append(entries, HistoryEntry{
|
||||
Sender: m.Sender,
|
||||
SenderID: m.SenderID,
|
||||
Body: m.Body,
|
||||
Timestamp: m.CreatedAt,
|
||||
MessageID: m.PlatformMsgID,
|
||||
Sender: m.Sender,
|
||||
SenderID: m.SenderID,
|
||||
Body: m.Body,
|
||||
ParentHistoryKey: m.ParentHistoryKey,
|
||||
Timestamp: m.CreatedAt,
|
||||
MessageID: m.PlatformMsgID,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -115,11 +115,12 @@ func CompactGroup(ctx context.Context, s store.PendingMessageStore, channelName,
|
||||
|
||||
// Step 4: Compact in DB (atomic tx: delete old + insert summary)
|
||||
summary := &store.PendingMessage{
|
||||
ChannelName: channelName,
|
||||
HistoryKey: historyKey,
|
||||
Sender: "[summary]",
|
||||
Body: resp.Content,
|
||||
IsSummary: true,
|
||||
ChannelName: channelName,
|
||||
HistoryKey: historyKey,
|
||||
ParentHistoryKey: firstParentHistoryKey(entries),
|
||||
Sender: "[summary]",
|
||||
Body: resp.Content,
|
||||
IsSummary: true,
|
||||
}
|
||||
if err := s.Compact(ctx, deleteIDs, summary); err != nil {
|
||||
return 0, fmt.Errorf("db compact: %w", err)
|
||||
@@ -136,6 +137,15 @@ func CompactGroup(ctx context.Context, s store.PendingMessageStore, channelName,
|
||||
return remaining, nil
|
||||
}
|
||||
|
||||
func firstParentHistoryKey(entries []store.PendingMessage) string {
|
||||
for _, entry := range entries {
|
||||
if entry.ParentHistoryKey != "" {
|
||||
return entry.ParentHistoryKey
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// sweepCompaction checks DB for groups exceeding compaction threshold.
|
||||
// Called periodically from flushLoop as a safety net for post-restart scenarios
|
||||
// where RAM is empty but DB has accumulated messages.
|
||||
@@ -212,10 +222,11 @@ func (ph *PendingHistory) runCompaction(historyKey string, cfg *CompactionConfig
|
||||
rebuilt := make([]HistoryEntry, 0, len(fresh))
|
||||
for _, f := range fresh {
|
||||
rebuilt = append(rebuilt, HistoryEntry{
|
||||
Sender: f.Sender,
|
||||
Body: f.Body,
|
||||
Timestamp: f.CreatedAt,
|
||||
MessageID: f.PlatformMsgID,
|
||||
Sender: f.Sender,
|
||||
Body: f.Body,
|
||||
ParentHistoryKey: f.ParentHistoryKey,
|
||||
Timestamp: f.CreatedAt,
|
||||
MessageID: f.PlatformMsgID,
|
||||
})
|
||||
}
|
||||
ph.entries[historyKey] = rebuilt
|
||||
|
||||
@@ -10,26 +10,28 @@ import (
|
||||
// PendingMessage represents a buffered group chat message (or LLM-generated summary)
|
||||
// stored in channel_pending_messages table.
|
||||
type PendingMessage struct {
|
||||
ID uuid.UUID `json:"id" db:"id"`
|
||||
ChannelName string `json:"channel_name" db:"channel_name"`
|
||||
HistoryKey string `json:"history_key" db:"history_key"`
|
||||
Sender string `json:"sender" db:"sender"`
|
||||
SenderID string `json:"sender_id" db:"sender_id"`
|
||||
Body string `json:"body" db:"body"`
|
||||
PlatformMsgID string `json:"platform_msg_id" db:"platform_msg_id"`
|
||||
IsSummary bool `json:"is_summary" db:"is_summary"`
|
||||
CreatedAt time.Time `json:"created_at" db:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at" db:"updated_at"`
|
||||
ID uuid.UUID `json:"id" db:"id"`
|
||||
ChannelName string `json:"channel_name" db:"channel_name"`
|
||||
HistoryKey string `json:"history_key" db:"history_key"`
|
||||
ParentHistoryKey string `json:"parent_history_key,omitempty" db:"parent_history_key"`
|
||||
Sender string `json:"sender" db:"sender"`
|
||||
SenderID string `json:"sender_id" db:"sender_id"`
|
||||
Body string `json:"body" db:"body"`
|
||||
PlatformMsgID string `json:"platform_msg_id" db:"platform_msg_id"`
|
||||
IsSummary bool `json:"is_summary" db:"is_summary"`
|
||||
CreatedAt time.Time `json:"created_at" db:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at" db:"updated_at"`
|
||||
}
|
||||
|
||||
// PendingMessageGroup is a summary row for the grouped overview page.
|
||||
type PendingMessageGroup struct {
|
||||
ChannelName string `json:"channel_name" db:"channel_name"`
|
||||
HistoryKey string `json:"history_key" db:"history_key"`
|
||||
GroupTitle string `json:"group_title,omitempty" db:"group_title"`
|
||||
MessageCount int `json:"message_count" db:"message_count"`
|
||||
HasSummary bool `json:"has_summary" db:"has_summary"`
|
||||
LastActivity time.Time `json:"last_activity" db:"last_activity"`
|
||||
ChannelName string `json:"channel_name" db:"channel_name"`
|
||||
HistoryKey string `json:"history_key" db:"history_key"`
|
||||
ParentHistoryKey string `json:"parent_history_key,omitempty" db:"parent_history_key"`
|
||||
GroupTitle string `json:"group_title,omitempty" db:"group_title"`
|
||||
MessageCount int `json:"message_count" db:"message_count"`
|
||||
HasSummary bool `json:"has_summary" db:"has_summary"`
|
||||
LastActivity time.Time `json:"last_activity" db:"last_activity"`
|
||||
}
|
||||
|
||||
// PendingMessageStore persists group chat messages for context when bot is mentioned.
|
||||
|
||||
@@ -27,8 +27,8 @@ func (s *PGPendingMessageStore) AppendBatch(ctx context.Context, msgs []store.Pe
|
||||
return nil
|
||||
}
|
||||
|
||||
// Build multi-row INSERT: VALUES ($1,$2,...,$11), ($12,$13,...,$22), ...
|
||||
const cols = 11
|
||||
// Build multi-row INSERT: VALUES ($1,$2,...,$12), ($13,$14,...), ...
|
||||
const cols = 12
|
||||
placeholders := make([]string, len(msgs))
|
||||
args := make([]any, 0, len(msgs)*cols)
|
||||
now := time.Now()
|
||||
@@ -39,14 +39,14 @@ func (s *PGPendingMessageStore) AppendBatch(ctx context.Context, msgs []store.Pe
|
||||
msgs[i].ID = uuid.Must(uuid.NewV7())
|
||||
}
|
||||
base := i * cols
|
||||
placeholders[i] = fmt.Sprintf("($%d,$%d,$%d,$%d,$%d,$%d,$%d,$%d,$%d,$%d,$%d)",
|
||||
base+1, base+2, base+3, base+4, base+5, base+6, base+7, base+8, base+9, base+10, base+11)
|
||||
args = append(args, msgs[i].ID, msgs[i].ChannelName, msgs[i].HistoryKey,
|
||||
placeholders[i] = fmt.Sprintf("($%d,$%d,$%d,$%d,$%d,$%d,$%d,$%d,$%d,$%d,$%d,$%d)",
|
||||
base+1, base+2, base+3, base+4, base+5, base+6, base+7, base+8, base+9, base+10, base+11, base+12)
|
||||
args = append(args, msgs[i].ID, msgs[i].ChannelName, msgs[i].HistoryKey, msgs[i].ParentHistoryKey,
|
||||
msgs[i].Sender, msgs[i].SenderID, msgs[i].Body, msgs[i].PlatformMsgID, msgs[i].IsSummary, now, now, tid)
|
||||
}
|
||||
|
||||
_, err := s.db.ExecContext(ctx,
|
||||
`INSERT INTO channel_pending_messages (id, channel_name, history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at, tenant_id)
|
||||
`INSERT INTO channel_pending_messages (id, channel_name, history_key, parent_history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at, tenant_id)
|
||||
VALUES `+strings.Join(placeholders, ","),
|
||||
args...,
|
||||
)
|
||||
@@ -60,7 +60,7 @@ func (s *PGPendingMessageStore) ListByKey(ctx context.Context, channelName, hist
|
||||
}
|
||||
var result []store.PendingMessage
|
||||
err = pkgSqlxDB.SelectContext(ctx, &result,
|
||||
`SELECT id, channel_name, history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at
|
||||
`SELECT id, channel_name, history_key, parent_history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at
|
||||
FROM channel_pending_messages
|
||||
WHERE channel_name = $1 AND history_key = $2`+tClause+`
|
||||
ORDER BY created_at ASC, id ASC`,
|
||||
@@ -120,9 +120,9 @@ func (s *PGPendingMessageStore) Compact(ctx context.Context, deleteIDs []uuid.UU
|
||||
}
|
||||
now := time.Now()
|
||||
_, err = tx.ExecContext(ctx,
|
||||
`INSERT INTO channel_pending_messages (id, channel_name, history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at, tenant_id)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)`,
|
||||
summary.ID, summary.ChannelName, summary.HistoryKey, summary.Sender, summary.SenderID, summary.Body, summary.PlatformMsgID, true, now, now, tenantIDForInsert(ctx),
|
||||
`INSERT INTO channel_pending_messages (id, channel_name, history_key, parent_history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at, tenant_id)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)`,
|
||||
summary.ID, summary.ChannelName, summary.HistoryKey, summary.ParentHistoryKey, summary.Sender, summary.SenderID, summary.Body, summary.PlatformMsgID, true, now, now, tenantIDForInsert(ctx),
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("compact insert summary: %w", err)
|
||||
@@ -155,7 +155,7 @@ func (s *PGPendingMessageStore) ListGroups(ctx context.Context) ([]store.Pending
|
||||
}
|
||||
var result []store.PendingMessageGroup
|
||||
err = pkgSqlxDB.SelectContext(ctx, &result,
|
||||
`SELECT channel_name, history_key,
|
||||
`SELECT channel_name, history_key, MAX(parent_history_key) AS parent_history_key,
|
||||
COUNT(*) AS message_count,
|
||||
BOOL_OR(is_summary)
|
||||
AND NOT EXISTS (
|
||||
|
||||
@@ -28,7 +28,7 @@ func (s *SQLitePendingMessageStore) AppendBatch(ctx context.Context, msgs []stor
|
||||
return nil
|
||||
}
|
||||
|
||||
const cols = 11
|
||||
const cols = 12
|
||||
placeholders := make([]string, len(msgs))
|
||||
args := make([]any, 0, len(msgs)*cols)
|
||||
now := time.Now()
|
||||
@@ -38,13 +38,13 @@ func (s *SQLitePendingMessageStore) AppendBatch(ctx context.Context, msgs []stor
|
||||
if msgs[i].ID == uuid.Nil {
|
||||
msgs[i].ID = uuid.Must(uuid.NewV7())
|
||||
}
|
||||
placeholders[i] = "(?,?,?,?,?,?,?,?,?,?,?)"
|
||||
args = append(args, msgs[i].ID, msgs[i].ChannelName, msgs[i].HistoryKey,
|
||||
placeholders[i] = "(?,?,?,?,?,?,?,?,?,?,?,?)"
|
||||
args = append(args, msgs[i].ID, msgs[i].ChannelName, msgs[i].HistoryKey, msgs[i].ParentHistoryKey,
|
||||
msgs[i].Sender, msgs[i].SenderID, msgs[i].Body, msgs[i].PlatformMsgID, msgs[i].IsSummary, now, now, tid)
|
||||
}
|
||||
|
||||
_, err := s.db.ExecContext(ctx,
|
||||
`INSERT INTO channel_pending_messages (id, channel_name, history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at, tenant_id)
|
||||
`INSERT INTO channel_pending_messages (id, channel_name, history_key, parent_history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at, tenant_id)
|
||||
VALUES `+strings.Join(placeholders, ","),
|
||||
args...,
|
||||
)
|
||||
@@ -58,7 +58,7 @@ func (s *SQLitePendingMessageStore) ListByKey(ctx context.Context, channelName,
|
||||
}
|
||||
args := append([]any{channelName, historyKey}, tArgs...)
|
||||
rows, err := s.db.QueryContext(ctx,
|
||||
`SELECT id, channel_name, history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at
|
||||
`SELECT id, channel_name, history_key, parent_history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at
|
||||
FROM channel_pending_messages
|
||||
WHERE channel_name = ? AND history_key = ?`+tClause+`
|
||||
ORDER BY created_at ASC, id ASC`,
|
||||
@@ -73,7 +73,7 @@ func (s *SQLitePendingMessageStore) ListByKey(ctx context.Context, channelName,
|
||||
for rows.Next() {
|
||||
var m store.PendingMessage
|
||||
createdAt, updatedAt := scanTimePair()
|
||||
if err := rows.Scan(&m.ID, &m.ChannelName, &m.HistoryKey, &m.Sender, &m.SenderID, &m.Body, &m.PlatformMsgID, &m.IsSummary, createdAt, updatedAt); err != nil {
|
||||
if err := rows.Scan(&m.ID, &m.ChannelName, &m.HistoryKey, &m.ParentHistoryKey, &m.Sender, &m.SenderID, &m.Body, &m.PlatformMsgID, &m.IsSummary, createdAt, updatedAt); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
m.CreatedAt = createdAt.Time
|
||||
@@ -132,9 +132,9 @@ func (s *SQLitePendingMessageStore) Compact(ctx context.Context, deleteIDs []uui
|
||||
}
|
||||
now := time.Now()
|
||||
_, err = tx.ExecContext(ctx,
|
||||
`INSERT INTO channel_pending_messages (id, channel_name, history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at, tenant_id)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
||||
summary.ID, summary.ChannelName, summary.HistoryKey, summary.Sender, summary.SenderID, summary.Body, summary.PlatformMsgID, true, now, now, tenantIDForInsert(ctx),
|
||||
`INSERT INTO channel_pending_messages (id, channel_name, history_key, parent_history_key, sender, sender_id, body, platform_msg_id, is_summary, created_at, updated_at, tenant_id)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
||||
summary.ID, summary.ChannelName, summary.HistoryKey, summary.ParentHistoryKey, summary.Sender, summary.SenderID, summary.Body, summary.PlatformMsgID, true, now, now, tenantIDForInsert(ctx),
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("compact insert summary: %w", err)
|
||||
@@ -162,7 +162,7 @@ func (s *SQLitePendingMessageStore) ListGroups(ctx context.Context) ([]store.Pen
|
||||
return nil, err
|
||||
}
|
||||
// SQLite: BOOL_OR → MAX(is_summary), EXISTS subquery logic preserved
|
||||
q := `SELECT channel_name, history_key,
|
||||
q := `SELECT channel_name, history_key, MAX(parent_history_key) AS parent_history_key,
|
||||
COUNT(*) AS message_count,
|
||||
MAX(is_summary) AND NOT EXISTS (
|
||||
SELECT 1 FROM channel_pending_messages n
|
||||
@@ -192,7 +192,7 @@ func (s *SQLitePendingMessageStore) ListGroups(ctx context.Context) ([]store.Pen
|
||||
var result []store.PendingMessageGroup
|
||||
for rows.Next() {
|
||||
var g store.PendingMessageGroup
|
||||
if err := rows.Scan(&g.ChannelName, &g.HistoryKey, &g.MessageCount, &g.HasSummary, &g.LastActivity); err != nil {
|
||||
if err := rows.Scan(&g.ChannelName, &g.HistoryKey, &g.ParentHistoryKey, &g.MessageCount, &g.HasSummary, &g.LastActivity); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
result = append(result, g)
|
||||
|
||||
@@ -16,7 +16,7 @@ var schemaSQL string
|
||||
|
||||
// SchemaVersion is the current SQLite schema version.
|
||||
// Bump this when adding new migration steps below.
|
||||
const SchemaVersion = 54
|
||||
const SchemaVersion = 55
|
||||
|
||||
// migrations maps version → SQL to apply when upgrading FROM that version.
|
||||
// schema.sql always represents the LATEST full schema (for fresh DBs).
|
||||
@@ -893,6 +893,11 @@ ALTER TABLE usage_event_rollups ADD COLUMN cache_create_tokens BIGINT NOT NULL D
|
||||
ALTER TABLE usage_event_rollups ADD COLUMN thinking_tokens BIGINT NOT NULL DEFAULT 0;`,
|
||||
// Version 53 → 54: dedupe passive memory extraction items across runs for the same channel instance.
|
||||
53: addChannelMemoryItemChannelHashUnique,
|
||||
// Version 54 → 55: preserve Discord thread parent channel for passive-memory excludes.
|
||||
54: `ALTER TABLE channel_pending_messages ADD COLUMN parent_history_key VARCHAR(200) NOT NULL DEFAULT '';
|
||||
CREATE INDEX IF NOT EXISTS idx_channel_pending_messages_parent
|
||||
ON channel_pending_messages(channel_name, parent_history_key)
|
||||
WHERE parent_history_key <> '';`,
|
||||
}
|
||||
|
||||
const addUsageEventAnalyticsTables = `
|
||||
@@ -1452,6 +1457,12 @@ func EnsureSchema(db *sql.DB) error {
|
||||
return fmt.Errorf("inspect usage event token columns: %w", err)
|
||||
}
|
||||
}
|
||||
if v == 54 {
|
||||
patch, err = sqlitePendingMessageParentMigrationPatch(db)
|
||||
if err != nil {
|
||||
return fmt.Errorf("inspect channel pending message parent column: %w", err)
|
||||
}
|
||||
}
|
||||
// Migrations that rebuild a table referenced by another table's FK
|
||||
// require foreign_keys=OFF per SQLite altertable §7. The pragma is
|
||||
// a no-op inside a transaction, so toggle it around BEGIN/COMMIT.
|
||||
@@ -1595,6 +1606,21 @@ func sqliteUsageEventTokenMigrationPatch(db *sql.DB) (string, error) {
|
||||
return patch, nil
|
||||
}
|
||||
|
||||
func sqlitePendingMessageParentMigrationPatch(db *sql.DB) (string, error) {
|
||||
hasColumn, err := sqliteColumnExists(db, "channel_pending_messages", "parent_history_key")
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
patch := ""
|
||||
if !hasColumn {
|
||||
patch += "ALTER TABLE channel_pending_messages ADD COLUMN parent_history_key VARCHAR(200) NOT NULL DEFAULT '';\n"
|
||||
}
|
||||
patch += `CREATE INDEX IF NOT EXISTS idx_channel_pending_messages_parent
|
||||
ON channel_pending_messages(channel_name, parent_history_key)
|
||||
WHERE parent_history_key <> '';`
|
||||
return patch, nil
|
||||
}
|
||||
|
||||
func sqliteColumnExists(db *sql.DB, tableName, columnName string) (bool, error) {
|
||||
rows, err := db.Query("PRAGMA table_info(" + tableName + ")")
|
||||
if err != nil {
|
||||
|
||||
@@ -1166,6 +1166,7 @@ CREATE TABLE IF NOT EXISTS channel_pending_messages (
|
||||
id TEXT NOT NULL PRIMARY KEY,
|
||||
channel_name VARCHAR(100) NOT NULL,
|
||||
history_key VARCHAR(200) NOT NULL,
|
||||
parent_history_key VARCHAR(200) NOT NULL DEFAULT '',
|
||||
sender VARCHAR(255) NOT NULL,
|
||||
sender_id VARCHAR(255) NOT NULL DEFAULT '',
|
||||
body TEXT NOT NULL,
|
||||
@@ -1177,6 +1178,7 @@ CREATE TABLE IF NOT EXISTS channel_pending_messages (
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_channel_pending_messages_lookup ON channel_pending_messages(channel_name, history_key, created_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_channel_pending_messages_parent ON channel_pending_messages(channel_name, parent_history_key) WHERE parent_history_key <> '';
|
||||
CREATE INDEX IF NOT EXISTS idx_channel_pending_messages_tenant ON channel_pending_messages(tenant_id);
|
||||
|
||||
-- ============================================================
|
||||
|
||||
@@ -2,4 +2,4 @@ package upgrade
|
||||
|
||||
// RequiredSchemaVersion is the schema migration version this binary requires.
|
||||
// Bump this whenever adding a new SQL migration file.
|
||||
const RequiredSchemaVersion uint = 90
|
||||
const RequiredSchemaVersion uint = 91
|
||||
@@ -0,0 +1,4 @@
|
||||
DROP INDEX IF EXISTS idx_channel_pending_messages_parent;
|
||||
|
||||
ALTER TABLE channel_pending_messages
|
||||
DROP COLUMN IF EXISTS parent_history_key;
|
||||
@@ -0,0 +1,6 @@
|
||||
ALTER TABLE channel_pending_messages
|
||||
ADD COLUMN IF NOT EXISTS parent_history_key VARCHAR(200) NOT NULL DEFAULT '';
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_channel_pending_messages_parent
|
||||
ON channel_pending_messages (channel_name, parent_history_key)
|
||||
WHERE parent_history_key <> '';
|
||||
@@ -233,7 +233,7 @@
|
||||
"runAllProgress": "Run all groups progress",
|
||||
"runAllProgressStats": "{{runs}} done, {{skipped}} skipped, {{errors}} failed, {{messages}} messages",
|
||||
"groupCompleted": "Completed group {{historyKey}}",
|
||||
"groupSkipped": "Skipped group {{historyKey}} below the minimum message count",
|
||||
"groupSkipped": "Skipped group {{historyKey}} below the minimum message count ({{messageCount}} messages)",
|
||||
"groupFailed": "Group {{historyKey}} failed: {{error}}",
|
||||
"itemUpdated": "Memory item updated",
|
||||
"itemFailed": "Memory item action failed",
|
||||
|
||||
@@ -232,7 +232,7 @@
|
||||
"runAllProgress": "Tiến độ chạy tất cả nhóm",
|
||||
"runAllProgressStats": "{{runs}} xong, {{skipped}} bỏ qua, {{errors}} lỗi, {{messages}} tin nhắn",
|
||||
"groupCompleted": "Đã xử lý nhóm {{historyKey}}",
|
||||
"groupSkipped": "Bỏ qua nhóm {{historyKey}} vì chưa đủ số tin nhắn tối thiểu",
|
||||
"groupSkipped": "Bỏ qua nhóm {{historyKey}} vì chưa đủ số tin nhắn tối thiểu ({{messageCount}} tin nhắn)",
|
||||
"groupFailed": "Nhóm {{historyKey}} lỗi: {{error}}",
|
||||
"itemUpdated": "Đã cập nhật mục bộ nhớ",
|
||||
"itemFailed": "Thao tác mục bộ nhớ thất bại",
|
||||
|
||||
@@ -232,7 +232,7 @@
|
||||
"runAllProgress": "运行全部群组进度",
|
||||
"runAllProgressStats": "{{runs}} 完成,{{skipped}} 跳过,{{errors}} 失败,{{messages}} 条消息",
|
||||
"groupCompleted": "已完成群组 {{historyKey}}",
|
||||
"groupSkipped": "群组 {{historyKey}} 因消息数不足而跳过",
|
||||
"groupSkipped": "群组 {{historyKey}} 因消息数不足而跳过({{messageCount}} 条消息)",
|
||||
"groupFailed": "群组 {{historyKey}} 失败:{{error}}",
|
||||
"itemUpdated": "记忆项已更新",
|
||||
"itemFailed": "记忆项操作失败",
|
||||
|
||||
@@ -90,7 +90,10 @@ export function PassiveMemorySection({ instanceId, channelType }: PassiveMemoryS
|
||||
excluded: true,
|
||||
}
|
||||
));
|
||||
const availableGroupOptions = groupOptions.filter((group) => !excludedHistoryKeys.includes(group.history_key));
|
||||
const availableGroupOptions = groupOptions.filter((group) => {
|
||||
return !excludedHistoryKeys.includes(group.history_key) &&
|
||||
!excludedHistoryKeys.includes(group.parent_history_key ?? "");
|
||||
});
|
||||
|
||||
const updateConfig = (patch: Partial<ChannelMemoryConfig>) => {
|
||||
setConfig((current) => ({ ...current, ...patch }));
|
||||
@@ -366,6 +369,7 @@ function formatGroupLabel(group: { history_key: string; group_title?: string })
|
||||
function formatRunAllEvent(event: {
|
||||
type: string;
|
||||
history_key?: string;
|
||||
group_message_count?: number;
|
||||
error?: string;
|
||||
}, t: (key: string, opts?: Record<string, unknown>) => string) {
|
||||
const historyKey = event.history_key || "-";
|
||||
@@ -373,7 +377,7 @@ function formatRunAllEvent(event: {
|
||||
case "group_completed":
|
||||
return t("detail.passiveMemory.groupCompleted", { historyKey });
|
||||
case "group_skipped":
|
||||
return t("detail.passiveMemory.groupSkipped", { historyKey });
|
||||
return t("detail.passiveMemory.groupSkipped", { historyKey, messageCount: event.group_message_count ?? 0 });
|
||||
case "group_failed":
|
||||
return t("detail.passiveMemory.groupFailed", { historyKey, error: event.error || "" });
|
||||
default:
|
||||
|
||||
@@ -94,6 +94,7 @@ export interface ChannelMemoryProcessAllEvent {
|
||||
type: "group_completed" | "group_skipped" | "group_failed" | "final" | "error";
|
||||
channel_name?: string;
|
||||
history_key?: string;
|
||||
group_message_count?: number;
|
||||
run?: ChannelMemoryExtractionRun;
|
||||
error?: string;
|
||||
run_count: number;
|
||||
@@ -113,6 +114,7 @@ export interface ChannelMemoryItemsResponse {
|
||||
export interface ChannelMemoryGroupOption {
|
||||
channel_name: string;
|
||||
history_key: string;
|
||||
parent_history_key?: string;
|
||||
group_title?: string;
|
||||
message_count: number;
|
||||
last_activity: string;
|
||||
|
||||
Reference in new issue
Block a user