Skip to content
Open
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
103 changes: 102 additions & 1 deletion apps/webapp/src/script/main/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
// Polyfill for "tsyringe" dependency injection

import type {WallClock} from '@enormora/wall-clock/wall-clock';
import {isNonEmptyArray} from '@sindresorhus/is';
import {Context} from '@wireapp/api-client/lib/auth';
import {ClientClassification, ClientType} from '@wireapp/api-client/lib/client/';
import {FEATURE_KEY, FEATURE_STATUS, FeatureList} from '@wireapp/api-client/lib/team';
Expand All @@ -31,6 +32,7 @@ import 'core-js/full/reflect';
import pWaitFor from 'p-wait-for';
import platform from 'platform';
import {pdfjs} from 'react-pdf';
import {task} from 'true-myth';
import {container} from 'tsyringe';

import {Runtime} from '@wireapp/commons';
Expand Down Expand Up @@ -98,6 +100,7 @@ import {DebugUtil} from 'Util/debugUtil';
import {Environment} from 'Util/environment';
import {type Translate} from 'Util/localizerUtil';
import {getLogger, Logger} from 'Util/logger';
import {matchQualifiedIds} from 'Util/qualifiedId';
import {durationFrom, formatCoarseDuration, TIME_IN_MILLIS} from 'Util/timeUtil';
import {AppInitializationStep, checkIndexedDb, InitializationEventLogger} from 'Util/util';

