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
463 changes: 463 additions & 0 deletions docs/SYNC-merge-vectors.json

Large diffs are not rendered by default.

8 changes: 6 additions & 2 deletions src/js/backupImport.js
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@
*/

import { Logger } from './logger.js';
import { mergeChannels } from './syncMerge.js';
import { mergeChannels, mergeSentDeletedAt, withoutDeleted } from './syncMerge.js';

/**
* Apply backup state to the unlocked secure storage and the live managers.
Expand Down Expand Up @@ -96,9 +96,13 @@ export async function importBackupData(data, { secureStorage, channelManager, id
}

// ---- Sent DM messages (per stream, only when absent) ----
if (data.sentDeletedAt) {
cache.sentDeletedAt = mergeSentDeletedAt(cache.sentDeletedAt, data.sentDeletedAt);
}
if (data.sentMessages) {
if (!cache.sentMessages) cache.sentMessages = {};
for (const [streamId, msgs] of Object.entries(data.sentMessages)) {
const restored = withoutDeleted(data.sentMessages, cache.sentDeletedAt);
for (const [streamId, msgs] of Object.entries(restored)) {
if (!cache.sentMessages[streamId]) {
cache.sentMessages[streamId] = msgs;
summary.dmHistories++;
Expand Down
24 changes: 16 additions & 8 deletions src/js/secureStorage.js
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,7 @@ class SecureStorage {
blockedPeers: [],
dmLeftAt: {},
channelsLeftAt: {},
sentDeletedAt: {},
pendingInvites: [],
sliceTs: {},
version: 2
Expand Down Expand Up @@ -1031,20 +1032,21 @@ class SecureStorage {
}

/**
* Remove a sent message from local storage (for deletes)
* Delete a sent message: drop it here and record the deletion, which the
* sync carries to the account's other devices.
* @param {string} streamId - Stream ID
* @param {string} messageId - Message ID to remove
*/
async removeSentMessage(streamId, messageId) {
if (!this.isUnlocked) return;
if (!this.cache.sentDeletedAt) this.cache.sentDeletedAt = {};
if (!this.cache.sentDeletedAt[streamId]) this.cache.sentDeletedAt[streamId] = {};
this.cache.sentDeletedAt[streamId][messageId] = Date.now();
const messages = this.cache.sentMessages?.[streamId];
if (!messages) return;
const idx = messages.findIndex(m => m.id === messageId);
if (idx >= 0) {
messages.splice(idx, 1);
await this.saveToStorage();
this.onSentDataChanged?.({ type: 'sentMessage', streamId });
}
const idx = messages ? messages.findIndex(m => m.id === messageId) : -1;
if (idx >= 0) messages.splice(idx, 1);
await this.saveToStorage();
this.onSentDataChanged?.({ type: 'sentMessage', streamId });
}

