Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 35 additions & 11 deletions ts/session/apis/snode_api/swarmPolling.ts
Original file line number Diff line number Diff line change
Expand Up @@ -126,11 +126,20 @@ export class SwarmPolling {
* lastHashes[snode_edkey][pubkey_polled][namespace_polled] = last_hash
*/
private readonly lastHashes: Record<string, Record<string, Record<number, string>>>;
/**
* 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<string, number>;
private hasStarted = false;

constructor() {
this.groupPolling = [];
this.lastHashes = {};
this.cursorResets = new Map();
}

public async start(waitForFirstPoll = false): Promise<void> {
Expand Down Expand Up @@ -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 => {
Expand All @@ -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....`
);
Expand Down Expand Up @@ -786,6 +791,7 @@ export class SwarmPolling {
namespace: namespaces[index],
hash: lastMessage.hash,
expiration: lastMessage.expiration,
resetsAtStart,
});
})
);
Expand Down Expand Up @@ -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<void> {
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,
Expand All @@ -906,6 +920,9 @@ export class SwarmPolling {
});
}

if (this.cursorResetCount(pubkey) !== resetsAtStart) {
return;
}
if (!this.lastHashes[edkey]) {
this.lastHashes[edkey] = {};
}
Expand All @@ -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<string> {
if (!this.lastHashes[nodeEdKey]?.[pubkey]?.[namespace]) {
const lastHash = await Data.getLastHashBySnode(pubkey, nodeEdKey, namespace);
Expand All @@ -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++) {
Expand Down
198 changes: 198 additions & 0 deletions ts/test/session/unit/swarm_polling/SwarmPolling_cursorReset_test.ts
Original file line number Diff line number Diff line change
@@ -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<void>(resolve => {
release = resolve;
});
retrieveStub.callsFake(async (...args: Parameters<typeof answerWithNewMessage>) => {
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<string>) => 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<void> | 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('');
});
});
Loading