@@ -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
199219func parseHiddenChat (evt * event.Event ) (bool , error ) {
@@ -215,7 +235,7 @@ func parseHiddenChat(evt *event.Event) (bool, error) {
215235
216236func (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
237260func (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
269292func (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
0 commit comments