/**
Expand Down Expand Up @@ -1691,6 +1693,7 @@ class SecureStorage {

return {
sentMessages,
sentDeletedAt: this.cache.sentDeletedAt || {},
sentReactions: this.cache.sentReactions || {},
channels: this.cache.channels || [],
channelsLeftAt: this.cache.channelsLeftAt || {},
Expand Down Expand Up @@ -1736,6 +1739,10 @@ class SecureStorage {
changes.sentMessagesUpdated = true;
changes.hasChanges = true;
}
if (data.sentDeletedAt !== undefined && !isEqual(this.cache.sentDeletedAt, data.sentDeletedAt)) {
this.cache.sentDeletedAt = data.sentDeletedAt;
changes.hasChanges = true;
}
if (data.sentReactions !== undefined && !isEqual(this.cache.sentReactions, data.sentReactions)) {
this.cache.sentReactions = data.sentReactions;
changes.sentMessagesUpdated = true;
Expand Down Expand Up @@ -2036,6 +2043,7 @@ class SecureStorage {

return {
sentMessages: this.cache.sentMessages || {},
sentDeletedAt: this.cache.sentDeletedAt || {},
sentReactions: this.cache.sentReactions || {},
channels: this.cache.channels || [],
channelsLeftAt: this.cache.channelsLeftAt || {},
Expand Down
5 changes: 4 additions & 1 deletion src/js/syncManager.js
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ import { dmCrypto } from './dmCrypto.js';
import { channelManager } from './channels.js';
import { dmManager } from './dm.js';
import { identityManager } from './identity.js';
import { mergePayloadSeries as mergeSyncPayloadSeries, mergeSentMessages as mergeSyncSentMessages, mergeSentReactions as mergeSyncSentReactions, mergeState as mergeSyncState, mergeChannels as mergeSyncChannels, mergeEpochKeys as mergeSyncEpochKeys } from './syncMerge.js';
import { mergePayloadSeries as mergeSyncPayloadSeries, mergeSentMessages as mergeSyncSentMessages, mergeSentReactions as mergeSyncSentReactions, mergeState as mergeSyncState, mergeChannels as mergeSyncChannels, mergeEpochKeys as mergeSyncEpochKeys, mergeSentDeletedAt as mergeSyncSentDeletedAt, withoutDeleted as withoutDeletedSent } from './syncMerge.js';
import { syncWorkerClient } from './workers/syncWorkerClient.js';
import { cryptoManager } from './crypto.js';
import { splitSyncPayload, reassembleSyncPayloads } from './syncChunks.js';
Expand Down Expand Up @@ -829,6 +829,9 @@ class SyncManager {
// would be dropped by the slice replace, and for a paid gate the
// network never re-serves it.
merged.epochKeys = mergeSyncEpochKeys(liveState.epochKeys, merged.epochKeys);
// And a sent DM deleted while the merge ran would come back.
merged.sentDeletedAt = mergeSyncSentDeletedAt(liveState.sentDeletedAt, merged.sentDeletedAt);
merged.sentMessages = withoutDeletedSent(merged.sentMessages, merged.sentDeletedAt);

// Apply final merged state
const importStartedAt = getNow();
Expand Down
55 changes: 53 additions & 2 deletions src/js/syncMerge.js
Original file line number Diff line number Diff line change
@@ -1,6 +1,19 @@
import { CONFIG } from './config.js';

export function mergeSentMessages(local, remote, maxSentMessages = CONFIG.dm.maxSentMessages) {
const editedAt = (message) => (typeof message?._editedAt === 'number' ? message._editedAt : 0);

/**
* Union of sent messages by id. A copy edited later than the other replaces
* it, and any message the deletion map names is left out, whichever side
* still holds it.
*
* @param {Object} local - { streamId: [...messages] }
* @param {Object} remote - Same shape
* @param {number} [maxSentMessages] - Kept per conversation, newest last
* @param {Object} [deleted] - { streamId: { messageId: deletedAt } }
* @returns {Object}
*/
export function mergeSentMessages(local, remote, maxSentMessages = CONFIG.dm.maxSentMessages, deleted = {}) {
const result = {};

for (const [streamId, messages] of Object.entries(local || {})) {
Expand All @@ -18,6 +31,10 @@ export function mergeSentMessages(local, remote, maxSentMessages = CONFIG.dm.max
const existing = localById.get(message.id);
if (!existing) {
result[streamId].push({ ...message });
} else if (editedAt(message) > editedAt(existing)) {
const { imageData } = existing;
Object.assign(existing, message);
if (imageData && !message.imageData) existing.imageData = imageData;
} else if (message.type === 'image' && message.imageData && !existing.imageData) {
existing.imageData = message.imageData;
}
Expand All @@ -29,9 +46,39 @@ export function mergeSentMessages(local, remote, maxSentMessages = CONFIG.dm.max
}
}

return withoutDeleted(result, deleted);
}

/**
* Union of two sent-message deletion maps ({ streamId: { messageId:
* deletedAt } }), the latest time per message. A message id is never reused,
* so an entry is never retracted and never pruned.
*/
export function mergeSentDeletedAt(base, incoming) {
const result = {};
for (const source of [base || {}, incoming || {}]) {
for (const [streamId, ids] of Object.entries(source)) {
if (!ids || typeof ids !== 'object') continue;
const out = result[streamId] || (result[streamId] = {});
for (const [messageId, ts] of Object.entries(ids)) {
if (typeof ts !== 'number') continue;
if (!(messageId in out) || ts > out[messageId]) out[messageId] = ts;
}
}
}
return result;
}

/** The sent messages without every one the deletion map names. */
export function withoutDeleted(sentMessages, deleted) {
const out = {};
for (const [streamId, messages] of Object.entries(sentMessages || {})) {
const gone = deleted?.[streamId];
out[streamId] = gone ? messages.filter(message => !Object.hasOwn(gone, message.id)) : messages;
}
return out;
}

export function mergeSentReactions(local, remote) {
const result = {};
const allStreamIds = new Set([...Object.keys(local || {}), ...Object.keys(remote || {})]);
Expand Down Expand Up @@ -277,12 +324,16 @@ export function mergeState(base, incoming, maxSentMessages = CONFIG.dm.maxSentMe
graphApiKey = pickSlice('graphApiKey') || null;
}

const sentDeletedAt = mergeSentDeletedAt(base?.sentDeletedAt, incoming?.sentDeletedAt);

return {
sentMessages: mergeSentMessages(
base?.sentMessages || {},
incoming?.sentMessages || {},
maxSentMessages
maxSentMessages,
sentDeletedAt
),
sentDeletedAt,
sentReactions: mergeSentReactions(
base?.sentReactions || {},
incoming?.sentReactions || {}
Expand Down
15 changes: 15 additions & 0 deletions tests/unit/backupImport.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,21 @@ describe('importBackupData', () => {
expect(summary.dmHistories).toBe(2);
});

it('does not restore a sent DM deleted here or in the backup, and keeps both deletion records', async () => {
const { secureStorage, channelManager, cache } = makeDeps({
sentDeletedAt: { s2: { gone: 500 } }
});
const data = {
sentMessages: { s2: [{ id: 'gone' }, { id: 'kept' }, { id: 'also-gone' }] },
sentDeletedAt: { s2: { 'also-gone': 600 } }
};

await importBackupData(data, { secureStorage, channelManager });

expect(cache.sentMessages.s2.map(m => m.id)).toEqual(['kept']);
expect(cache.sentDeletedAt).toEqual({ s2: { gone: 500, 'also-gone': 600 } });
});

it('reloads the identity manager for the slices it caches in memory', async () => {
const { secureStorage, channelManager, identityManager } = makeDeps();
const data = {
Expand Down
51 changes: 51 additions & 0 deletions tests/unit/secureStorage.sentDeletedAt.test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
/**
* Deleting a sent DM has to reach the account's other devices: the local
* copy goes, and the deletion is recorded in the state the sync carries, so
* a device still holding the message drops it at the next merge.
*/

import { describe, it, expect, beforeEach, afterEach } from 'vitest';
import { secureStorage } from '../../src/js/secureStorage.js';

const DM = '0xpeer/Pombo-DM-1';
const text = (id) => ({ id, type: 'text', text: `text of ${id}`, timestamp: 100 });

describe('sent DM deletions', () => {
let saved;

beforeEach(() => {
saved = { cache: secureStorage.cache, isUnlocked: secureStorage.isUnlocked, isGuestMode: secureStorage.isGuestMode, address: secureStorage.address };
secureStorage.initAsGuest('0x1234567890abcdef1234567890abcdef12345678');
secureStorage.cache.sentMessages = { [DM]: [text('m1'), text('m2')] };
});

afterEach(() => Object.assign(secureStorage, saved));

it('drops the local copy and records the deletion', async () => {
await secureStorage.removeSentMessage(DM, 'm2');

expect(secureStorage.cache.sentMessages[DM].map(m => m.id)).toEqual(['m1']);
expect(typeof secureStorage.cache.sentDeletedAt[DM].m2).toBe('number');
});

it('records the deletion of a message this device never held', async () => {
await secureStorage.removeSentMessage(DM, 'elsewhere');

expect(secureStorage.cache.sentMessages[DM]).toHaveLength(2);
expect(secureStorage.cache.sentDeletedAt[DM]).toHaveProperty('elsewhere');
});

it('carries the deletions in the sync and backup exports', async () => {
await secureStorage.removeSentMessage(DM, 'm2');

expect(secureStorage.exportForSync().sentDeletedAt[DM]).toHaveProperty('m2');
expect(secureStorage.exportForBackup().sentDeletedAt[DM]).toHaveProperty('m2');
});

it('takes the merged deletions from a pull', async () => {
const changes = await secureStorage.importFromSync({ sentDeletedAt: { [DM]: { m9: 900 } } });

expect(changes.hasChanges).toBe(true);
expect(secureStorage.cache.sentDeletedAt).toEqual({ [DM]: { m9: 900 } });
});
});
39 changes: 39 additions & 0 deletions tests/unit/syncManager.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -635,6 +635,45 @@ describe('syncManager', () => {
expect(ids).toContain('ch-imported');
expect(ids).toContain('ch-remote');
});

it('should keep a sent DM deleted while the merge ran deleted', async () => {
authManager.wallet = { privateKey: '0x1234' };
authManager.getAddress.mockReturnValue('0xabc123');

const dm = '0xpeer/Pombo-DM-1';
const message = { id: 'm1', type: 'text', text: 'hi', timestamp: 100 };
const state = {
sentMessages: { [dm]: [message] },
sentReactions: {},
channels: [],
channelsLeftAt: {},
blockedPeers: [],
dmLeftAt: {},
trustedContacts: {},
ensCache: {},
username: null,
graphApiKey: null
};
secureStorage.exportForBackup
.mockReturnValueOnce(state)
.mockReturnValueOnce({ ...state, sentMessages: { [dm]: [] }, sentDeletedAt: { [dm]: { m1: 500 } } });

dmCrypto.decrypt.mockResolvedValue({
type: 'sync',
v: 1,
ts: 1000,
data: { sentMessages: { [dm]: [message] } }
});
streamrController.fetchPartitionHistory.mockResolvedValue([
{ content: { ct: 'enc' }, publisherId: '0xABC123', timestamp: 1000 }
]);

await syncManager.pullSync();

const imported = secureStorage.importFromSync.mock.calls[0][0];
expect(imported.sentMessages[dm]).toEqual([]);
expect(imported.sentDeletedAt).toEqual({ [dm]: { m1: 500 } });
});
});

describe('mergeState', () => {
Expand Down
8 changes: 8 additions & 0 deletions tests/unit/syncMerge.vectors.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -26,4 +26,12 @@ describe('sync merge parity vectors', () => {
expect(stampedSliceTs(v.state)).toEqual(v.sliceTs);
});
}

for (const v of vectors.sent) {
it(v.what, () => {
const merged = mergeState(v.base, v.incoming);
expect(merged.sentMessages).toEqual(v.expected.sentMessages);
expect(merged.sentDeletedAt).toEqual(v.expected.sentDeletedAt);
});
}
});
46 changes: 44 additions & 2 deletions tests/vectors/gen_sync_merge_vectors.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,8 @@
// every device.
//
// The vectors fix the slice outcome of one merge step (base = this device,
// incoming = a remote snapshot) and the stamping of unstamped values before a
// state leaves the device.
// incoming = a remote snapshot), the stamping of unstamped values before a
// state leaves the device, and how sent DMs carry deletions and edits.
import { mergeState, stampedSliceTs } from '../../src/js/syncMerge.js';

const SLICES = ['blockedPeers', 'dmLeftAt', 'trustedContacts', 'username', 'graphApiKey'];
Expand Down Expand Up @@ -37,6 +37,16 @@ const merge = (what, base, incoming) => ({
expected: slicesOf(mergeState(base, incoming))
});

const DM = '0x00000000000000000000000000000000000000a1/Pombo-DM-1';
const text = (id, timestamp, extra = {}) => ({ id, type: 'text', text: `text of ${id}`, timestamp, ...extra });

// Sent DMs: a deletion on any device removes the message on every device,
// and the latest edit wins.
const sent = (what, base, incoming) => {
const merged = mergeState(base, incoming);
return { what, base, incoming, expected: { sentMessages: merged.sentMessages, sentDeletedAt: merged.sentDeletedAt } };
};

console.log(JSON.stringify({
merge: [
merge('an unstamped empty snapshot never erases unstamped values', RESTORED, EMPTY),
Expand Down Expand Up @@ -64,5 +74,37 @@ console.log(JSON.stringify({
state: EMPTY,
sliceTs: stampedSliceTs(EMPTY)
}
],
sent: [
sent('a message deleted here stays deleted against a copy that still has it',
{ sentMessages: { [DM]: [text('m1', 1000)] }, sentDeletedAt: { [DM]: { m2: 5000 } } },
{ sentMessages: { [DM]: [text('m1', 1000), text('m2', 2000)] } }),
sent('a deletion made on another device removes the copy held here',
{ sentMessages: { [DM]: [text('m1', 1000), text('m2', 2000)] } },
{ sentMessages: { [DM]: [text('m1', 1000)] }, sentDeletedAt: { [DM]: { m2: 5000 } } }),
sent('deletions from both sides are joined, the latest time per message',
{ sentMessages: {}, sentDeletedAt: { [DM]: { m2: 5000, m3: 100 } } },
{ sentMessages: {}, sentDeletedAt: { [DM]: { m2: 6000, m4: 200 } } }),
sent('a snapshot from a client that knows no deletions keeps the ones held here',
{ sentMessages: { [DM]: [text('m1', 1000)] }, sentDeletedAt: { [DM]: { m2: 5000 } } },
{ sentMessages: { [DM]: [text('m1', 1000), text('m2', 2000)] }, sentDeletedAt: undefined }),
sent('the deletion of a message does not touch one sent later under a new id',
{ sentMessages: { [DM]: [text('m3', 7000)] }, sentDeletedAt: { [DM]: { m2: 5000 } } },
{ sentMessages: { [DM]: [text('m2', 2000), text('m3', 7000)] } }),
sent('the later edit wins, whichever side holds it',
{ sentMessages: { [DM]: [text('m1', 1000, { text: 'first', _edited: true, _editedAt: 3000 })] } },
{ sentMessages: { [DM]: [text('m1', 1000, { text: 'second', _edited: true, _editedAt: 4000 })] } }),
sent('an older edit arriving does not undo a newer one here',
{ sentMessages: { [DM]: [text('m1', 1000, { text: 'second', _edited: true, _editedAt: 4000 })] } },
{ sentMessages: { [DM]: [text('m1', 1000, { text: 'first', _edited: true, _editedAt: 3000 })] } }),
sent('an edit replaces the copy that was never edited',
{ sentMessages: { [DM]: [text('m1', 1000)] } },
{ sentMessages: { [DM]: [text('m1', 1000, { text: 'fixed', _edited: true, _editedAt: 3000 })] } }),
sent('a deletion wins over a later edit made on another device',
{ sentMessages: { [DM]: [] }, sentDeletedAt: { [DM]: { m1: 5000 } } },
{ sentMessages: { [DM]: [text('m1', 1000, { text: 'later', _edited: true, _editedAt: 9000 })] } }),
sent('a deletion from another device wins over a later edit held here',
{ sentMessages: { [DM]: [text('m1', 1000, { text: 'later', _edited: true, _editedAt: 9000 })] } },
{ sentMessages: { [DM]: [] }, sentDeletedAt: { [DM]: { m1: 5000 } } })
]
}, null, 2));
Loading