Skip to content

Commit 429c6fb

Browse files
Fizzadarclaude
andcommitted
Skip bad events in sync and backfill instead of aborting
handleSync aborted the whole batch on any room or event error and runSync retried it forever, so one deterministic failure stopped bridging for the entire account. Log and skip failures per room and per event, only holding the cursor when the sync is cancelled. FetchMessages likewise failed the whole page on one unconvertible message. Messages that fail to parse or convert now backfill as a placeholder notice mapped to their ID, edits are skipped, and a message whose reactions can't be fetched is backfilled without them. Cursor and identity problems still fail the page. Reaction images that fail to download fall back to the matching unicode emoji or shortcode text. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 5846ee3 commit 429c6fb

6 files changed

Lines changed: 103 additions & 62 deletions

File tree

‎pkg/connector/backfill.go‎

Lines changed: 23 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import (
77
"sort"
88
"time"
99

10+
"github.com/rs/zerolog"
1011
"maunium.net/go/mautrix"
1112
"maunium.net/go/mautrix/bridgev2"
1213
"maunium.net/go/mautrix/bridgev2/database"
@@ -111,10 +112,11 @@ func (r *RedditClient) FetchMessages(ctx context.Context, params bridgev2.FetchM
111112
return nil, errors.New("reddit history event has no stable identity")
112113
}
113114
if err = evt.Content.ParseRaw(event.EventMessage); err != nil && !errors.Is(err, event.ErrContentAlreadyParsed) {
114-
return nil, err
115+
result.Messages = append(result.Messages, r.historyPlaceholder(ctx, evt, err))
116+
continue
115117
}
116118
if evt.Content.AsMessage().RelatesTo.GetReplaceID() != "" {
117-
return nil, errors.New("reddit history edit reconciliation has not been implemented")
119+
continue
118120
}
119121
if params.ThreadRoot != "" && evt.Unsigned.RedactedBecause == nil && makeMessageID(evt.Content.AsMessage().RelatesTo.GetThreadParent()) != params.ThreadRoot {
120122
return nil, errors.New("reddit thread history contains a reply to another root")
@@ -125,13 +127,14 @@ func (r *RedditClient) FetchMessages(ctx context.Context, params bridgev2.FetchM
125127
}
126128
converted, err := r.convertMessage(ctx, params.Portal, intent, evt)
127129
if err != nil {
128-
return nil, err
130+
result.Messages = append(result.Messages, r.historyPlaceholder(ctx, evt, err))
131+
continue
129132
}
130133
var reactions []*bridgev2.BackfillReaction
131134
if evt.Unsigned.RedactedBecause == nil {
132135
reactions, err = r.backfillReactions(ctx, params.Portal.PortalKey, makeMessageID(evt.ID))
133136
if err != nil {
134-
return nil, err
137+
zerolog.Ctx(ctx).Warn().Err(err).Stringer("event_id", evt.ID).Msg("Failed to fetch Reddit history reactions, backfilling message without them")
135138
}
136139
}
137140
result.Messages = append(result.Messages, &bridgev2.BackfillMessage{ConvertedMessage: converted, ID: makeMessageID(evt.ID), Sender: r.makeSender(evt.Sender), Timestamp: time.UnixMilli(evt.Timestamp), Reactions: reactions})
@@ -146,6 +149,9 @@ func (r *RedditClient) FetchMessages(ctx context.Context, params bridgev2.FetchM
146149
if (params.Forward || seekBackwardAnchor) && anchor != nil && !anchorFound {
147150
return nil, errors.New("reddit history ended before the known message; cannot locate the history anchor")
148151
}
152+
if err := ctx.Err(); err != nil {
153+
return nil, err
154+
}
149155
sort.SliceStable(result.Messages, func(i, j int) bool {
150156
left, right := result.Messages[i], result.Messages[j]
151157
if left.Timestamp.Equal(right.Timestamp) {
@@ -155,3 +161,16 @@ func (r *RedditClient) FetchMessages(ctx context.Context, params bridgev2.FetchM
155161
})
156162
return result, nil
157163
}
164+
165+
// Map the placeholder to the native ID so replies and reactions to it resolve.
166+
func (r *RedditClient) historyPlaceholder(ctx context.Context, evt *event.Event, err error) *bridgev2.BackfillMessage {
167+
zerolog.Ctx(ctx).Warn().Err(err).Stringer("event_id", evt.ID).Msg("Failed to convert Reddit history message, using placeholder")
168+
return &bridgev2.BackfillMessage{
169+
ConvertedMessage: &bridgev2.ConvertedMessage{Parts: []*bridgev2.ConvertedMessagePart{{
170+
Type: event.EventMessage,
171+
Content: &event.MessageEventContent{MsgType: event.MsgNotice, Body: "An error occurred while processing an incoming message", Mentions: &event.Mentions{}},
172+
Extra: map[string]any{"fi.mau.bridge.internal_error": err.Error()},
173+
}}},
174+
ID: makeMessageID(evt.ID), Sender: r.makeSender(evt.Sender), Timestamp: time.UnixMilli(evt.Timestamp),
175+
}
176+
}

‎pkg/connector/chatsync.go‎

Lines changed: 55 additions & 52 deletions
Original file line numberDiff line numberDiff line change
@@ -86,16 +86,9 @@ func (r *RedditClient) runSync(ctx context.Context, done chan struct{}) {
8686
continue
8787
}
8888

89-
err = r.handleSync(ctx, resp)
90-
if err != nil {
91-
log.Error().Err(err).Msg("Failed to process Reddit sync; retaining cursor")
92-
r.userLogin.BridgeState.Send(status.BridgeState{StateEvent: status.StateTransientDisconnect, Error: "reddit-sync-apply-failed", Message: "Unable to synchronize Reddit chat state. Retrying."})
93-
select {
94-
case <-time.After(backoff):
95-
case <-ctx.Done():
96-
return
97-
}
98-
continue
89+
r.handleSync(ctx, resp)
90+
if ctx.Err() != nil {
91+
return
9992
}
10093
// All framework event calls have returned. Advance the native token;
10194
// event delivery and crash safety belong to bridgev2 and its runtime.
@@ -118,90 +111,79 @@ func (r *RedditClient) runSync(ctx context.Context, done chan struct{}) {
118111
}
119112
}
120113

121-
func (r *RedditClient) handleSync(ctx context.Context, resp *redditchat.SyncResponse) error {
122-
unhidden, err := r.applyHiddenChatData(ctx, resp)
123-
if err != nil {
124-
return err
125-
}
114+
// Failures are isolated per room and event: a deterministic failure must not
115+
// retain the cursor and stop every other chat on the account.
116+
func (r *RedditClient) handleSync(ctx context.Context, resp *redditchat.SyncResponse) {
117+
unhidden := r.applyHiddenChatData(ctx, resp)
126118
for roomID, joined := range resp.Rooms.Join {
127119
if r.meta.HiddenRooms[roomID] {
128120
continue
129121
}
130-
if err := r.handleRoomUpdate(ctx, roomID, joined); err != nil {
131-
return err
132-
}
122+
logRoomError(ctx, roomID, r.handleRoomUpdate(ctx, roomID, joined))
133123
}
134124
for roomID, invited := range resp.Rooms.Invite {
135125
if r.meta.HiddenRooms[roomID] {
136126
continue
137127
}
138-
if err := r.handleInvitedRoom(ctx, roomID, invited); err != nil {
139-
return err
140-
}
128+
logRoomError(ctx, roomID, r.handleInvitedRoom(ctx, roomID, invited))
141129
}
142130
// Pending requests receive subsequent messages in rooms.peek. Process
143131
// invitation state first, then use the same queue/backfill path as joined
144132
// timelines. Membership comes from state, never from the response bucket.
145133
for roomID, peek := range resp.Rooms.Peek {
146-
if r.meta.HiddenRooms[roomID] || resp.Rooms.Join[roomID] != nil || resp.Rooms.Leave[roomID] != nil {
134+
if peek == nil || r.meta.HiddenRooms[roomID] || resp.Rooms.Join[roomID] != nil || resp.Rooms.Leave[roomID] != nil {
147135
continue
148136
}
149137
if len(peek.State.Events)+len(peek.Timeline.Events)+len(peek.Ephemeral.Events) == 0 {
150138
continue
151139
}
152-
if err := r.handleRoomUpdate(ctx, roomID, peek); err != nil {
153-
return err
154-
}
140+
logRoomError(ctx, roomID, r.handleRoomUpdate(ctx, roomID, peek))
155141
}
156142
for roomID := range resp.Rooms.Leave {
157-
if err := r.handleLeftRoom(ctx, roomID); err != nil {
158-
return err
159-
}
143+
logRoomError(ctx, roomID, r.handleLeftRoom(ctx, roomID))
160144
}
161145
if len(unhidden) > 0 {
162146
state, err := r.fetchRoomSnapshots(ctx, unhidden)
163147
if err != nil {
164-
return err
148+
if ctx.Err() == nil {
149+
zerolog.Ctx(ctx).Err(err).Msg("Failed to fetch unhidden Reddit chats, skipping")
150+
}
151+
return
165152
}
166-
return r.handleSync(ctx, state)
153+
r.handleSync(ctx, state)
154+
}
155+
}
156+
157+
func logRoomError(ctx context.Context, roomID id.RoomID, err error) {
158+
if err != nil && ctx.Err() == nil {
159+
zerolog.Ctx(ctx).Err(err).Stringer("room_id", roomID).Msg("Failed to handle Reddit room update, skipping")
167160
}
168-
return nil
169161
}
170162

171-
func (r *RedditClient) applyHiddenChatData(ctx context.Context, resp *redditchat.SyncResponse) ([]id.RoomID, error) {
163+
func (r *RedditClient) applyHiddenChatData(ctx context.Context, resp *redditchat.SyncResponse) []id.RoomID {
172164
var unhidden []id.RoomID
173165
// Apply account data before invitations: full native sync includes hidden
174166
// DMs in both peek and invite. Never mistake a preview for membership.
175167
for _, rooms := range []map[id.RoomID]*mautrix.SyncJoinedRoom{resp.Rooms.Join, resp.Rooms.Peek} {
176168
for roomID, room := range rooms {
177169
if room == nil {
178-
return nil, errors.New("reddit sync contains a null room")
170+
continue
179171
}
180172
for _, evt := range room.AccountData.Events {
181173
if evt == nil || evt.Type.Type != "com.reddit.hidden_chat" {
182174
continue
183175
}
184-
var content struct {
185-
Hidden *bool `json:"hidden"`
186-
}
187-
raw, err := json.Marshal(&evt.Content)
176+
hidden, err := parseHiddenChat(evt)
188177
if err != nil {
189-
return nil, err
190-
}
191-
if err = json.Unmarshal(raw, &content); err != nil {
192-
return nil, err
193-
}
194-
if content.Hidden == nil {
195-
return nil, errors.New("reddit hidden chat state has no hidden flag")
178+
logRoomError(ctx, roomID, err)
179+
continue
196180
}
197-
if *content.Hidden {
181+
if hidden {
198182
if r.meta.HiddenRooms == nil {
199183
r.meta.HiddenRooms = make(map[id.RoomID]bool)
200184
}
201185
r.meta.HiddenRooms[roomID] = true
202-
if err = r.handleLeftRoom(ctx, roomID); err != nil {
203-
return nil, err
204-
}
186+
logRoomError(ctx, roomID, r.handleLeftRoom(ctx, roomID))
205187
} else {
206188
delete(r.meta.HiddenRooms, roomID)
207189
if resp.Rooms.Join[roomID] == nil && resp.Rooms.Invite[roomID] == nil {
@@ -211,10 +193,30 @@ func (r *RedditClient) applyHiddenChatData(ctx context.Context, resp *redditchat
211193
}
212194
}
213195
}
214-
return unhidden, nil
196+
return unhidden
197+
}
198+
199+
func parseHiddenChat(evt *event.Event) (bool, error) {
200+
var content struct {
201+
Hidden *bool `json:"hidden"`
202+
}
203+
raw, err := json.Marshal(&evt.Content)
204+
if err != nil {
205+
return false, err
206+
}
207+
if err = json.Unmarshal(raw, &content); err != nil {
208+
return false, err
209+
}
210+
if content.Hidden == nil {
211+
return false, errors.New("reddit hidden chat state has no hidden flag")
212+
}
213+
return *content.Hidden, nil
215214
}
216215

217216
func (r *RedditClient) handleRoomUpdate(ctx context.Context, roomID id.RoomID, joined *mautrix.SyncJoinedRoom) error {
217+
if joined == nil {
218+
return errors.New("reddit sync contains a null room")
219+
}
218220
state, portalKey, err := r.resolveRoomState(ctx, roomID, joined)
219221
if err != nil {
220222
return err
@@ -224,11 +226,12 @@ func (r *RedditClient) handleRoomUpdate(ctx context.Context, roomID id.RoomID, j
224226
r.userLogin.QueueRemoteEvent(resync)
225227

226228
for _, evt := range joined.Timeline.Events {
227-
if err = r.handleTimelineEvent(ctx, portalKey, evt); err != nil {
228-
return err
229+
if err = r.handleTimelineEvent(ctx, portalKey, evt); err != nil && ctx.Err() == nil {
230+
zerolog.Ctx(ctx).Err(err).Stringer("room_id", roomID).Stringer("event_id", evt.ID).Msg("Failed to handle Reddit event, skipping")
229231
}
230232
}
231-
return r.handleEphemeralEvents(ctx, portalKey, joined.Ephemeral.Events)
233+
r.handleEphemeralEvents(ctx, portalKey, joined.Ephemeral.Events)
234+
return nil
232235
}
233236

234237
func (r *RedditClient) handleInvitedRoom(ctx context.Context, roomID id.RoomID, invited *mautrix.SyncInvitedRoom) error {

‎pkg/connector/emoji.go‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import (
77
"strings"
88
"time"
99

10+
"github.com/rs/zerolog"
1011
"go.mau.fi/util/variationselector"
1112
"maunium.net/go/mautrix/bridgev2"
1213
"maunium.net/go/mautrix/bridgev2/database"
@@ -95,6 +96,22 @@ func (r *RedditClient) reactionImage(ctx context.Context, key string) (string, s
9596
return mxc, ":" + asset.Name + ":", nil
9697
}
9798

99+
// Bridged reactions must not depend on the image CDN: fall back to text.
100+
func (r *RedditClient) reactionEmoji(ctx context.Context, key string) (string, string, error) {
101+
emoji, shortcode, err := r.reactionImage(ctx, key)
102+
if err == nil || ctx.Err() != nil {
103+
return emoji, shortcode, err
104+
}
105+
asset, _ := redditchat.ReactionAssetByKey(key)
106+
zerolog.Ctx(ctx).Warn().Err(err).Str("reaction", asset.Name).Msg("Failed to fetch Reddit reaction image, using text")
107+
for unicode, name := range reactionAliases {
108+
if name == asset.Name {
109+
return unicode, "", nil
110+
}
111+
}
112+
return ":" + asset.Name + ":", "", nil
113+
}
114+
98115
func (r *RedditClient) resolveOutgoingReaction(ctx context.Context, value string) (string, error) {
99116
if strings.HasPrefix(value, "mxc://") {
100117
key := r.main.Bridge.DB.KV.Get(ctx, database.Key("reddit_emoji_v1/mxc/"+value))

‎pkg/connector/reactionhistory.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -81,7 +81,7 @@ func (r *RedditClient) backfillReactions(ctx context.Context, key networkid.Port
8181
continue
8282
}
8383
seen[identity] = evt.ID
84-
emoji, shortcode, err := r.reactionImage(ctx, key)
84+
emoji, shortcode, err := r.reactionEmoji(ctx, key)
8585
if err != nil {
8686
return nil, err
8787
}

‎pkg/connector/reactions.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ func (r *RedditClient) convertReaction(ctx context.Context, key networkid.Portal
1616
if !ok || content.RelatesTo.EventID == "" || content.RelatesTo.Type != event.RelAnnotation || content.RelatesTo.Key == "" {
1717
return nil, errors.New("reaction is missing its target")
1818
}
19-
emoji, shortcode, err := r.reactionImage(ctx, content.RelatesTo.Key)
19+
emoji, shortcode, err := r.reactionEmoji(ctx, content.RelatesTo.Key)
2020
if err != nil {
2121
return nil, err
2222
}

‎pkg/connector/receipts.go‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import (
77
"strings"
88
"time"
99

10+
"github.com/rs/zerolog"
1011
"maunium.net/go/mautrix/bridgev2"
1112
"maunium.net/go/mautrix/bridgev2/networkid"
1213
"maunium.net/go/mautrix/bridgev2/simplevent"
@@ -70,15 +71,16 @@ func (r *RedditClient) convertReadReceipts(key networkid.PortalKey, evt *event.E
7071
return out, nil
7172
}
7273

73-
func (r *RedditClient) handleEphemeralEvents(ctx context.Context, key networkid.PortalKey, events []*event.Event) error {
74+
func (r *RedditClient) handleEphemeralEvents(ctx context.Context, key networkid.PortalKey, events []*event.Event) {
7475
for _, evt := range events {
7576
if evt == nil {
7677
continue
7778
}
7879
if evt.Type == event.EphemeralEventTyping {
7980
typings, err := r.convertTyping(key, evt)
8081
if err != nil {
81-
return err
82+
zerolog.Ctx(ctx).Err(err).Stringer("portal_key", key).Msg("Failed to convert Reddit typing event, skipping")
83+
continue
8284
}
8385
for _, typing := range typings {
8486
r.userLogin.QueueRemoteEvent(typing)
@@ -90,7 +92,8 @@ func (r *RedditClient) handleEphemeralEvents(ctx context.Context, key networkid.
9092
}
9193
receipts, err := r.convertReadReceipts(key, evt)
9294
if err != nil {
93-
return err
95+
zerolog.Ctx(ctx).Err(err).Stringer("portal_key", key).Msg("Failed to convert Reddit read receipt, skipping")
96+
continue
9497
}
9598
for _, receipt := range receipts {
9699
if receipt.Timestamp.IsZero() {
@@ -99,5 +102,4 @@ func (r *RedditClient) handleEphemeralEvents(ctx context.Context, key networkid.
99102
r.userLogin.QueueRemoteEvent(receipt)
100103
}
101104
}
102-
return nil
103105
}

0 commit comments

Comments
 (0)