Expand All @@ -121,7 +124,7 @@ import {
startNewVersionPolling,
} from '../lifecycle/newVersionHandler';
import {scheduleApiVersionUpdate, updateApiVersion} from '../lifecycle/updateRemoteConfigs';
import {initialiseSelfAndTeamConversations, initMLSGroupConversations} from '../mls';
import {initialiseSelfAndTeamConversations, initMLSGroupConversations, recoverMLSConversationsInBatches} from '../mls';
import {joinConversationsAfterMigrationFinalisation} from '../mls/MLSMigration/migrationFinaliser';
import type {ApplicationObservability} from '../observability/applicationObservability';
import type {ApplicationStartupReport} from '../observability/applicationStartupReport';
Expand Down Expand Up @@ -190,6 +193,7 @@ export class App {
debug?: DebugUtil;
util?: {debug: DebugUtil};
private newVersionPollingCleanup: (() => void) | undefined;
private mlsConversationRecoveryCleanup: (() => void) | undefined;

static get CONFIG() {
return {
Expand Down Expand Up @@ -788,6 +792,12 @@ export class App {
// resume the notification queue now that we're fully initialized
this.core.resumeNotificationQueue();

this.initializeMLSConversationRecovery({
conversationRepository,
eventRepository,
fireAndForgetInvoker,
});

return selfUser;
} catch (error: unknown) {
return reportStartupFailure(error, {
Expand All @@ -808,6 +818,92 @@ export class App {
}
}

private initializeMLSConversationRecovery({
conversationRepository,
eventRepository,
fireAndForgetInvoker,
}: {
conversationRepository: ConversationRepository;
eventRepository: EventRepository;
fireAndForgetInvoker: FireAndForgetInvoker;
}): void {
const mlsService = this.core.service?.mls;
if (!mlsService) {
return;
}

let recoveryInProgress = false;
const isApplicationActive = () => document.visibilityState === 'visible';
const isNotificationSyncLive = () =>
eventRepository.notificationHandlingState() === NOTIFICATION_HANDLING_STATE.WEB_SOCKET;

const recoverConversations = async (): Promise<void> => {
// Atomic check-and-set to prevent concurrent recovery attempts
if (recoveryInProgress || !isApplicationActive() || !isNotificationSyncLive()) {
return;
}

recoveryInProgress = true;
const recoveryTask = await task.tryOrElse(
error => error,
async () => {
if (!(await mlsService.prepareMLSConversationRecovery(this.core.clientId))) {
return;
}

const allConversations = await conversationRepository.refreshConversationsForMLSRecovery();
const pendingIds = await mlsService.getPendingRecoveryConversationIds();

// Only process conversations that haven't been recovered yet
const conversations = isNonEmptyArray(pendingIds)
? allConversations.filter(conv => pendingIds.some(pending => matchQualifiedIds(pending, conv.qualifiedId)))
: allConversations;

const result = await recoverMLSConversationsInBatches({
conversations,
conversationRepository,
core: this.core,
isActive: isApplicationActive,
mlsService,
});

if (result.completed && isApplicationActive()) {
await mlsService.completeMLSConversationRecovery();
this.logger.info('Completed MLS conversation recovery', {
recoveredConversationCount: result.recoveredConversationCount,
});
}
},
);
recoveryInProgress = false;

if (recoveryTask.isErr) {
this.logger.error('Failed to run MLS conversation recovery', recoveryTask.error);
}
};

const triggerRecovery = () => fireAndForgetInvoker.fireAndForget(recoverConversations);
const handleVisibilityChange = () => {
if (isApplicationActive()) {
triggerRecovery();
}
};

mlsService.on(MLSServiceEvents.MLS_CONVERSATION_RECOVERY_REQUIRED, triggerRecovery);
window.addEventListener('focus', triggerRecovery);
window.addEventListener('online', triggerRecovery);
document.addEventListener('visibilitychange', handleVisibilityChange);

this.mlsConversationRecoveryCleanup = () => {
mlsService.off(MLSServiceEvents.MLS_CONVERSATION_RECOVERY_REQUIRED, triggerRecovery);
window.removeEventListener('focus', triggerRecovery);
window.removeEventListener('online', triggerRecovery);
document.removeEventListener('visibilitychange', handleVisibilityChange);
};

triggerRecovery();
}

private _appInitFailure(error: BaseError) {
const {message, type} = error;
let logMessage = `Could not initialize app version '${Environment.version(false)}'`;
Expand Down Expand Up @@ -925,6 +1021,11 @@ export class App {
this.newVersionPollingCleanup = undefined;
}

if (this.mlsConversationRecoveryCleanup !== undefined) {
this.mlsConversationRecoveryCleanup();
this.mlsConversationRecoveryCleanup = undefined;
}

if (selfUser.isActivatedAccount()) {
this.repository.storage.terminate('window.onunload');
} else {
Expand Down
93 changes: 93 additions & 0 deletions apps/webapp/src/script/mls/MLSConversations.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ import {
initMLSGroupConversations,
initialiseSelfAndTeamConversations,
readLocalMLSState,
recoverMLSConversationsInBatches,
} from './MLSConversations';

function createMLSConversation(type?: CONVERSATION_TYPE, epoch = 0): MLSConversation {
Expand Down Expand Up @@ -140,6 +141,98 @@ describe('MLSConversations', () => {
});
});

describe('recoverMLSConversationsInBatches', () => {
it('recovers active MLS and mixed conversations while skipping Proteus, past-member, and established groups', async () => {
const mlsGroup = createMLSConversation(CONVERSATION_TYPE.REGULAR, 1);
const mlsOneToOne = createMLSConversation(CONVERSATION_TYPE.ONE_TO_ONE, 1);
const mixedSelf = new Conversation(
randomUUID(),
'',
CONVERSATION_PROTOCOL.MIXED,
translateForTest,
) as MLSConversation;
mixedSelf.groupId = `groupid-${randomUUID()}`;
mixedSelf.type(CONVERSATION_TYPE.SELF);
const established = createMLSConversation(CONVERSATION_TYPE.REGULAR, 1);
const pastMember = createMLSConversation(CONVERSATION_TYPE.REGULAR, 1);
pastMember.status(ConversationStatus.PAST_MEMBER);
const proteus = new Conversation(randomUUID(), '', CONVERSATION_PROTOCOL.PROTEUS, translateForTest);

const conversationRepository = await testFactory.exposeConversationActors();
const repositoryCore = conversationRepository['core'];
jest
.spyOn(repositoryCore.service!.conversation, 'mlsGroupExistsLocally')
.mockImplementation(async groupId => groupId === established.groupId);
const recoverSpy = jest
.spyOn(conversationRepository, 'safeEnsureConversationExists')
.mockReturnValue(task.resolve(undefined));

const result = await recoverMLSConversationsInBatches({
conversations: [mlsGroup, mlsOneToOne, mixedSelf, established, pastMember, proteus],
conversationRepository,
core: repositoryCore,
isActive: () => true,
batchSize: 2,
});

expect(result).toEqual({completed: true, failedConversationCount: 0, recoveredConversationCount: 3});
expect(recoverSpy).toHaveBeenCalledTimes(3);
expect(recoverSpy).toHaveBeenCalledWith({
conversationId: mlsOneToOne.qualifiedId,
groupId: mlsOneToOne.groupId,
core: repositoryCore,
});
expect(recoverSpy).toHaveBeenCalledWith({
conversationId: mixedSelf.qualifiedId,
groupId: mixedSelf.groupId,
core: repositoryCore,
});
});

it('pauses before the next batch when the application becomes inactive', async () => {
const conversations = createMLSConversations(12, CONVERSATION_TYPE.REGULAR);
const conversationRepository = await testFactory.exposeConversationActors();
const repositoryCore = conversationRepository['core'];
jest.spyOn(repositoryCore.service!.conversation, 'mlsGroupExistsLocally').mockResolvedValue(false);
const recoverSpy = jest
.spyOn(conversationRepository, 'safeEnsureConversationExists')
.mockReturnValue(task.resolve(undefined));
const isActive = jest.fn().mockReturnValueOnce(true).mockReturnValue(false);

const result = await recoverMLSConversationsInBatches({
conversations,
conversationRepository,
core: repositoryCore,
isActive,
batchSize: 5,
});

expect(result).toEqual({completed: false, failedConversationCount: 0, recoveredConversationCount: 5});
expect(recoverSpy).toHaveBeenCalledTimes(5);
});

it('keeps recovery incomplete after an individual failure and continues auditing the batch', async () => {
const conversations = createMLSConversations(3, CONVERSATION_TYPE.REGULAR);
const conversationRepository = await testFactory.exposeConversationActors();
const repositoryCore = conversationRepository['core'];
jest.spyOn(repositoryCore.service!.conversation, 'mlsGroupExistsLocally').mockResolvedValue(false);
const recoverSpy = jest
.spyOn(conversationRepository, 'safeEnsureConversationExists')
.mockReturnValueOnce(task.reject(new Error('join failed')))
.mockReturnValue(task.resolve(undefined));

const result = await recoverMLSConversationsInBatches({
conversations,
conversationRepository,
core: repositoryCore,
isActive: () => true,
});

expect(result).toEqual({completed: false, failedConversationCount: 1, recoveredConversationCount: 2});
expect(recoverSpy).toHaveBeenCalledTimes(3);
});
});

it('schedules key renewal intervals for all already established mls groups', async () => {
const core = new Account();
const nbMLSConversations = 5 + Math.ceil(Math.random() * 10);
Expand Down
Loading
Loading