discrawl/internal/syncer/tail.go
2026-05-05 10:07:56 +01:00

123 lines
3.4 KiB
Go

package syncer
import (
"context"
"time"
"github.com/bwmarrin/discordgo"
"github.com/openclaw/discrawl/internal/store"
)
func (s *Syncer) RunTail(ctx context.Context, guildIDs []string, repairEvery time.Duration) error {
handler := &tailHandler{
guilds: makeGuildSet(guildIDs),
store: s.store,
client: s.client,
attachmentTextEnabled: s.attachmentTextEnabled,
}
if repairEvery <= 0 {
return s.client.Tail(ctx, handler)
}
errCh := make(chan error, 2)
go func() {
errCh <- s.client.Tail(ctx, handler)
}()
ticker := time.NewTicker(repairEvery)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return nil
case err := <-errCh:
return err
case <-ticker.C:
if _, err := s.Sync(ctx, SyncOptions{GuildIDs: guildIDs, Full: false, RepairReason: "tail_repair"}); err != nil {
s.logger.Warn("repair sync failed", "err", err)
}
}
}
}
type tailHandler struct {
guilds map[string]struct{}
store *store.Store
client Client
attachmentTextEnabled bool
}
func (t *tailHandler) OnMessageCreate(ctx context.Context, msg *discordgo.Message) error {
if !t.allowGuild(msg.GuildID) {
return nil
}
mutation, err := buildMessageMutation(ctx, msg, "", "", false, t.attachmentTextEnabled)
if err != nil {
return err
}
if err := t.store.UpsertMessages(ctx, []store.MessageMutation{mutation}); err != nil {
return err
}
if err := t.store.AppendMessageEvent(ctx, msg.GuildID, msg.ChannelID, msg.ID, "create", msg); err != nil {
return err
}
if err := t.store.SetSyncState(ctx, "tail:last_event", msg.ID); err != nil {
return err
}
return t.store.SetSyncState(ctx, channelLatestScope(msg.ChannelID), msg.ID)
}
func (t *tailHandler) OnMessageUpdate(ctx context.Context, msg *discordgo.Message) error {
if !t.allowGuild(msg.GuildID) {
return nil
}
mutation, err := buildMessageMutation(ctx, msg, "", "", false, t.attachmentTextEnabled)
if err != nil {
return err
}
if err := t.store.UpsertMessages(ctx, []store.MessageMutation{mutation}); err != nil {
return err
}
if err := t.store.AppendMessageEvent(ctx, msg.GuildID, msg.ChannelID, msg.ID, "update", msg); err != nil {
return err
}
return t.store.SetSyncState(ctx, "tail:last_event", msg.ID)
}
func (t *tailHandler) OnMessageDelete(ctx context.Context, evt *discordgo.MessageDelete) error {
if !t.allowGuild(evt.GuildID) {
return nil
}
if err := t.store.MarkMessageDeleted(ctx, evt.GuildID, evt.ChannelID, evt.ID, evt); err != nil {
return err
}
return t.store.SetSyncState(ctx, "tail:last_event", evt.ID)
}
func (t *tailHandler) OnChannelUpsert(ctx context.Context, channel *discordgo.Channel) error {
if !t.allowGuild(channel.GuildID) {
return nil
}
return t.store.UpsertChannel(ctx, toChannelRecord(channel, marshalJSONString(channel, "{}")))
}
func (t *tailHandler) OnMemberUpsert(ctx context.Context, guildID string, member *discordgo.Member) error {
if !t.allowGuild(guildID) || member == nil || member.User == nil {
return nil
}
return t.store.UpsertMember(ctx, toMemberRecord(guildID, member))
}
func (t *tailHandler) OnMemberDelete(ctx context.Context, guildID, userID string) error {
if !t.allowGuild(guildID) {
return nil
}
return t.store.DeleteMember(ctx, guildID, userID)
}
func (t *tailHandler) allowGuild(guildID string) bool {
if len(t.guilds) == 0 {
return true
}
_, ok := t.guilds[guildID]
return ok
}