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
1,340 changes: 1,339 additions & 1 deletion docs/SYNC-merge-vectors.json

Large diffs are not rendered by default.

4 changes: 3 additions & 1 deletion src/js/channels.js
Original file line number Diff line number Diff line change
Expand Up @@ -518,14 +518,16 @@ class ChannelManager {
record.fieldTs = stamps;
this.channels.get(record.messageStreamId).fieldTs = stamps;
}
const changed = channelsData.length !== this._persisted.size || channelsData.some(record =>
JSON.stringify(record) !== JSON.stringify(this._persisted.get(record.messageStreamId)));
this._rememberPersisted(channelsData);
Logger.debug('Saving channels to secure storage:', channelsData.length);

await secureStorage.setChannels(channelsData);
Logger.debug('Channels saved to secure storage (metadata only)');

// Schedule auto-push to sync (debounced 30s)
this.onChannelsSaved?.();
if (changed) this.onChannelsSaved?.();
} catch (error) {
Logger.error('Failed to save channels:', error);
throw new StorageError(
Expand Down
2 changes: 1 addition & 1 deletion src/js/identity.js
Original file line number Diff line number Diff line change
Expand Up @@ -1053,7 +1053,7 @@ class IdentityManager {
*/
async removeTrustedContact(address) {
const normalizedAddress = address.toLowerCase();
this.trustedContacts.delete(normalizedAddress);
if (!this.trustedContacts.delete(normalizedAddress)) return;
await this.saveTrustedContacts();
this.onTrustedContactsChanged?.({ type: 'remove', address: normalizedAddress });
Logger.info('Removed trusted contact:', address);
Expand Down
11 changes: 8 additions & 3 deletions src/js/secureStorage.js
Original file line number Diff line number Diff line change
Expand Up @@ -624,8 +624,9 @@ class SecureStorage {
*/
async setTrustedContacts(contacts) {
if (!this.isUnlocked) return;
const changed = JSON.stringify(this.cache.trustedContacts) !== JSON.stringify(contacts);
this.cache.trustedContacts = contacts;
this._stampSliceTs('trustedContacts');
if (changed) this._stampSliceTs('trustedContacts');
await this.saveToStorage();
}

Expand Down Expand Up @@ -659,8 +660,8 @@ class SecureStorage {
*/
async setUsername(username) {
if (!this.isUnlocked) return;
if (this.cache.username !== username) this._stampSliceTs('username');
this.cache.username = username;
this._stampSliceTs('username');
// Also store in plain localStorage for pre-unlock display (unlock modal)
// Skip in guest mode — guest sessions are memory-only
if (this.address && !this.isGuestMode) {
Expand Down Expand Up @@ -751,8 +752,8 @@ class SecureStorage {
*/
async setGraphApiKey(apiKey) {
if (!this.isUnlocked) return;
if (this.cache.graphApiKey !== apiKey) this._stampSliceTs('graphApiKey');
this.cache.graphApiKey = apiKey;
this._stampSliceTs('graphApiKey');
await this.saveToStorage();
}

Expand Down Expand Up @@ -1026,6 +1027,7 @@ class SecureStorage {
if (!messages) return;
const msg = messages.find(m => m.id === messageId);
if (!msg) return;
if (Object.entries(fields).every(([key, value]) => JSON.stringify(msg[key]) === JSON.stringify(value))) return;
Object.assign(msg, fields);
await this.saveToStorage();
this.onSentDataChanged?.({ type: 'sentMessage', streamId });
Expand Down Expand Up @@ -1081,6 +1083,9 @@ class SecureStorage {
*/
async addSentReaction(streamId, messageId, emoji, user, action = 'add') {
if (!this.isUnlocked) return;
const held = (this.cache.sentReactions?.[streamId]?.[messageId]?.[emoji] || [])
.some(u => u.toLowerCase() === user.toLowerCase());
if (held === (action !== 'remove')) return;
if (!this.cache.sentReactions) {
this.cache.sentReactions = {};
}
Expand Down
14 changes: 11 additions & 3 deletions src/js/syncManager.js
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import { mergePayloadSeries as mergeSyncPayloadSeries, mergeSentMessages as merg
import { syncWorkerClient } from './workers/syncWorkerClient.js';
import { cryptoManager } from './crypto.js';
import { splitSyncPayload, reassembleSyncPayloads } from './syncChunks.js';
import { syncStateKey } from './syncStateKey.js';

/** A snapshot is a RUN of messages, so the window must hold several of them. */
const SYNC_FETCH_COUNT = 60;
Expand Down Expand Up @@ -95,6 +96,10 @@ class SyncManager {
this._ownRowKeys = new Set();
this._publishing = false;
this._confirmRun = 0;
// A pushed state counts as sent while its read-back runs. Only a
// confirmed one is remembered across restarts: a push storage never
// kept has to go out again.
this._pendingHash = null;
this._confirmTimer = null;
this._blobLeaveTimer = null;
this._snapshotWatch = null;
Expand Down Expand Up @@ -240,7 +245,7 @@ class SyncManager {
}

async _stateHash(state) {
const digest = await crypto.subtle.digest('SHA-256', new TextEncoder().encode(JSON.stringify(state)));
const digest = await crypto.subtle.digest('SHA-256', new TextEncoder().encode(syncStateKey(state)));
return Array.from(new Uint8Array(digest), (b) => b.toString(16).padStart(2, '0')).join('');
}

Expand Down Expand Up @@ -449,8 +454,8 @@ class SyncManager {
// Gather state from secureStorage
const state = secureStorage.exportForSync();
const hash = await this._stateHash(state);
if (this._isConfirmedState(hash)) {
Logger.info('Sync: State unchanged since the last confirmed push, not publishing');
if (this._isConfirmedState(hash) || hash === this._pendingHash) {
Logger.info('Sync: State unchanged since the last push, not publishing');
if (!this.autoPushTimeout && !this.pushQueued) {
this.clearDirty();
}
Expand Down Expand Up @@ -522,6 +527,7 @@ class SyncManager {
*/
_confirmPush(inboxStreamId, hash, rows) {
const run = ++this._confirmRun;
this._pendingHash = hash;
clearTimeout(this._confirmTimer);
const wanted = new Set(rows.map(rowKey));
const startedAt = Date.now();
Expand All @@ -540,6 +546,7 @@ class SyncManager {
Math.max(0, startedAt + SYNC_CONFIRM_AT_MS[index + 1] - Date.now()));
return;
}
this._pendingHash = null;
if (confirmed) {
this._writeConfirmedState(hash);
Logger.info('Sync: Push confirmed by storage');
Expand Down Expand Up @@ -608,6 +615,7 @@ class SyncManager {
/** On disconnect: a confirmation still running would act on the next account. */
cancelPushConfirmation() {
this._confirmRun++;
this._pendingHash = null;
clearTimeout(this._confirmTimer);
clearTimeout(this._blobLeaveTimer);
this._pulledRowTs = 0;
Expand Down
68 changes: 68 additions & 0 deletions src/js/syncStateKey.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
/**
* What a sync state carries that is not worth a push of its own: caches,
* bookkeeping the network or the next real change brings back, and stream ids
* derived from the channel's own id. Everything else is news, including a
* slice this client does not know. Android's SyncStateKey.kt keeps the same
* lists; the "publish" parity vectors in docs/SYNC-merge-vectors.json lock them.
*/
export const SYNC_STATE_IGNORED = Object.freeze({
slices: Object.freeze(['ensCache', 'sliceTs']),
epochKeyFields: Object.freeze([
'announces', 'pendingRequests', 'helloEpochs', 'helloName', 'helloTs',
'seenRequesters', 'pubAnnounce', 'intAnnounce'
]),
channelFields: Object.freeze([
'ephemeralStreamId', 'adminStreamId', 'keysStreamId', 'interactionsStreamId',
'inboxStreamId', 'storageProvider'
])
});

const isEmpty = (value) => value === null || value === undefined || value === '' || value === false
|| (Array.isArray(value) && value.length === 0)
|| (typeof value === 'object' && !Array.isArray(value) && Object.keys(value).length === 0);

/** Sorted keys, and an empty entry (null, '', false, [], {}) dropped: it is the same state as a missing one. */
function canonical(value) {
if (Array.isArray(value)) return value.map(canonical);
if (value && typeof value === 'object') {
const out = {};
for (const key of Object.keys(value).sort()) {
const entry = canonical(value[key]);
if (!isEmpty(entry)) out[key] = entry;
}
return out;
}
return value;
}

function without(record, fields) {
if (!record || typeof record !== 'object' || Array.isArray(record)) return record;
const out = { ...record };
for (const field of fields) delete out[field];
return out;
}

/**
* The news a sync state carries, as one canonical string: two states with the
* same key need no push between them. The order of the channel list is local
* and is not part of it.
* @param {Object} state - A sync export (exportForSync shape)
* @returns {string}
*/
export function syncStateKey(state) {
const projected = {};
for (const [slice, value] of Object.entries(state || {})) {
if (SYNC_STATE_IGNORED.slices.includes(slice)) continue;
if (slice === 'channels' && Array.isArray(value)) {
projected.channels = value
.map(channel => without(channel, SYNC_STATE_IGNORED.channelFields))
.sort((a, b) => String(a?.messageStreamId).localeCompare(String(b?.messageStreamId)));
} else if (slice === 'epochKeys' && value && typeof value === 'object') {
projected.epochKeys = Object.fromEntries(Object.entries(value)
.map(([streamId, entry]) => [streamId, without(entry, SYNC_STATE_IGNORED.epochKeyFields)]));
} else {
projected[slice] = value;
}
}
return JSON.stringify(canonical(projected));
}
19 changes: 19 additions & 0 deletions tests/unit/channels.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -1255,6 +1255,25 @@ describe('ChannelManager', () => {
expect(saved().fieldTs).toEqual({ name: 5000 });
});

it('asks the sync for a push only when a save changed a record', async () => {
channelManager.onChannelsSaved = vi.fn();
secureStorage.getChannels.mockReturnValue([stored()]);
channelManager.loadChannels();

await channelManager.saveChannels();
expect(channelManager.onChannelsSaved).not.toHaveBeenCalled();

channelManager.channels.get('stream1').name = 'Renamed';
await channelManager.saveChannels();
expect(channelManager.onChannelsSaved).toHaveBeenCalledTimes(1);

channelManager.channels.set('stream2', stored({ messageStreamId: 'stream2' }));
await channelManager.saveChannels();
channelManager.channels.delete('stream2');
await channelManager.saveChannels();
expect(channelManager.onChannelsSaved).toHaveBeenCalledTimes(3);
});

it('does not stamp a record created here', async () => {
secureStorage.getChannels.mockReturnValue([]);
channelManager.loadChannels();
Expand Down
9 changes: 7 additions & 2 deletions tests/unit/identity.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -302,9 +302,14 @@ describe('IdentityManager', () => {
expect(identityManager.trustedContacts.has('0xremove')).toBe(false);
});

it('should handle removing non-existent contact', async () => {
// Should not throw
it('should handle removing non-existent contact without saving or asking for a push', async () => {
identityManager.onTrustedContactsChanged = vi.fn();
vi.clearAllMocks();

await identityManager.removeTrustedContact('0xNonExistent');

expect(secureStorage.setTrustedContacts).not.toHaveBeenCalled();
expect(identityManager.onTrustedContactsChanged).not.toHaveBeenCalled();
});

it('should save contacts after removing', async () => {
Expand Down
53 changes: 53 additions & 0 deletions tests/unit/secureStorage.noopChanges.test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
/**
* A write that changes nothing must not ask the sync for a push, and must not
* stamp its slice newer: a fresh stamp on an unchanged value beats a real
* change another device made in the meantime.
*/

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

const DM = '0xpeer/Pombo-DM-1';

describe('writes that change nothing', () => {
let saved;

beforeEach(() => {
saved = { cache: secureStorage.cache, isUnlocked: secureStorage.isUnlocked, isGuestMode: secureStorage.isGuestMode, address: secureStorage.address, onSentDataChanged: secureStorage.onSentDataChanged };
secureStorage.initAsGuest('0x1234567890abcdef1234567890abcdef12345678');
secureStorage.cache.sentMessages = { [DM]: [{ id: 'm1', type: 'text', text: 'hello', timestamp: 100 }] };
secureStorage.onSentDataChanged = vi.fn();
});

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

it('an edit that sets the same text asks for no push', async () => {
await secureStorage.updateSentMessage(DM, 'm1', { text: 'hello' });
expect(secureStorage.onSentDataChanged).not.toHaveBeenCalled();

await secureStorage.updateSentMessage(DM, 'm1', { text: 'edited' });
expect(secureStorage.onSentDataChanged).toHaveBeenCalledTimes(1);
});

it('a reaction already held, or removed when absent, asks for no push and leaves no empty entry', async () => {
await secureStorage.addSentReaction(DM, 'm1', '👍', '0xAbc', 'remove');
expect(secureStorage.cache.sentReactions?.[DM]).toBeUndefined();

await secureStorage.addSentReaction(DM, 'm1', '👍', '0xAbc', 'add');
await secureStorage.addSentReaction(DM, 'm1', '👍', '0xabc', 'add');
expect(secureStorage.onSentDataChanged).toHaveBeenCalledTimes(1);
});

it('setting the same username or Graph key does not stamp the slice again', async () => {
await secureStorage.setUsername('Bob');
await secureStorage.setGraphApiKey('key');
await secureStorage.setTrustedContacts({ '0xc1': { nickname: 'Carol' } });
Object.assign(secureStorage.cache.sliceTs, { username: 1, graphApiKey: 1, trustedContacts: 1 });

await secureStorage.setUsername('Bob');
await secureStorage.setGraphApiKey('key');
await secureStorage.setTrustedContacts({ '0xc1': { nickname: 'Carol' } });

expect(secureStorage.cache.sliceTs).toMatchObject({ username: 1, graphApiKey: 1, trustedContacts: 1 });
});
});
10 changes: 10 additions & 0 deletions tests/unit/syncManager.extended.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,7 @@ describe('syncManager extended', () => {
syncManager.autoPushTimeout = null;
syncManager.pushQueued = false;
syncManager.autoPushRetryCount = 0;
syncManager.cancelPushConfirmation();
localStorage.clear();
authManager.wallet = { privateKey: '0x1234567890abcdef1234567890abcdef1234567890abcdef1234567890abcdef' };
authManager.isGuestMode.mockReturnValue(false);
Expand Down Expand Up @@ -750,6 +751,15 @@ describe('syncManager extended', () => {
expect(dmManager.sealAndPublish).not.toHaveBeenCalled();
});

it('does not publish the same state again while its read-back runs', async () => {
await syncManager.pushSync();
dmManager.sealAndPublish.mockClear();

await syncManager.pushSync();

expect(dmManager.sealAndPublish).not.toHaveBeenCalled();
});

it('leaves an unconfirmed push too, and sends that state again', async () => {
await syncManager.pushSync();
streamrController.fetchPartitionHistory.mockResolvedValue([]);
Expand Down
1 change: 1 addition & 0 deletions tests/unit/syncManager.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,7 @@ describe('syncManager', () => {
syncManager.handlers = [];
syncManager.pushQueued = false;
syncManager.autoPushRetryCount = 0;
syncManager.cancelPushConfirmation();
// Reset authManager state
authManager.wallet = { privateKey: '0x1234567890abcdef1234567890abcdef1234567890abcdef1234567890abcdef' };
authManager.isGuestMode.mockReturnValue(false);
Expand Down
22 changes: 22 additions & 0 deletions tests/unit/syncMerge.vectors.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,28 @@ import { readFileSync } from 'fs';
import { fileURLToPath } from 'url';
import { dirname, join } from 'path';
import { mergeState, stampedSliceTs } from '../../src/js/syncMerge.js';
import { syncStateKey } from '../../src/js/syncStateKey.js';

const vectors = JSON.parse(readFileSync(
join(dirname(fileURLToPath(import.meta.url)), '..', '..',
'docs', 'SYNC-merge-vectors.json'), 'utf8'));

const SLICES = ['blockedPeers', 'dmLeftAt', 'trustedContacts', 'username', 'graphApiKey'];

// The publish cases are patches on a shared base (see the generator).
const at = (root, path) => path.reduce((node, key) => node[key], root);
function applyPatch(state, patch = {}) {
const out = structuredClone(state);
for (const [path, value] of patch.set || []) at(out, path.slice(0, -1))[path.at(-1)] = structuredClone(value);
for (const path of patch.reverse || []) at(out, path).reverse();
for (const path of patch.reverseKeys || []) {
const reversed = Object.fromEntries(Object.entries(at(out, path)).reverse());
if (!path.length) return reversed;
at(out, path.slice(0, -1))[path.at(-1)] = reversed;
}
return out;
}

describe('sync merge parity vectors', () => {
for (const v of vectors.merge) {
it(v.what, () => {
Expand All @@ -35,6 +50,13 @@ describe('sync merge parity vectors', () => {
});
}

for (const v of vectors.publish.cases) {
it(v.what, () => {
const from = applyPatch(vectors.publish.base, v.basePatch);
expect(syncStateKey(from) === syncStateKey(applyPatch(from, v.patch))).toBe(v.same);
});
}

for (const v of vectors.sent) {
it(v.what, () => {
const merged = mergeState(v.base, v.incoming);
Expand Down
Loading
Loading