Skip to content

Commit 04a7cd2

Browse files
Fizzadarclaude
andcommitted
Only skip Reddit data that can never be bridged
Network, database and other transient failures retry the sync batch with the cursor retained, or fail the backfill page. Errors explicitly marked unbridgeable (invalid native data, 404/410, unsupported or oversized media) are skipped, or become a placeholder in history. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 429c6fb commit 04a7cd2

11 files changed

Lines changed: 136 additions & 72 deletions

File tree

‎pkg/connector/backfill.go‎

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -126,15 +126,19 @@ func (r *RedditClient) FetchMessages(ctx context.Context, params bridgev2.FetchM
126126
intent, _ = params.Portal.GetIntentFor(ctx, r.makeSender(evt.Sender), r.userLogin, bridgev2.RemoteEventMessage)
127127
}
128128
converted, err := r.convertMessage(ctx, params.Portal, intent, evt)
129-
if err != nil {
129+
if isUnbridgeable(err) {
130130
result.Messages = append(result.Messages, r.historyPlaceholder(ctx, evt, err))
131131
continue
132+
} else if err != nil {
133+
return nil, err
132134
}
133135
var reactions []*bridgev2.BackfillReaction
134136
if evt.Unsigned.RedactedBecause == nil {
135137
reactions, err = r.backfillReactions(ctx, params.Portal.PortalKey, makeMessageID(evt.ID))
136-
if err != nil {
137-
zerolog.Ctx(ctx).Warn().Err(err).Stringer("event_id", evt.ID).Msg("Failed to fetch Reddit history reactions, backfilling message without them")
138+
if isUnbridgeable(err) {
139+
zerolog.Ctx(ctx).Warn().Err(err).Stringer("event_id", evt.ID).Msg("Skipping unbridgeable Reddit history reactions")
140+
} else if err != nil {
141+
return nil, err
138142
}
139143
}
140144
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})
@@ -149,9 +153,6 @@ func (r *RedditClient) FetchMessages(ctx context.Context, params bridgev2.FetchM
149153
if (params.Forward || seekBackwardAnchor) && anchor != nil && !anchorFound {
150154
return nil, errors.New("reddit history ended before the known message; cannot locate the history anchor")
151155
}
152-
if err := ctx.Err(); err != nil {
153-
return nil, err
154-
}
155156
sort.SliceStable(result.Messages, func(i, j int) bool {
156157
left, right := result.Messages[i], result.Messages[j]
157158
if left.Timestamp.Equal(right.Timestamp) {

‎pkg/connector/chatsync.go‎

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

89-
r.handleSync(ctx, resp)
90-
if ctx.Err() != nil {
91-
return
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
9299
}
93100
// All framework event calls have returned. Advance the native token;
94101
// event delivery and crash safety belong to bridgev2 and its runtime.
@@ -111,21 +118,28 @@ func (r *RedditClient) runSync(ctx context.Context, done chan struct{}) {
111118
}
112119
}
113120

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)
121+
// Unbridgeable native data is skipped so it cannot hold the cursor forever.
122+
// Any other failure retries the whole batch.
123+
func (r *RedditClient) handleSync(ctx context.Context, resp *redditchat.SyncResponse) error {
124+
unhidden, err := r.applyHiddenChatData(ctx, resp)
125+
if err != nil {
126+
return err
127+
}
118128
for roomID, joined := range resp.Rooms.Join {
119129
if r.meta.HiddenRooms[roomID] {
120130
continue
121131
}
122-
logRoomError(ctx, roomID, r.handleRoomUpdate(ctx, roomID, joined))
132+
if err = skipUnbridgeable(ctx, r.handleRoomUpdate(ctx, roomID, joined), roomID, ""); err != nil {
133+
return err
134+
}
123135
}
124136
for roomID, invited := range resp.Rooms.Invite {
125137
if r.meta.HiddenRooms[roomID] {
126138
continue
127139
}
128-
logRoomError(ctx, roomID, r.handleInvitedRoom(ctx, roomID, invited))
140+
if err = skipUnbridgeable(ctx, r.handleInvitedRoom(ctx, roomID, invited), roomID, ""); err != nil {
141+
return err
142+
}
129143
}
130144
// Pending requests receive subsequent messages in rooms.peek. Process
131145
// invitation state first, then use the same queue/backfill path as joined
@@ -137,30 +151,34 @@ func (r *RedditClient) handleSync(ctx context.Context, resp *redditchat.SyncResp
137151
if len(peek.State.Events)+len(peek.Timeline.Events)+len(peek.Ephemeral.Events) == 0 {
138152
continue
139153
}
140-
logRoomError(ctx, roomID, r.handleRoomUpdate(ctx, roomID, peek))
154+
if err = skipUnbridgeable(ctx, r.handleRoomUpdate(ctx, roomID, peek), roomID, ""); err != nil {
155+
return err
156+
}
141157
}
142158
for roomID := range resp.Rooms.Leave {
143-
logRoomError(ctx, roomID, r.handleLeftRoom(ctx, roomID))
159+
if err = r.handleLeftRoom(ctx, roomID); err != nil {
160+
return err
161+
}
144162
}
145163
if len(unhidden) > 0 {
146164
state, err := r.fetchRoomSnapshots(ctx, unhidden)
147165
if err != nil {
148-
if ctx.Err() == nil {
149-
zerolog.Ctx(ctx).Err(err).Msg("Failed to fetch unhidden Reddit chats, skipping")
150-
}
151-
return
166+
return err
152167
}
153-
r.handleSync(ctx, state)
168+
return r.handleSync(ctx, state)
154169
}
170+
return nil
155171
}
156172

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")
173+
func skipUnbridgeable(ctx context.Context, err error, roomID id.RoomID, eventID id.EventID) error {
174+
if !isUnbridgeable(err) {
175+
return err
160176
}
177+
zerolog.Ctx(ctx).Warn().Err(err).Stringer("room_id", roomID).Stringer("event_id", eventID).Msg("Skipping unbridgeable Reddit data")
178+
return nil
161179
}
162180

163-
func (r *RedditClient) applyHiddenChatData(ctx context.Context, resp *redditchat.SyncResponse) []id.RoomID {
181+
func (r *RedditClient) applyHiddenChatData(ctx context.Context, resp *redditchat.SyncResponse) ([]id.RoomID, error) {
164182
var unhidden []id.RoomID
165183
// Apply account data before invitations: full native sync includes hidden
166184
// DMs in both peek and invite. Never mistake a preview for membership.
@@ -175,15 +193,17 @@ func (r *RedditClient) applyHiddenChatData(ctx context.Context, resp *redditchat
175193
}
176194
hidden, err := parseHiddenChat(evt)
177195
if err != nil {
178-
logRoomError(ctx, roomID, err)
196+
zerolog.Ctx(ctx).Warn().Err(err).Stringer("room_id", roomID).Stringer("event_id", evt.ID).Msg("Skipping unbridgeable Reddit hidden chat state")
179197
continue
180198
}
181199
if hidden {
182200
if r.meta.HiddenRooms == nil {
183201
r.meta.HiddenRooms = make(map[id.RoomID]bool)
184202
}
185203
r.meta.HiddenRooms[roomID] = true
186-
logRoomError(ctx, roomID, r.handleLeftRoom(ctx, roomID))
204+
if err = r.handleLeftRoom(ctx, roomID); err != nil {
205+
return nil, err
206+
}
187207
} else {
188208
delete(r.meta.HiddenRooms, roomID)
189209
if resp.Rooms.Join[roomID] == nil && resp.Rooms.Invite[roomID] == nil {
@@ -193,7 +213,7 @@ func (r *RedditClient) applyHiddenChatData(ctx context.Context, resp *redditchat
193213
}
194214
}
195215
}
196-
return unhidden
216+
return unhidden, nil
197217
}
198218

199219
func parseHiddenChat(evt *event.Event) (bool, error) {
@@ -215,7 +235,7 @@ func parseHiddenChat(evt *event.Event) (bool, error) {
215235

216236
func (r *RedditClient) handleRoomUpdate(ctx context.Context, roomID id.RoomID, joined *mautrix.SyncJoinedRoom) error {
217237
if joined == nil {
218-
return errors.New("reddit sync contains a null room")
238+
return unbridgeable(errors.New("reddit sync contains a null room"))
219239
}
220240
state, portalKey, err := r.resolveRoomState(ctx, roomID, joined)
221241
if err != nil {
@@ -226,8 +246,11 @@ func (r *RedditClient) handleRoomUpdate(ctx context.Context, roomID id.RoomID, j
226246
r.userLogin.QueueRemoteEvent(resync)
227247

228248
for _, evt := range joined.Timeline.Events {
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")
249+
if evt == nil {
250+
continue
251+
}
252+
if err = skipUnbridgeable(ctx, r.handleTimelineEvent(ctx, portalKey, evt), roomID, evt.ID); err != nil {
253+
return err
231254
}
232255
}
233256
r.handleEphemeralEvents(ctx, portalKey, joined.Ephemeral.Events)
@@ -237,7 +260,7 @@ func (r *RedditClient) handleRoomUpdate(ctx context.Context, roomID id.RoomID, j
237260
func (r *RedditClient) handleInvitedRoom(ctx context.Context, roomID id.RoomID, invited *mautrix.SyncInvitedRoom) error {
238261
joined, err := r.inviteSnapshot(invited)
239262
if err != nil {
240-
return err
263+
return unbridgeable(err)
241264
}
242265
state, key, err := r.roomState(ctx, roomID, joined)
243266
if err != nil {
@@ -267,15 +290,12 @@ func (r *RedditClient) handleLeftRoom(ctx context.Context, roomID id.RoomID) err
267290
}
268291

269292
func (r *RedditClient) handleTimelineEvent(ctx context.Context, portalKey networkid.PortalKey, evt *event.Event) error {
270-
if evt == nil {
271-
return nil
272-
}
273293
if evt.Type == event.EventMessage || evt.Type == event.EventReaction || evt.Type == event.EventRedaction {
274294
if evt.ID == "" || evt.Sender == "" || (evt.RoomID != "" && makePortalID(evt.RoomID) != portalKey.ID) {
275-
return errors.New("reddit timeline event has invalid room or event identity")
295+
return unbridgeable(errors.New("reddit timeline event has invalid room or event identity"))
276296
}
277297
if err := evt.Content.ParseRaw(evt.Type); err != nil && !errors.Is(err, event.ErrContentAlreadyParsed) {
278-
return err
298+
return unbridgeable(err)
279299
}
280300
}
281301
switch evt.Type {
@@ -308,7 +328,7 @@ func (r *RedditClient) queueMessage(ctx context.Context, key networkid.PortalKey
308328
if evt.Unsigned.RedactedBecause != nil {
309329
removal, err := r.convertDeletedEvent(key, evt)
310330
if err != nil {
311-
return err
331+
return unbridgeable(err)
312332
}
313333
r.userLogin.QueueRemoteEvent(removal)
314334
return nil
@@ -339,7 +359,7 @@ func (r *RedditClient) queueReaction(ctx context.Context, key networkid.PortalKe
339359
// Old sync responses may retain the original reaction content.
340360
removal, err := r.convertDeletedEvent(key, evt)
341361
if err != nil {
342-
return err
362+
return unbridgeable(err)
343363
}
344364
r.userLogin.QueueRemoteEvent(removal)
345365
return nil

‎pkg/connector/emoji.go‎

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ package connector
33
import (
44
"context"
55
"errors"
6+
"fmt"
67
"net/http"
78
"strings"
89
"time"
@@ -69,8 +70,10 @@ func (r *RedditClient) reactionImage(ctx context.Context, key string) (string, s
6970
return "", "", mediaError(ctx, "Reddit reaction download", err)
7071
}
7172
defer resp.Body.Close()
72-
if resp.StatusCode != http.StatusOK {
73-
return "", "", errors.New("reddit reaction image is unavailable")
73+
if resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusGone {
74+
return "", "", unbridgeable(errors.New("reddit reaction image no longer exists"))
75+
} else if resp.StatusCode != http.StatusOK {
76+
return "", "", fmt.Errorf("reddit reaction image is unavailable (HTTP %d)", resp.StatusCode)
7477
}
7578
data, err := readMediaBounded(resp.Body, 2*1024*1024)
7679
if err != nil {
@@ -96,10 +99,10 @@ func (r *RedditClient) reactionImage(ctx context.Context, key string) (string, s
9699
return mxc, ":" + asset.Name + ":", nil
97100
}
98101

99-
// Bridged reactions must not depend on the image CDN: fall back to text.
102+
// A reaction image that can never be fetched falls back to text.
100103
func (r *RedditClient) reactionEmoji(ctx context.Context, key string) (string, string, error) {
101104
emoji, shortcode, err := r.reactionImage(ctx, key)
102-
if err == nil || ctx.Err() != nil {
105+
if !isUnbridgeable(err) {
103106
return emoji, shortcode, err
104107
}
105108
asset, _ := redditchat.ReactionAssetByKey(key)

‎pkg/connector/media.go‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -71,13 +71,16 @@ func (r *RedditClient) downloadRedditMedia(ctx context.Context, uri id.ContentUR
7171
}
7272
parsed, err := uri.Parse()
7373
if err != nil || parsed.IsEmpty() || parsed.Homeserver != "reddit.com" {
74-
return nil, errors.New("attachment is not a Reddit media URI")
74+
return nil, unbridgeable(errors.New("attachment is not a Reddit media URI"))
7575
}
7676
resp, err := r.remote().DownloadMedia(ctx, parsed)
7777
if err != nil {
7878
if resp != nil && resp.Body != nil {
7979
resp.Body.Close()
8080
}
81+
if isNotFound(err) {
82+
return nil, unbridgeable(mediaError(ctx, "Reddit media download", err))
83+
}
8184
return nil, mediaError(ctx, "Reddit media download", err)
8285
}
8386
defer resp.Body.Close()
@@ -90,7 +93,7 @@ func (r *RedditClient) downloadRedditMedia(ctx context.Context, uri id.ContentUR
9093
}
9194
if file != nil {
9295
if err = file.DecryptInPlace(data); err != nil {
93-
return nil, errors.New("reddit attachment failed decryption or integrity verification")
96+
return nil, unbridgeable(errors.New("reddit attachment failed decryption or integrity verification"))
9497
}
9598
}
9699
return data, nil

‎pkg/connector/reactionhistory.go‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -30,22 +30,22 @@ func (r *RedditClient) fetchReactionEvents(ctx context.Context, key networkid.Po
3030
}
3131
for _, evt := range page.Chunk {
3232
if evt == nil || evt.Type != event.EventReaction || evt.ID == "" || (evt.RoomID != "" && makePortalID(evt.RoomID) != key.ID) {
33-
return nil, errors.New("reaction history has an invalid event identity")
33+
return nil, unbridgeable(errors.New("reaction history has an invalid event identity"))
3434
}
3535
local, server, err := evt.Sender.Parse()
3636
if err != nil || server != "reddit.com" || !strings.HasPrefix(local, "t2_") {
37-
return nil, errors.New("reaction history has an invalid sender")
37+
return nil, unbridgeable(errors.New("reaction history has an invalid sender"))
3838
}
3939
if err = evt.Content.ParseRaw(event.EventReaction); err != nil && !errors.Is(err, event.ErrContentAlreadyParsed) {
40-
return nil, err
40+
return nil, unbridgeable(err)
4141
}
4242
relation := evt.Content.AsReaction().RelatesTo
4343
if relation.Type != event.RelAnnotation || makeMessageID(relation.EventID) != target || relation.Key == "" {
44-
return nil, errors.New("reaction history belongs to another message")
44+
return nil, unbridgeable(errors.New("reaction history belongs to another message"))
4545
}
4646
if redaction := evt.Unsigned.RedactedBecause; redaction != nil {
4747
if redaction.Type != event.EventRedaction || redaction.Redacts != evt.ID || (redaction.RoomID != "" && makePortalID(redaction.RoomID) != key.ID) {
48-
return nil, errors.New("reaction history has an invalid removal identity")
48+
return nil, unbridgeable(errors.New("reaction history has an invalid removal identity"))
4949
}
5050
}
5151
events = append(events, evt)
@@ -54,7 +54,7 @@ func (r *RedditClient) fetchReactionEvents(ctx context.Context, key networkid.Po
5454
return events, nil
5555
}
5656
if page.NextBatch == cursor || seen[page.NextBatch] {
57-
return nil, errors.New("reddit reaction history cursor did not advance")
57+
return nil, unbridgeable(errors.New("reddit reaction history cursor did not advance"))
5858
}
5959
seen[page.NextBatch] = true
6060
cursor = page.NextBatch
@@ -76,7 +76,7 @@ func (r *RedditClient) backfillReactions(ctx context.Context, key networkid.Port
7676
identity := string(evt.Sender) + "\x00" + key
7777
if previous := seen[identity]; previous != "" {
7878
if previous != evt.ID {
79-
return nil, errors.New("reaction history has conflicting active identities")
79+
return nil, unbridgeable(errors.New("reaction history has conflicting active identities"))
8080
}
8181
continue
8282
}

‎pkg/connector/reactions.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ import (
1414
func (r *RedditClient) convertReaction(ctx context.Context, key networkid.PortalKey, evt *event.Event) (*simplevent.Reaction, error) {
1515
content, ok := evt.Content.Parsed.(*event.ReactionEventContent)
1616
if !ok || content.RelatesTo.EventID == "" || content.RelatesTo.Type != event.RelAnnotation || content.RelatesTo.Key == "" {
17-
return nil, errors.New("reaction is missing its target")
17+
return nil, unbridgeable(errors.New("reaction is missing its target"))
1818
}
1919
emoji, shortcode, err := r.reactionEmoji(ctx, content.RelatesTo.Key)
2020
if err != nil {

‎pkg/connector/redactions.go‎

Lines changed: 10 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -68,26 +68,29 @@ func validateTombstone(key networkid.PortalKey, target *event.Event) error {
6868
// event endpoint; the framework owns the corresponding local lookup/removal.
6969
func (r *RedditClient) convertRedaction(ctx context.Context, key networkid.PortalKey, removal *event.Event) (bridgev2.RemoteEvent, error) {
7070
if err := validateNativeEvent(key, removal); err != nil {
71-
return nil, err
71+
return nil, unbridgeable(err)
7272
}
7373
targetID, err := redactionTarget(removal)
7474
if err != nil {
75-
return nil, err
75+
return nil, unbridgeable(err)
7676
}
7777
target, err := r.remote().GetEvent(ctx, portalIDToRoomID(key.ID), targetID)
78-
if err != nil {
78+
if isNotFound(err) {
79+
return nil, unbridgeable(fmt.Errorf("resolve deleted Reddit event: %w", err))
80+
} else if err != nil {
7981
return nil, fmt.Errorf("resolve deleted Reddit event: %w", err)
8082
}
8183
if target == nil || target.ID != targetID {
82-
return nil, errors.New("reddit returned another deletion target")
84+
return nil, unbridgeable(errors.New("reddit returned another deletion target"))
8385
}
8486
if err = validateTombstone(key, target); err != nil {
85-
return nil, err
87+
return nil, unbridgeable(err)
8688
}
8789
if proof := target.Unsigned.RedactedBecause; proof.ID != removal.ID || proof.Sender != removal.Sender || proof.Timestamp != removal.Timestamp {
88-
return nil, errors.New("reddit returned a different deletion proof")
90+
return nil, unbridgeable(errors.New("reddit returned a different deletion proof"))
8991
}
90-
return r.convertDeletedEvent(key, target)
92+
removed, err := r.convertDeletedEvent(key, target)
93+
return removed, unbridgeable(err)
9194
}
9295

9396
func (r *RedditClient) convertDeletedEvent(key networkid.PortalKey, target *event.Event) (bridgev2.RemoteEvent, error) {

‎pkg/connector/requests.go‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -35,8 +35,6 @@ func (r *RedditClient) fetchRoomSnapshots(ctx context.Context, rooms []id.RoomID
3535
state.Rooms.Join[room] = joined
3636
} else if invited := full.Rooms.Invite[room]; invited != nil {
3737
state.Rooms.Invite[room] = invited
38-
} else {
39-
return nil, errors.New("reddit did not return the requested conversation state")
4038
}
4139
if peek := full.Rooms.Peek[room]; peek != nil {
4240
state.Rooms.Peek[room] = peek

0 commit comments

Comments
 (0)