Skip to content

Commit 808703d

Browse files
authored
Merge pull request Expensify#85920 from callstack-internal/callstack-internal/szymonzalarski/fix-seqquential-queue-issues-after-revert
2 parents 90d893a + 4f7328a commit 808703d

5 files changed

Lines changed: 263 additions & 118 deletions

File tree

src/libs/API/index.ts

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,10 @@ addMiddleware(SaveResponseInOnyx);
5353
// FraudMonitoring - Tags the request with the appropriate Fraud Protection event.
5454
addMiddleware(FraudMonitoring);
5555

56-
let requestIndex = 0;
56+
// Use timestamp-based IDs to avoid collisions between browser tabs.
57+
// Each tab has its own JS context with its own counter, so a simple
58+
// incrementing number would collide across tabs.
59+
let requestIndex = Date.now();
5760

5861
/**
5962
* Prepare the request to be sent. Bind data together with request metadata and apply optimistic Onyx data.
@@ -122,13 +125,13 @@ function prepareRequest<TCommand extends ApiCommand, TKey extends OnyxKey>(
122125
/**
123126
* Process a prepared request according to its type.
124127
*/
125-
function processRequest<TKey extends OnyxKey>(request: OnyxRequest<TKey>, type: ApiRequestType): Promise<void | Response<TKey>> {
128+
async function processRequest<TKey extends OnyxKey>(request: OnyxRequest<TKey>, type: ApiRequestType): Promise<void | Response<TKey>> {
126129
Log.info('[API] Processing request', false, {command: request.command, type});
127130
// Write commands can be saved and retried, so push it to the SequentialQueue
128131
if (type === CONST.API_REQUEST_TYPE.WRITE) {
129132
Log.info('[API] Write command. Pushing to SequentialQueue', false, {command: request.command});
130-
pushToSequentialQueue(request);
131-
return Promise.resolve();
133+
await pushToSequentialQueue(request);
134+
return;
132135
}
133136

134137
// Read requests are processed right away, but don't return the response to the caller
@@ -164,6 +167,7 @@ function write<TCommand extends WriteCommand, TKey extends OnyxKey>(
164167
): Promise<void | Response<TKey>> {
165168
Log.info('[API] Called API write', false, {command, ...apiCommandParameters});
166169
const request = prepareRequest(command, CONST.API_REQUEST_TYPE.WRITE, apiCommandParameters, onyxData, conflictResolver);
170+
167171
return processRequest(request, CONST.API_REQUEST_TYPE.WRITE);
168172
}
169173

src/libs/Network/SequentialQueue.ts

Lines changed: 21 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import {
66
endRequestAndRemoveFromQueue as endPersistedRequestAndRemoveFromQueue,
77
getAll as getAllPersistedRequests,
88
getCommands,
9+
onCrossTabRequestsMerged as onPersistedRequestsCrossTabMerge,
910
onInitialization as onPersistedRequestsInitialization,
1011
processNextRequest as processNextPersistedRequest,
1112
rollbackOngoingRequest as rollbackOngoingPersistedRequest,
@@ -436,7 +437,10 @@ onReconnection(flush);
436437
// Flush the queue when the persisted requests are initialized
437438
onPersistedRequestsInitialization(flush);
438439

439-
function handleConflictActions<TKey extends OnyxKey>(conflictAction: ConflictData, newRequest: OnyxRequest<TKey>) {
440+
// Flush the queue when another tab enqueues new requests
441+
onPersistedRequestsCrossTabMerge(flush);
442+
443+
async function handleConflictActions<TKey extends OnyxKey>(conflictAction: ConflictData, newRequest: OnyxRequest<TKey>): Promise<void> {
440444
Log.info('[SequentialQueue] handleConflictActions', false, {
441445
conflictType: conflictAction.type,
442446
newCommand: newRequest.command,
@@ -447,34 +451,34 @@ function handleConflictActions<TKey extends OnyxKey>(conflictAction: ConflictDat
447451
Log.info('[SequentialQueue] Conflict resolution: PUSH', false, {
448452
command: newRequest.command,
449453
});
450-
savePersistedRequest(newRequest);
454+
await savePersistedRequest(newRequest);
451455
} else if (conflictAction.type === 'replace') {
452456
Log.info('[SequentialQueue] Conflict resolution: REPLACE', false, {
453457
command: newRequest.command,
454458
replaceIndex: conflictAction.index,
455459
replacementRequest: conflictAction.request?.command ?? newRequest.command,
456460
});
457-
updatePersistedRequest(conflictAction.index, conflictAction.request ?? (newRequest as AnyRequest));
461+
await updatePersistedRequest(conflictAction.index, conflictAction.request ?? (newRequest as AnyRequest));
458462
} else if (conflictAction.type === 'delete') {
459463
Log.info('[SequentialQueue] Conflict resolution: DELETE', false, {
460464
command: newRequest.command,
461465
deleteIndices: conflictAction.indices,
462466
willPushNewRequest: conflictAction.pushNewRequest ?? false,
463467
hasNextAction: !!conflictAction.nextAction,
464468
});
465-
deletePersistedRequestsByIndices(conflictAction.indices);
469+
await deletePersistedRequestsByIndices(conflictAction.indices);
466470
if (conflictAction.pushNewRequest) {
467471
Log.info('[SequentialQueue] Pushing new request after delete', false, {
468472
command: newRequest.command,
469473
});
470-
savePersistedRequest(newRequest);
474+
await savePersistedRequest(newRequest);
471475
}
472476
if (conflictAction.nextAction) {
473477
Log.info('[SequentialQueue] Processing next conflict action', false, {
474478
command: newRequest.command,
475479
nextActionType: conflictAction.nextAction.type,
476480
});
477-
handleConflictActions(conflictAction.nextAction, newRequest);
481+
await handleConflictActions(conflictAction.nextAction, newRequest);
478482
}
479483
} else {
480484
Log.info('[SequentialQueue] No action performed, request ignored', false, {
@@ -484,7 +488,7 @@ function handleConflictActions<TKey extends OnyxKey>(conflictAction: ConflictDat
484488
}
485489
}
486490

487-
function push<TKey extends OnyxKey>(newRequest: OnyxRequest<TKey>) {
491+
function push<TKey extends OnyxKey>(newRequest: OnyxRequest<TKey>): Promise<void> {
488492
const currentRequests = getAllPersistedRequests();
489493
Log.info('[SequentialQueue] push() called', false, {
490494
command: newRequest.command,
@@ -494,6 +498,11 @@ function push<TKey extends OnyxKey>(newRequest: OnyxRequest<TKey>) {
494498
isSequentialQueueRunning,
495499
});
496500

501+
// Save the request to the persisted queue. The in-memory update inside save()
502+
// happens synchronously, so flush() below will see the new request immediately.
503+
// The returned promise resolves when disk persistence completes.
504+
let persistencePromise: Promise<void>;
505+
497506
if (newRequest.checkAndFixConflictingRequest) {
498507
const requests = currentRequests;
499508
Log.info('[SequentialQueue] Checking for conflicts', false, {
@@ -510,13 +519,13 @@ function push<TKey extends OnyxKey>(newRequest: OnyxRequest<TKey>) {
510519
// don't try to serialize a function.
511520
// eslint-disable-next-line no-param-reassign
512521
delete newRequest.checkAndFixConflictingRequest;
513-
handleConflictActions(conflictAction, newRequest);
522+
persistencePromise = handleConflictActions(conflictAction, newRequest);
514523
} else {
515524
Log.info('[SequentialQueue] No conflict action. Adding request to Persisted Requests', false, {
516525
command: newRequest.command,
517526
});
518527
// Add request to Persisted Requests so that it can be retried if it fails
519-
savePersistedRequest(newRequest);
528+
persistencePromise = savePersistedRequest(newRequest);
520529
}
521530

522531
// If we are offline we don't need to trigger the queue to empty as it will happen when we come back online
@@ -525,7 +534,7 @@ function push<TKey extends OnyxKey>(newRequest: OnyxRequest<TKey>) {
525534
command: newRequest.command,
526535
queueLength: getAllPersistedRequests().length,
527536
});
528-
return;
537+
return persistencePromise;
529538
}
530539

531540
// If the queue is running this request will run once it has finished processing the current batch
@@ -539,13 +548,14 @@ function push<TKey extends OnyxKey>(newRequest: OnyxRequest<TKey>) {
539548
});
540549
flush(true);
541550
});
542-
return;
551+
return persistencePromise;
543552
}
544553

545554
Log.info('[SequentialQueue] Queue is not running. Flushing the queue.', false, {
546555
command: newRequest.command,
547556
});
548557
flush(true);
558+
return persistencePromise;
549559
}
550560

551561
function getCurrentRequest(): Promise<void> {

0 commit comments

Comments
 (0)