Skip to content

Commit 04a6c4e

Browse files
authored
Make sliding sync linearize processing of sync requests (#3442)
* Make sliding sync linearize processing of sync requests * Iterate * Iterate * Iterate * Iterate
1 parent 0329824 commit 04a6c4e

File tree

4 files changed

+28
-10
lines changed

4 files changed

+28
-10
lines changed

spec/integ/sliding-sync-sdk.spec.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ import { fail } from "assert";
2020

2121
import { SlidingSync, SlidingSyncEvent, MSC3575RoomData, SlidingSyncState, Extension } from "../../src/sliding-sync";
2222
import { TestClient } from "../TestClient";
23-
import { IRoomEvent, IStateEvent } from "../../src/sync-accumulator";
23+
import { IRoomEvent, IStateEvent } from "../../src";
2424
import {
2525
MatrixClient,
2626
MatrixEvent,
@@ -39,7 +39,7 @@ import {
3939
} from "../../src";
4040
import { SlidingSyncSdk } from "../../src/sliding-sync-sdk";
4141
import { SyncApiOptions, SyncState } from "../../src/sync";
42-
import { IStoredClientOpts } from "../../src/client";
42+
import { IStoredClientOpts } from "../../src";
4343
import { logger } from "../../src/logger";
4444
import { emitPromise } from "../test-utils/test-utils";
4545
import { defer } from "../../src/utils";

src/models/typed-event-emitter.ts

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,24 @@ export class TypedEventEmitter<
8989
return super.emit(event, ...args);
9090
}
9191

92+
/**
93+
* Similar to `emit` but calls all listeners within a `Promise.all` and returns the promise chain
94+
* @param event - The name of the event to emit
95+
* @param args - Arguments to pass to the listener
96+
* @returns `true` if the event had listeners, `false` otherwise.
97+
*/
98+
public async emitPromised<T extends Events>(
99+
event: T,
100+
...args: Parameters<SuperclassArguments[T]>
101+
): Promise<boolean>;
102+
public async emitPromised<T extends Events>(event: T, ...args: Parameters<Arguments[T]>): Promise<boolean>;
103+
public async emitPromised<T extends Events>(event: T, ...args: any[]): Promise<boolean> {
104+
const listeners = this.listeners(event);
105+
return Promise.allSettled(listeners.map((l) => l(...args))).then(() => {
106+
return listeners.length > 0;
107+
});
108+
}
109+
92110
/**
93111
* Returns the number of listeners listening to the event named `event`.
94112
*

src/sliding-sync-sdk.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -376,7 +376,7 @@ export class SlidingSyncSdk {
376376
});
377377
}
378378

379-
private onRoomData(roomId: string, roomData: MSC3575RoomData): void {
379+
private async onRoomData(roomId: string, roomData: MSC3575RoomData): Promise<void> {
380380
let room = this.client.store.getRoom(roomId);
381381
if (!room) {
382382
if (!roomData.initial) {
@@ -385,7 +385,7 @@ export class SlidingSyncSdk {
385385
}
386386
room = _createAndReEmitRoom(this.client, roomId, this.opts);
387387
}
388-
this.processRoomData(this.client, room, roomData);
388+
await this.processRoomData(this.client, room!, roomData);
389389
}
390390

391391
private onLifecycle(state: SlidingSyncState, resp: MSC3575SlidingSyncResponse | null, err?: Error): void {

src/sliding-sync.ts

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -326,7 +326,7 @@ export enum SlidingSyncEvent {
326326
}
327327

328328
export type SlidingSyncEventHandlerMap = {
329-
[SlidingSyncEvent.RoomData]: (roomId: string, roomData: MSC3575RoomData) => void;
329+
[SlidingSyncEvent.RoomData]: (roomId: string, roomData: MSC3575RoomData) => Promise<void> | void;
330330
[SlidingSyncEvent.Lifecycle]: (
331331
state: SlidingSyncState,
332332
resp: MSC3575SlidingSyncResponse | null,
@@ -567,14 +567,14 @@ export class SlidingSync extends TypedEventEmitter<SlidingSyncEvent, SlidingSync
567567
* @param roomId - The room which received some data.
568568
* @param roomData - The raw sliding sync response JSON.
569569
*/
570-
private invokeRoomDataListeners(roomId: string, roomData: MSC3575RoomData): void {
570+
private async invokeRoomDataListeners(roomId: string, roomData: MSC3575RoomData): Promise<void> {
571571
if (!roomData.required_state) {
572572
roomData.required_state = [];
573573
}
574574
if (!roomData.timeline) {
575575
roomData.timeline = [];
576576
}
577-
this.emit(SlidingSyncEvent.RoomData, roomId, roomData);
577+
await this.emitPromised(SlidingSyncEvent.RoomData, roomId, roomData);
578578
}
579579

580580
/**
@@ -923,9 +923,9 @@ export class SlidingSync extends TypedEventEmitter<SlidingSyncEvent, SlidingSync
923923
}
924924
this.onPreExtensionsResponse(resp.extensions);
925925

926-
Object.keys(resp.rooms).forEach((roomId) => {
927-
this.invokeRoomDataListeners(roomId, resp!.rooms[roomId]);
928-
});
926+
for (const roomId in resp.rooms) {
927+
await this.invokeRoomDataListeners(roomId, resp!.rooms[roomId]);
928+
}
929929

930930
const listKeysWithUpdates: Set<string> = new Set();
931931
if (!doNotUpdateList) {

0 commit comments

Comments
 (0)