diff --git a/ts/session/apis/snode_api/swarmPolling.ts b/ts/session/apis/snode_api/swarmPolling.ts index 9212b2f57..f9f88899e 100644 --- a/ts/session/apis/snode_api/swarmPolling.ts +++ b/ts/session/apis/snode_api/swarmPolling.ts @@ -126,11 +126,20 @@ export class SwarmPolling { * lastHashes[snode_edkey][pubkey_polled][namespace_polled] = last_hash */ private readonly lastHashes: Record>>; + /** + * conversationId -> how many times its cursor has been reset. + * + * A poll compares this before and after its fetch rather than the cursor values themselves, + * because a value comparison cannot see a reset of a cursor that was already empty, and a poll + * that misses the reset writes its newest hash back and undoes it. + */ + private readonly cursorResets: Map; private hasStarted = false; constructor() { this.groupPolling = []; this.lastHashes = {}; + this.cursorResets = new Map(); } public async start(waitForFirstPoll = false): Promise { @@ -658,6 +667,8 @@ export class SwarmPolling { const snodeEdkey = node.pubkey_ed25519; try { + // taken before anything is read, so a reset at any point from here on is seen + const resetsAtStart = this.cursorResetCount(pubkey); const configHashesToBump = await this.getHashesToBump(type, pubkey); const namespacesAndLastHashes = await Promise.all( namespaces.map(async namespace => { @@ -680,17 +691,11 @@ export class SwarmPolling { allow401s ); - const namespacesAndLastHashesAfterFetch = await Promise.all( - namespaces.map(async namespace => { - const lastHash = await this.getLastHash(snodeEdkey, pubkey, namespace); - return { namespace, lastHash }; - }) - ); - - if ( - namespacesAndLastHashes.some(m => m) && - namespacesAndLastHashesAfterFetch.every(m => !m) - ) { + // The cursor was reset while this fetch was in flight, so it was made against a cursor that + // no longer exists. Writing its newest hash back would undo the reset and the history it + // asked for would never be fetched. Its messages are dropped with it: the next poll fetches + // them again from the start, and seen-message dedupe absorbs the overlap. + if (this.cursorResetCount(pubkey) !== resetsAtStart) { swarmLog( `SwarmPolling: hashes for ${ed25519Str(pubkey)} have been reset while we were fetching new messages. discarding them....` ); @@ -786,6 +791,7 @@ export class SwarmPolling { namespace: namespaces[index], hash: lastMessage.hash, expiration: lastMessage.expiration, + resetsAtStart, }); }) ); @@ -887,15 +893,23 @@ export class SwarmPolling { hash, namespace, pubkey, + resetsAtStart, }: { edkey: string; pubkey: string; namespace: number; hash: string; expiration: number; + /** the reset count when the poll writing this started; see cursorResets */ + resetsAtStart: number; }): Promise { const cached = await this.getLastHash(edkey, pubkey, namespace); + // Checked again before each write, not only once before the loop: a reset can land during the + // awaits in here, and either write landing after it undoes it. + if (this.cursorResetCount(pubkey) !== resetsAtStart) { + return; + } if (!cached || cached !== hash) { await Data.updateLastHash({ convoId: pubkey, @@ -906,6 +920,9 @@ export class SwarmPolling { }); } + if (this.cursorResetCount(pubkey) !== resetsAtStart) { + return; + } if (!this.lastHashes[edkey]) { this.lastHashes[edkey] = {}; } @@ -915,6 +932,10 @@ export class SwarmPolling { this.lastHashes[edkey][pubkey][namespace] = hash; } + private cursorResetCount(pubkey: string) { + return this.cursorResets.get(pubkey) ?? 0; + } + private async getLastHash(nodeEdKey: string, pubkey: string, namespace: number): Promise { if (!this.lastHashes[nodeEdKey]?.[pubkey]?.[namespace]) { const lastHash = await Data.getLastHashBySnode(pubkey, nodeEdKey, namespace); @@ -933,6 +954,9 @@ export class SwarmPolling { } public async resetLastHashesForConversation(conversationId: string) { + // Counted before the first await, so a poll checking at any point after this call starts sees + // it, including one whose cursor write would otherwise land between the clears below. + this.cursorResets.set(conversationId, this.cursorResetCount(conversationId) + 1); await Data.clearLastHashesForConvoId(conversationId); const snodeKeys = Object.keys(this.lastHashes); for (let index = 0; index < snodeKeys.length; index++) { diff --git a/ts/test/session/unit/swarm_polling/SwarmPolling_cursorReset_test.ts b/ts/test/session/unit/swarm_polling/SwarmPolling_cursorReset_test.ts new file mode 100644 index 000000000..a18644c97 --- /dev/null +++ b/ts/test/session/unit/swarm_polling/SwarmPolling_cursorReset_test.ts @@ -0,0 +1,198 @@ +import chai from 'chai'; +import { describe } from 'mocha'; +import Sinon from 'sinon'; + +import { getSwarmPollingInstance } from '../../../../session/apis/snode_api'; +import { SnodeAPIRetrieve } from '../../../../session/apis/snode_api/retrieveRequest'; +import { SwarmPolling } from '../../../../session/apis/snode_api/swarmPolling'; +import { SnodeNamespaces } from '../../../../session/apis/snode_api/namespaces'; +import { SnodePool } from '../../../../session/apis/snode_api/snodePool'; +import { PubKey } from '../../../../session/types'; +import { UserUtils } from '../../../../session/utils'; +import { UserSync } from '../../../../session/utils/job_runners/jobs/UserSyncJob'; +import { ConvoHub } from '../../../../session/conversations'; +import { ConversationTypeEnum } from '../../../../models/types'; +import { ReduxOnionSelectors } from '../../../../state/selectors/onions'; +import { TestUtils } from '../../../test-utils'; +import { generateFakeSnodes, stubData } from '../../../test-utils/utils'; + +const { expect } = chai; + +/** + * A cursor reset that happens while a poll of the same swarm is in flight must not be undone by + * that poll: once it completes, the next poll fetches from the beginning. + * + * Everything from pollOnceForKey down to the cursor write is production code. The retrieve is held + * open so the reset can land inside the window, and answers the way the real one does: one result + * per namespace asked about, in order. + */ +describe('SwarmPolling: a cursor reset during an in-flight poll', () => { + const ourNumber = TestUtils.generateFakePubKeyStr(); + const newestHash = 'newesthash'; + + let swarmPolling: SwarmPolling; + let retrieveStub: Sinon.SinonStub; + let updateLastHashStub: Sinon.SinonStub; + /** what the database holds as our cursor; the reset clears it */ + let storedCursor: string | undefined; + + function answerWithNewMessage( + _node: unknown, + _pubkey: unknown, + namespacesAndLastHashes: Array<{ namespace: SnodeNamespaces }> + ) { + return namespacesAndLastHashes.map(({ namespace }) => ({ + code: 200, + namespace, + messages: { + messages: + namespace === SnodeNamespaces.Default + ? [{ hash: newestHash, data: 'AQID', timestamp: 1, expiration: Date.now() + 60_000 }] + : [], + more: false, + t: 1, + }, + })) as any; + } + + /** a retrieve that does not answer until `release` is called */ + function holdRetrieveOpen() { + let release: () => void = () => {}; + const released = new Promise(resolve => { + release = resolve; + }); + retrieveStub.callsFake(async (...args: Parameters) => { + await released; + return answerWithNewMessage(...args); + }); + return () => release(); + } + + async function untilRetrieveCalls(count: number) { + while (retrieveStub.callCount < count) { + // eslint-disable-next-line no-await-in-loop + await new Promise(resolve => { + setImmediate(resolve); + }); + } + } + + /** the cursor the NEXT poll asks the Default namespace from */ + async function cursorTheNextPollUses() { + retrieveStub.resetHistory(); + retrieveStub.callsFake(answerWithNewMessage); + await swarmPolling.pollOnceForKey([ourNumber, ConversationTypeEnum.PRIVATE]); + const asked = retrieveStub.firstCall.args[2] as Array<{ namespace: number; lastHash: string }>; + return asked.find(n => n.namespace === SnodeNamespaces.Default)?.lastHash; + } + + beforeEach(async () => { + TestUtils.stubWindowFeatureFlags(); + TestUtils.stubWindowLog(); + ConvoHub.use().reset(); + Sinon.stub(UserSync, 'queueNewJobIfNeeded').resolves(); + Sinon.stub(UserUtils, 'getOurPubKeyStrFromCache').returns(ourNumber); + // read by pollOnceForKey once a poll returns messages + Sinon.stub(UserUtils, 'getUserED25519KeyPairBytes').resolves({ + pubKeyBytes: new Uint8Array(32), + privKeyBytes: new Uint8Array(64), + }); + TestUtils.stubLibSessionWorker(undefined); + TestUtils.stubUserGroupWrapper('getAllGroups', []); + TestUtils.stubUserGroupWrapper('getAllLegacyGroups', []); + + stubData('getAllConversations').resolves([]); + stubData('saveConversation').resolves(); + stubData('getSwarmNodesForPubkey').resolves(); + storedCursor = undefined; + stubData('getLastHashBySnode').callsFake(async () => storedCursor); + stubData('clearLastHashesForConvoId').callsFake(async () => { + storedCursor = undefined; + }); + updateLastHashStub = stubData('updateLastHash').callsFake(async ({ hash }: any) => { + storedCursor = hash; + }); + // everything fetched is already seen, so nothing past the cursor write is exercised + stubData('getSeenMessagesByHashList').callsFake(async (hashes: Array) => hashes); + + Sinon.stub(SnodePool, 'getSwarmFor').resolves(generateFakeSnodes(5)); + Sinon.stub(ReduxOnionSelectors, 'isOnlineOutsideRedux').returns(true); + TestUtils.stubWindow('inboxStore', undefined); + TestUtils.stubWindow('isOnline', true); + retrieveStub = Sinon.stub(SnodeAPIRetrieve, 'retrieveNextMessagesNoRetries'); + + await ConvoHub.use().load(); + ConvoHub.use().getOrCreate(PubKey.cast(ourNumber).key, ConversationTypeEnum.PRIVATE); + + swarmPolling = getSwarmPollingInstance(); + swarmPolling.resetSwarmPolling(); + // the cursor cache lives on the shared instance, so start every test without one + await swarmPolling.resetLastHashesForConversation(ourNumber); + }); + + afterEach(() => { + ConvoHub.use().reset(); + Sinon.restore(); + }); + + it('control: with no reset, the poll writes its newest hash as the cursor', async () => { + retrieveStub.callsFake(answerWithNewMessage); + + await swarmPolling.pollOnceForKey([ourNumber, ConversationTypeEnum.PRIVATE]); + + expect(updateLastHashStub.calledWithMatch({ hash: newestHash })).to.be.true; + expect(await cursorTheNextPollUses()).to.be.eq(newestHash); + }); + + it('a reset of a cursor that was already EMPTY is not undone', async () => { + // The case a before/after comparison of the cursor values cannot see: both read empty. + const release = holdRetrieveOpen(); + const poll = swarmPolling.pollOnceForKey([ourNumber, ConversationTypeEnum.PRIVATE]); + await untilRetrieveCalls(2); + + await swarmPolling.resetLastHashesForConversation(ourNumber); + release(); + await poll; + + expect(updateLastHashStub.called, 'the in-flight poll wrote no cursor').to.be.false; + expect(await cursorTheNextPollUses(), 'the next poll starts from the beginning').to.be.eq(''); + }); + + it('a reset of a cursor that was SET is not undone', async () => { + storedCursor = 'oldhash'; + const release = holdRetrieveOpen(); + const poll = swarmPolling.pollOnceForKey([ourNumber, ConversationTypeEnum.PRIVATE]); + await untilRetrieveCalls(2); + expect( + (retrieveStub.firstCall.args[2] as Array<{ lastHash: string }>).some( + n => n.lastHash === 'oldhash' + ), + 'PREMISE: the in-flight poll asked from the old cursor' + ).to.be.true; + + await swarmPolling.resetLastHashesForConversation(ourNumber); + release(); + await poll; + + expect(updateLastHashStub.called).to.be.false; + expect(await cursorTheNextPollUses()).to.be.eq(''); + }); + + it('a reset that lands while the cursor write is in progress is not undone', async () => { + retrieveStub.callsFake(answerWithNewMessage); + let reset: Promise | undefined; + updateLastHashStub.callsFake(async ({ hash }: any) => { + storedCursor = hash; + reset ??= swarmPolling.resetLastHashesForConversation(ourNumber); + await reset; + }); + + await swarmPolling.pollOnceForKey([ourNumber, ConversationTypeEnum.PRIVATE]); + + expect(reset, 'PREMISE: the reset happened inside the write').to.not.be.eq(undefined); + updateLastHashStub.callsFake(async ({ hash }: any) => { + storedCursor = hash; + }); + expect(await cursorTheNextPollUses(), 'the cached cursor was not written back').to.be.eq(''); + }); +});