From 3e6156ced775b4cf60cf0ffe515f2b84e058785e Mon Sep 17 00:00:00 2001 From: Fedor Indutny <79877362+indutny-signal@users.noreply.github.com> Date: Tue, 13 Sep 2022 09:16:01 -0700 Subject: [PATCH] Run checkForConflicts on a p-queue --- ts/ConversationController.ts | 325 ++++++++++++++++++----------------- 1 file changed, 169 insertions(+), 156 deletions(-) diff --git a/ts/ConversationController.ts b/ts/ConversationController.ts index fef7ebf48c..17cadc7225 100644 --- a/ts/ConversationController.ts +++ b/ts/ConversationController.ts @@ -664,7 +664,15 @@ export class ConversationController { return convoUuid; } - async checkForConflicts(): Promise { + checkForConflicts(): Promise { + return this._combineConversationsQueue.add(() => + this.doCheckForConflicts() + ); + } + + // Note: `doCombineConversations` is used within this function since both + // run on `_combineConversationsQueue` queue and we don't want deadlocks. + private async doCheckForConflicts(): Promise { log.info('checkForConflicts: starting...'); const byUuid = Object.create(null); const byE164 = Object.create(null); @@ -698,12 +706,12 @@ export class ConversationController { if (conversation.get('e164')) { // Keep new one // eslint-disable-next-line no-await-in-loop - await this.combineConversations(conversation, existing); + await this.doCombineConversations(conversation, existing); byUuid[uuid] = conversation; } else { // Keep existing - note that this applies if neither had an e164 // eslint-disable-next-line no-await-in-loop - await this.combineConversations(existing, conversation); + await this.doCombineConversations(existing, conversation); } } } @@ -719,12 +727,12 @@ export class ConversationController { if (conversation.get('e164') || conversation.get('pni')) { // Keep new one // eslint-disable-next-line no-await-in-loop - await this.combineConversations(conversation, existing); + await this.doCombineConversations(conversation, existing); byUuid[pni] = conversation; } else { // Keep existing - note that this applies if neither had an e164 // eslint-disable-next-line no-await-in-loop - await this.combineConversations(existing, conversation); + await this.doCombineConversations(existing, conversation); } } } @@ -759,12 +767,12 @@ export class ConversationController { if (conversation.get('uuid')) { // Keep new one // eslint-disable-next-line no-await-in-loop - await this.combineConversations(conversation, existing); + await this.doCombineConversations(conversation, existing); byE164[e164] = conversation; } else { // Keep existing - note that this applies if neither had a UUID // eslint-disable-next-line no-await-in-loop - await this.combineConversations(existing, conversation); + await this.doCombineConversations(existing, conversation); } } } @@ -799,11 +807,11 @@ export class ConversationController { !isGroupV2(existing.attributes) ) { // eslint-disable-next-line no-await-in-loop - await this.combineConversations(conversation, existing); + await this.doCombineConversations(conversation, existing); byGroupV2Id[groupV2Id] = conversation; } else { // eslint-disable-next-line no-await-in-loop - await this.combineConversations(existing, conversation); + await this.doCombineConversations(existing, conversation); } } } @@ -815,165 +823,170 @@ export class ConversationController { async combineConversations( current: ConversationModel, obsolete: ConversationModel + ): Promise { + return this._combineConversationsQueue.add(() => + this.doCombineConversations(current, obsolete) + ); + } + + private async doCombineConversations( + current: ConversationModel, + obsolete: ConversationModel ): Promise { const logId = `combineConversations/${obsolete.id}->${current.id}`; - return this._combineConversationsQueue.add(async () => { - const conversationType = current.get('type'); + const conversationType = current.get('type'); - if (!this.get(obsolete.id)) { - log.warn(`${logId}: Already combined obsolete conversation`); - } + if (!this.get(obsolete.id)) { + log.warn(`${logId}: Already combined obsolete conversation`); + } - if (obsolete.get('type') !== conversationType) { - assert( - false, - `${logId}: cannot combine a private and group conversation. Doing nothing` - ); - return; - } - - const dataToCopy: Partial = pick( - obsolete.attributes, - [ - 'conversationColor', - 'customColor', - 'customColorId', - 'draftAttachments', - 'draftBodyRanges', - 'draftTimestamp', - 'messageCount', - 'messageRequestResponseType', - 'quotedMessageId', - 'sentMessageCount', - ] + if (obsolete.get('type') !== conversationType) { + assert( + false, + `${logId}: cannot combine a private and group conversation. Doing nothing` ); + return; + } - const keys = Object.keys(dataToCopy) as Array< - keyof ConversationAttributesType - >; - keys.forEach(key => { - if (current.get(key) === undefined) { - current.set(key, dataToCopy[key]); + const dataToCopy: Partial = pick( + obsolete.attributes, + [ + 'conversationColor', + 'customColor', + 'customColorId', + 'draftAttachments', + 'draftBodyRanges', + 'draftTimestamp', + 'messageCount', + 'messageRequestResponseType', + 'quotedMessageId', + 'sentMessageCount', + ] + ); - // To ensure that any files on disk don't get deleted out from under us - if (key === 'draftAttachments') { - obsolete.set(key, undefined); - } - } - }); + const keys = Object.keys(dataToCopy) as Array< + keyof ConversationAttributesType + >; + keys.forEach(key => { + if (current.get(key) === undefined) { + current.set(key, dataToCopy[key]); - if (obsolete.get('isPinned')) { - obsolete.unpin(); - - if (!current.get('isPinned')) { - current.pin(); + // To ensure that any files on disk don't get deleted out from under us + if (key === 'draftAttachments') { + obsolete.set(key, undefined); } } - - const obsoleteId = obsolete.get('id'); - const obsoleteUuid = obsolete.getUuid(); - const currentId = current.get('id'); - log.warn( - `${logId}: Combining two conversations -`, - `old: ${obsolete.idForLogging()} -> new: ${current.idForLogging()}` - ); - - if (conversationType === 'private' && obsoleteUuid) { - if (!current.get('profileKey') && obsolete.get('profileKey')) { - log.warn(`${logId}: Copying profile key from old to new contact`); - - const profileKey = obsolete.get('profileKey'); - - if (profileKey) { - await current.setProfileKey(profileKey); - } - } - - log.warn(`${logId}: Delete all sessions tied to old conversationId`); - const ourACI = window.textsecure.storage.user.getUuid(UUIDKind.ACI); - const ourPNI = window.textsecure.storage.user.getUuid(UUIDKind.PNI); - await Promise.all( - [ourACI, ourPNI].map(async ourUuid => { - if (!ourUuid) { - return; - } - const deviceIds = - await window.textsecure.storage.protocol.getDeviceIds({ - ourUuid, - identifier: obsoleteUuid.toString(), - }); - await Promise.all( - deviceIds.map(async deviceId => { - const addr = new QualifiedAddress( - ourUuid, - new Address(obsoleteUuid, deviceId) - ); - await window.textsecure.storage.protocol.removeSession(addr); - }) - ); - }) - ); - - log.warn( - `${logId}: Delete all identity information tied to old conversationId` - ); - - if (obsoleteUuid) { - await window.textsecure.storage.protocol.removeIdentityKey( - obsoleteUuid - ); - } - - log.warn( - `${logId}: Ensure that all V1 groups have new conversationId instead of old` - ); - const groups = await this.getAllGroupsInvolvingUuid(obsoleteUuid); - groups.forEach(group => { - const members = group.get('members'); - const withoutObsolete = without(members, obsoleteId); - const currentAdded = uniq([...withoutObsolete, currentId]); - - group.set({ - members: currentAdded, - }); - updateConversation(group.attributes); - }); - } - - // Note: we explicitly don't want to update V2 groups - - log.warn(`${logId}: Delete the obsolete conversation from the database`); - await removeConversation(obsoleteId); - - log.warn(`${logId}: Update cached messages in MessageController`); - window.MessageController.update((message: MessageModel) => { - if (message.get('conversationId') === obsoleteId) { - message.set({ conversationId: currentId }); - } - }); - - log.warn(`${logId}: Update messages table`); - await migrateConversationMessages(obsoleteId, currentId); - - log.warn( - `${logId}: Emit refreshConversation event to close old/open new` - ); - window.Whisper.events.trigger('refreshConversation', { - newId: currentId, - oldId: obsoleteId, - }); - - log.warn( - `${logId}: Eliminate old conversation from ConversationController lookups` - ); - this._conversations.remove(obsolete); - this._conversations.resetLookups(); - - current.captureChange('combineConversations'); - - log.warn(`${logId}: Complete!`); }); + + if (obsolete.get('isPinned')) { + obsolete.unpin(); + + if (!current.get('isPinned')) { + current.pin(); + } + } + + const obsoleteId = obsolete.get('id'); + const obsoleteUuid = obsolete.getUuid(); + const currentId = current.get('id'); + log.warn( + `${logId}: Combining two conversations -`, + `old: ${obsolete.idForLogging()} -> new: ${current.idForLogging()}` + ); + + if (conversationType === 'private' && obsoleteUuid) { + if (!current.get('profileKey') && obsolete.get('profileKey')) { + log.warn(`${logId}: Copying profile key from old to new contact`); + + const profileKey = obsolete.get('profileKey'); + + if (profileKey) { + await current.setProfileKey(profileKey); + } + } + + log.warn(`${logId}: Delete all sessions tied to old conversationId`); + const ourACI = window.textsecure.storage.user.getUuid(UUIDKind.ACI); + const ourPNI = window.textsecure.storage.user.getUuid(UUIDKind.PNI); + await Promise.all( + [ourACI, ourPNI].map(async ourUuid => { + if (!ourUuid) { + return; + } + const deviceIds = + await window.textsecure.storage.protocol.getDeviceIds({ + ourUuid, + identifier: obsoleteUuid.toString(), + }); + await Promise.all( + deviceIds.map(async deviceId => { + const addr = new QualifiedAddress( + ourUuid, + new Address(obsoleteUuid, deviceId) + ); + await window.textsecure.storage.protocol.removeSession(addr); + }) + ); + }) + ); + + log.warn( + `${logId}: Delete all identity information tied to old conversationId` + ); + + if (obsoleteUuid) { + await window.textsecure.storage.protocol.removeIdentityKey( + obsoleteUuid + ); + } + + log.warn( + `${logId}: Ensure that all V1 groups have new conversationId instead of old` + ); + const groups = await this.getAllGroupsInvolvingUuid(obsoleteUuid); + groups.forEach(group => { + const members = group.get('members'); + const withoutObsolete = without(members, obsoleteId); + const currentAdded = uniq([...withoutObsolete, currentId]); + + group.set({ + members: currentAdded, + }); + updateConversation(group.attributes); + }); + } + + // Note: we explicitly don't want to update V2 groups + + log.warn(`${logId}: Delete the obsolete conversation from the database`); + await removeConversation(obsoleteId); + + log.warn(`${logId}: Update cached messages in MessageController`); + window.MessageController.update((message: MessageModel) => { + if (message.get('conversationId') === obsoleteId) { + message.set({ conversationId: currentId }); + } + }); + + log.warn(`${logId}: Update messages table`); + await migrateConversationMessages(obsoleteId, currentId); + + log.warn(`${logId}: Emit refreshConversation event to close old/open new`); + window.Whisper.events.trigger('refreshConversation', { + newId: currentId, + oldId: obsoleteId, + }); + + log.warn( + `${logId}: Eliminate old conversation from ConversationController lookups` + ); + this._conversations.remove(obsolete); + this._conversations.resetLookups(); + + current.captureChange('combineConversations'); + + log.warn(`${logId}: Complete!`); } /**