From 1b8b3fa23f52a0bbf5d8331d5e1ec226cc7eebea Mon Sep 17 00:00:00 2001 From: Sebastien Tardif Date: Sat, 29 Aug 2026 11:13:05 -0700 Subject: [PATCH 1/2] fix(syncer): reject non-advancing message page cursors A full 100-message page whose last or newest ID is empty or repeats made bootstrap, forward, and unlimited backfill reprint forever. Fail closed on a stuck before/after cursor, matching GuildMembers. Signed-off-by: Sebastien Tardif --- internal/syncer/message_sync.go | 29 +++- internal/syncer/message_sync_cursor_test.go | 162 ++++++++++++++++++++ 2 files changed, 187 insertions(+), 4 deletions(-) create mode 100644 internal/syncer/message_sync_cursor_test.go diff --git a/internal/syncer/message_sync.go b/internal/syncer/message_sync.go index 84e7784f..5a14335a 100644 --- a/internal/syncer/message_sync.go +++ b/internal/syncer/message_sync.go @@ -417,7 +417,6 @@ func (s *Syncer) bootstrapChannelHistory(ctx context.Context, channel *discordgo } break } - before = page[len(page)-1].ID if len(page) < 100 { if newest != "" { if err := s.store.SetSyncState(ctx, channelHistoryCompleteScope(channel.ID), "1"); err != nil { @@ -426,6 +425,14 @@ func (s *Syncer) bootstrapChannelHistory(ctx context.Context, channel *discordgo } break } + nextBefore := page[len(page)-1].ID + if nextBefore == "" { + return messageCount, fmt.Errorf("channel %s message page missing id", channel.ID) + } + if nextBefore == before { + return messageCount, fmt.Errorf("channel %s message page cursor did not advance", channel.ID) + } + before = nextBefore } if newest != "" { if err := s.advanceChannelLatest(ctx, channel.ID, newest); err != nil { @@ -451,7 +458,6 @@ func (s *Syncer) syncForwardPages(ctx context.Context, channel *discordgo.Channe return messageCount, newest, err } progress.touch(channel, len(page)) - after = maxSnowflake(after, pageNewest) newest = maxSnowflake(newest, pageNewest) messageCount += len(page) if err := s.advanceChannelLatest(ctx, channel.ID, newest); err != nil { @@ -460,6 +466,14 @@ func (s *Syncer) syncForwardPages(ctx context.Context, channel *discordgo.Channe if len(page) < 100 { break } + nextAfter := maxSnowflake(after, pageNewest) + if nextAfter == "" { + return messageCount, newest, fmt.Errorf("channel %s message page missing id", channel.ID) + } + if nextAfter == after { + return messageCount, newest, fmt.Errorf("channel %s message page cursor did not advance", channel.ID) + } + after = nextAfter } return messageCount, newest, nil } @@ -507,8 +521,8 @@ func (s *Syncer) syncBackfillPages(ctx context.Context, channel *discordgo.Chann } break } - before = page[len(page)-1].ID - if err := s.store.SetSyncState(ctx, channelBackfillScope(channel.ID), before); err != nil { + nextBefore := page[len(page)-1].ID + if err := s.store.SetSyncState(ctx, channelBackfillScope(channel.ID), nextBefore); err != nil { return messageCount, newest, err } if len(page) < 100 { @@ -520,6 +534,13 @@ func (s *Syncer) syncBackfillPages(ctx context.Context, channel *discordgo.Chann if pageLimit > 0 && pages >= pageLimit { break } + if nextBefore == "" { + return messageCount, newest, fmt.Errorf("channel %s message page missing id", channel.ID) + } + if nextBefore == before { + return messageCount, newest, fmt.Errorf("channel %s message page cursor did not advance", channel.ID) + } + before = nextBefore } return messageCount, newest, nil } diff --git a/internal/syncer/message_sync_cursor_test.go b/internal/syncer/message_sync_cursor_test.go new file mode 100644 index 00000000..39b41364 --- /dev/null +++ b/internal/syncer/message_sync_cursor_test.go @@ -0,0 +1,162 @@ +package syncer + +import ( + "context" + "fmt" + "path/filepath" + "testing" + "time" + + "github.com/bwmarrin/discordgo" + "github.com/stretchr/testify/require" + + "github.com/openclaw/discrawl/internal/store" +) + +type repeatingMessagePageClient struct { + fakeClient + page []*discordgo.Message + requests int +} + +func (c *repeatingMessagePageClient) ChannelMessages(_ context.Context, _ string, _ int, _, _ string) ([]*discordgo.Message, error) { + c.requests++ + if c.requests > 5 { + return nil, fmt.Errorf("message pager called %d times without stopping", c.requests) + } + return c.page, nil +} + +func fullMessagePage(lastID string) []*discordgo.Message { + page := make([]*discordgo.Message, 100) + now := time.Now().UTC() + author := &discordgo.User{ID: "u1", Username: "user"} + for i := range 99 { + page[i] = &discordgo.Message{ + ID: fmt.Sprintf("%d", i+1), + GuildID: "g1", + ChannelID: "c1", + Content: "msg", + Timestamp: now, + Author: author, + } + } + page[99] = &discordgo.Message{ + ID: lastID, + GuildID: "g1", + ChannelID: "c1", + Content: "tail", + Timestamp: now, + Author: author, + } + return page +} + +func emptyIDMessagePage() []*discordgo.Message { + page := make([]*discordgo.Message, 100) + now := time.Now().UTC() + author := &discordgo.User{ID: "u1", Username: "user"} + for i := range page { + page[i] = &discordgo.Message{ + ID: "", + GuildID: "g1", + ChannelID: "c1", + Content: "msg", + Timestamp: now, + Author: author, + } + } + return page +} + +func TestMessagePagesErrorWhenCursorDoesNotAdvance(t *testing.T) { + t.Parallel() + + channel := &discordgo.Channel{ID: "c1", GuildID: "g1", Name: "general", Type: discordgo.ChannelTypeGuildText} + + tests := []struct { + name string + page []*discordgo.Message + run func(*Syncer, *discordgo.Channel) error + wantErr string + wantCalls int + }{ + { + name: "bootstrap repeats last id", + page: fullMessagePage("100"), + run: func(svc *Syncer, channel *discordgo.Channel) error { + _, err := svc.bootstrapChannelHistory(context.Background(), channel, false, time.Time{}, nil) + return err + }, + wantErr: "message page cursor did not advance", + wantCalls: 2, + }, + { + name: "bootstrap empty last id", + page: fullMessagePage(""), + run: func(svc *Syncer, channel *discordgo.Channel) error { + _, err := svc.bootstrapChannelHistory(context.Background(), channel, false, time.Time{}, nil) + return err + }, + wantErr: "message page missing id", + wantCalls: 1, + }, + { + name: "forward repeats newest id", + page: fullMessagePage("100"), + run: func(svc *Syncer, channel *discordgo.Channel) error { + _, _, err := svc.syncForwardPages(context.Background(), channel, "50", false, nil) + return err + }, + wantErr: "message page cursor did not advance", + wantCalls: 2, + }, + { + name: "forward empty newest id", + page: emptyIDMessagePage(), + run: func(svc *Syncer, channel *discordgo.Channel) error { + _, _, err := svc.syncForwardPages(context.Background(), channel, "100", false, nil) + return err + }, + wantErr: "message page cursor did not advance", + wantCalls: 1, + }, + { + name: "unlimited backfill repeats last id", + page: fullMessagePage("100"), + run: func(svc *Syncer, channel *discordgo.Channel) error { + _, _, err := svc.syncBackfillPages(context.Background(), channel, "", "", channel.Name, false, time.Time{}, 0, nil) + return err + }, + wantErr: "message page cursor did not advance", + wantCalls: 2, + }, + { + name: "unlimited backfill empty last id", + page: fullMessagePage(""), + run: func(svc *Syncer, channel *discordgo.Channel) error { + _, _, err := svc.syncBackfillPages(context.Background(), channel, "", "", channel.Name, false, time.Time{}, 0, nil) + return err + }, + wantErr: "message page missing id", + wantCalls: 1, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + ctx := context.Background() + s, err := store.Open(ctx, filepath.Join(t.TempDir(), "discrawl.db")) + require.NoError(t, err) + t.Cleanup(func() { _ = s.Close() }) + + client := &repeatingMessagePageClient{page: tt.page} + svc := New(client, s, nil) + err = tt.run(svc, channel) + require.ErrorContains(t, err, tt.wantErr) + require.Equal(t, tt.wantCalls, client.requests) + }) + } +} From afb8466d59578d0aad939e42602aa4bed64ced24 Mon Sep 17 00:00:00 2001 From: Sebastien Tardif Date: Sat, 29 Aug 2026 18:15:53 -0700 Subject: [PATCH 2/2] fix(syncer): use strconv.Itoa in cursor test IDs Signed-off-by: Sebastien Tardif --- internal/syncer/message_sync_cursor_test.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/internal/syncer/message_sync_cursor_test.go b/internal/syncer/message_sync_cursor_test.go index 39b41364..344851a1 100644 --- a/internal/syncer/message_sync_cursor_test.go +++ b/internal/syncer/message_sync_cursor_test.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "path/filepath" + "strconv" "testing" "time" @@ -33,7 +34,7 @@ func fullMessagePage(lastID string) []*discordgo.Message { author := &discordgo.User{ID: "u1", Username: "user"} for i := range 99 { page[i] = &discordgo.Message{ - ID: fmt.Sprintf("%d", i+1), + ID: strconv.Itoa(i + 1), GuildID: "g1", ChannelID: "c1", Content: "msg",