diff --git a/src/js/channels.js b/src/js/channels.js index e6b823b..c1b4c91 100644 --- a/src/js/channels.js +++ b/src/js/channels.js @@ -1922,6 +1922,10 @@ class ChannelManager { this.messageFlow.resetPagingBackoff(channel); } + if (channel && channel.type !== 'dm') { + await this.messageFlow.restoreFailedOutbox(channel); + } + // Skip network subscription for write-only channels (no subscribe permission) // Instead, load locally persisted sent messages and reactions if (channel?.writeOnly) { @@ -2815,6 +2819,7 @@ class ChannelManager { try { await streamrController.unsubscribe(keysStreamId); } catch { /* not subscribed */ } await epochKeyManager.forgetChannel(messageStreamId); } + await secureStorage.clearFailedOutbox(messageStreamId); // Cancel any pending batch verifications for this channel this.cancelPendingVerifications(messageStreamId); diff --git a/src/js/channels/DeliveryConfirm.js b/src/js/channels/DeliveryConfirm.js index bfe273b..6728996 100644 --- a/src/js/channels/DeliveryConfirm.js +++ b/src/js/channels/DeliveryConfirm.js @@ -134,5 +134,6 @@ export class DeliveryConfirm { this.manager.notifyHandlers('message_failed', { streamId: messageStreamId, messageId: message.id, message, error: UNDELIVERED_REASON }); + this.manager.messageFlow?.keepForRetry(messageStreamId, message); } } diff --git a/src/js/channels/MessageFlow.js b/src/js/channels/MessageFlow.js index b634684..013a654 100644 --- a/src/js/channels/MessageFlow.js +++ b/src/js/channels/MessageFlow.js @@ -17,6 +17,7 @@ import { adminStatePoller } from '../adminStatePoller.js'; import { messageTime } from '../utils/messageTime.js'; import { storageFetch } from '../storageFetch.js'; import { NO_NETWORK, isOffline } from '../utils/network.js'; +import { dropLocalState } from '../publisherProof.js'; const PAGING_RETRY_DELAYS_MS = [5_000, 15_000, 30_000, 60_000]; @@ -261,6 +262,15 @@ export class MessageFlow { existing._timestamp = data._timestamp; if (Number.isFinite(data._seq)) existing._seq = data._seq; } + if (existing.failed) { + // A send marked failed reached the network after all. + existing.failed = false; + existing.pending = false; + delete existing.failError; + delete existing.undelivered; + secureStorage.removeFailedOutbox(streamId, existing.id).catch(() => {}); + this.manager.notifyHandlers('message_confirmed', { streamId, messageId: existing.id, message: existing }); + } Logger.debug('Message already exists, skipping duplicate:', data.id); return; } @@ -639,6 +649,46 @@ export class MessageFlow { } } + /** + * Keep a failed text send so its bubble and Retry come back after a + * restart. Never fails the caller: the send's own error is what surfaces. + * @param {string} messageStreamId + * @param {Object} message + */ + async keepForRetry(messageStreamId, message) { + try { + await secureStorage.putFailedOutbox(messageStreamId, { + ...dropLocalState({ ...message }), + failError: message.failError || null, + undelivered: !!message.undelivered, + failedAt: Date.now() + }); + } catch (error) { + Logger.warn('Could not keep the failed send for a retry:', error?.message); + } + } + + /** + * Put back the failed sends this timeline does not hold. Runs before + * history: an id history then delivers clears itself in handleTextMessage. + * @param {Object} channel + */ + async restoreFailedOutbox(channel) { + const shown = new Set(channel.messages.map(m => m.id)); + const entries = secureStorage.getFailedOutbox(channel.messageStreamId) + .filter(entry => entry.id && !shown.has(entry.id)); + if (!entries.length) return; + for (const entry of entries) { + channel.messages.push({ + ...entry, + failed: true, + pending: false, + verified: { valid: true, trustLevel: await identityManager.getTrustLevel(entry.sender) } + }); + } + this.manager.sortMessagesByTimestamp(channel); + } + /** * Publish a message already on the timeline and settle its send state. * @param {Object} channel @@ -658,12 +708,14 @@ export class MessageFlow { message.failed = true; message.failError = error.message; this.manager.notifyHandlers('message_failed', { streamId: messageStreamId, messageId: message.id, message, error: error.message }); + await this.keepForRetry(messageStreamId, message); throw error; } message.pending = false; message.failed = false; delete message.failError; + await secureStorage.removeFailedOutbox(messageStreamId, message.id).catch(() => {}); // For write-only channels, persist sent messages locally // (no subscribe permission = can't fetch history from network) diff --git a/src/js/channels/MessageOverrides.js b/src/js/channels/MessageOverrides.js index b9220be..6acb93a 100644 --- a/src/js/channels/MessageOverrides.js +++ b/src/js/channels/MessageOverrides.js @@ -297,6 +297,7 @@ export class MessageOverrides { const idx = channel.messages.indexOf(original); if (idx >= 0) channel.messages.splice(idx, 1); this.rememberDeleted(channel, targetId); + await secureStorage.removeFailedOutbox(streamId, targetId); this.manager.notifyHandlers('message_deleted', { streamId, targetId }); diff --git a/src/js/dm.js b/src/js/dm.js index 0de315a..de0842d 100644 --- a/src/js/dm.js +++ b/src/js/dm.js @@ -1384,6 +1384,7 @@ class DMManager { // Persist locally (we can't read from peer's inbox, so save our sent copy) await secureStorage.addSentMessage(peerInboxStreamId, message); + await secureStorage.removeFailedOutbox(peerInboxStreamId, message.id).catch(() => {}); channelManager.notifyHandlers('message_confirmed', { streamId: peerInboxStreamId, @@ -1408,6 +1409,7 @@ class DMManager { message, error: error.message }); + await channelManager.messageFlow.keepForRetry(peerInboxStreamId, message); throw error; } } @@ -1484,6 +1486,7 @@ class DMManager { const idx = channel.messages.indexOf(original); if (idx >= 0) channel.messages.splice(idx, 1); secureStorage.removeSentMessage(peerInboxStreamId, targetId); + secureStorage.removeFailedOutbox(peerInboxStreamId, targetId); channelManager.notifyHandlers('message_deleted', { streamId: peerInboxStreamId, targetId }); @@ -1526,9 +1529,14 @@ class DMManager { return !from || from === normalizedPeer; }); + // Sends that failed come last, so a sent or received copy of the same + // id wins: it reached the network after all. + const failedSends = secureStorage.getFailedOutbox(channelStreamId) + .map(entry => ({ ...entry, failed: true, pending: false, _dmSent: true })); + // Merge: combine sent + received, deduplicate by id, sort by timestamp // Prefer versions that have imageData (storage-backed messages may have it when sentMessages lost it) - const allMessages = [...sent, ...received]; + const allMessages = [...sent, ...received, ...failedSends]; const unique = new Map(); for (const msg of allMessages) { if (!msg.id) continue; diff --git a/src/js/secureStorage.js b/src/js/secureStorage.js index dae1306..7b92f11 100644 --- a/src/js/secureStorage.js +++ b/src/js/secureStorage.js @@ -15,6 +15,8 @@ import { StorageError } from './utils/errors.js'; import { cryptoWorkerPool } from './workers/cryptoWorkerPool.js'; import { stampedSliceTs } from './syncMerge.js'; +const FAILED_OUTBOX_MAX = 20; + class SecureStorage { constructor() { this.storageKey = null; // AES-256-GCM key @@ -1051,6 +1053,56 @@ class SecureStorage { this.onSentDataChanged?.({ type: 'sentMessage', streamId }); } + /** + * Failed text sends of a conversation, kept so the "Not sent" bubble and + * its Retry survive a restart. Device-local: neither exportForSync nor + * exportForBackup carries this slice, or every device of the account + * would offer a Retry for the same message. + * @param {string} streamId + * @returns {Array} + */ + getFailedOutbox(streamId) { + if (!this.isUnlocked) return []; + return this.cache.failedOutbox?.[streamId] || []; + } + + /** + * Add or replace the entry with the same id; past the cap the oldest go. + * @param {string} streamId + * @param {Object} entry - The message as published, plus failError/undelivered/failedAt + */ + async putFailedOutbox(streamId, entry) { + if (!this.isUnlocked || !entry?.id) return; + if (!this.cache.failedOutbox) this.cache.failedOutbox = {}; + const kept = (this.cache.failedOutbox[streamId] || []).filter(m => m.id !== entry.id); + kept.push(entry); + kept.sort((a, b) => (a.timestamp || 0) - (b.timestamp || 0)); + this.cache.failedOutbox[streamId] = kept.slice(-FAILED_OUTBOX_MAX); + await this.saveToStorage(); + } + + /** + * @param {string} streamId + * @param {string} messageId + */ + async removeFailedOutbox(streamId, messageId) { + const entries = this.isUnlocked ? this.cache.failedOutbox?.[streamId] : null; + if (!entries?.some(m => m.id === messageId)) return; + const kept = entries.filter(m => m.id !== messageId); + if (kept.length) this.cache.failedOutbox[streamId] = kept; + else delete this.cache.failedOutbox[streamId]; + await this.saveToStorage(); + } + + /** + * @param {string} streamId + */ + async clearFailedOutbox(streamId) { + if (!this.isUnlocked || !this.cache.failedOutbox?.[streamId]) return; + delete this.cache.failedOutbox[streamId]; + await this.saveToStorage(); + } + /** * Clear sent reactions for a channel * @param {string} streamId - Stream ID diff --git a/tests/unit/channels.extended.test.js b/tests/unit/channels.extended.test.js index 808ac17..0b3c48b 100644 --- a/tests/unit/channels.extended.test.js +++ b/tests/unit/channels.extended.test.js @@ -105,6 +105,10 @@ vi.mock('../../src/js/secureStorage.js', () => ({ addSentMessage: vi.fn().mockResolvedValue(undefined), getSentMessages: vi.fn().mockReturnValue([]), getSentReactions: vi.fn().mockReturnValue({}), + getFailedOutbox: vi.fn().mockReturnValue([]), + putFailedOutbox: vi.fn().mockResolvedValue(undefined), + removeFailedOutbox: vi.fn().mockResolvedValue(undefined), + clearFailedOutbox: vi.fn().mockResolvedValue(undefined), getEpochKeys: vi.fn().mockReturnValue(null), setEpochKeys: vi.fn().mockResolvedValue(undefined), clearEpochKeys: vi.fn().mockResolvedValue(undefined) diff --git a/tests/unit/channels.test.js b/tests/unit/channels.test.js index 79beb71..5efb1be 100644 --- a/tests/unit/channels.test.js +++ b/tests/unit/channels.test.js @@ -112,6 +112,10 @@ vi.mock('../../src/js/secureStorage.js', () => ({ addBlockedPeer: vi.fn().mockResolvedValue(undefined), addToChannelOrder: vi.fn().mockResolvedValue(undefined), addSentMessage: vi.fn().mockResolvedValue(undefined), + getFailedOutbox: vi.fn().mockReturnValue([]), + putFailedOutbox: vi.fn().mockResolvedValue(undefined), + removeFailedOutbox: vi.fn().mockResolvedValue(undefined), + clearFailedOutbox: vi.fn().mockResolvedValue(undefined), // Members-only creation adopts the publish key into the epoch slice getEpochKeys: vi.fn().mockReturnValue(null), setEpochKeys: vi.fn().mockResolvedValue(undefined) diff --git a/tests/unit/dm.test.js b/tests/unit/dm.test.js index 6be9b42..2345709 100644 --- a/tests/unit/dm.test.js +++ b/tests/unit/dm.test.js @@ -77,7 +77,9 @@ vi.mock('../../src/js/secureStorage.js', () => ({ getDMLeftAt: vi.fn().mockReturnValue(null), clearDMLeftAt: vi.fn().mockResolvedValue(undefined), updateSentMessage: vi.fn().mockResolvedValue(undefined), - removeSentMessage: vi.fn().mockResolvedValue(undefined) + removeSentMessage: vi.fn().mockResolvedValue(undefined), + getFailedOutbox: vi.fn().mockReturnValue([]), + removeFailedOutbox: vi.fn().mockResolvedValue(undefined) } })); @@ -169,7 +171,8 @@ vi.mock('../../src/js/channels.js', () => ({ handleControlMessage: vi.fn(), handleOverrideMessage: vi.fn(), storeReaction: vi.fn(), - sendWakeSignals: vi.fn().mockResolvedValue(undefined) + sendWakeSignals: vi.fn().mockResolvedValue(undefined), + messageFlow: { keepForRetry: vi.fn().mockResolvedValue(undefined) } } })); @@ -504,6 +507,26 @@ describe('DMManager', () => { expect(channel.messages.map(m => m.id)).toEqual(['ok']); }); + it('puts back the DMs that failed, unless a sent copy of the same id exists', async () => { + const peer = '0xpeer666666666666666666666666666666666666'; + const streamId = peer + '/Pombo-DM-1'; + const channel = { messageStreamId: streamId, type: 'dm', peerAddress: peer, messages: [] }; + channelManager.channels.set(streamId, channel); + dmManager.conversations.set(peer, streamId); + secureStorage.getSentMessages.mockReturnValueOnce([ + { id: 'delivered', text: 'went out on a retry', timestamp: 1 } + ]); + secureStorage.getFailedOutbox.mockReturnValueOnce([ + { id: 'delivered', text: 'went out on a retry', timestamp: 1, failError: 'No network' }, + { id: 'failed', text: 'still not sent', timestamp: 2, failError: 'No network' } + ]); + + await dmManager.loadDMTimeline(peer); + + expect(channel.messages.map(m => [m.id, !!m.failed])).toEqual([['delivered', false], ['failed', true]]); + expect(channel.messages[1]).toMatchObject({ failError: 'No network', _dmSent: true }); + }); + /** * The map that used to answer this question is rebuilt from each * record's peerAddress, so a record whose two halves disagree made @@ -1237,6 +1260,25 @@ describe('DMManager', () => { expect(streamrController.publishAs).not.toHaveBeenCalled(); }); + it('keeps a DM that failed for its retry, and drops it once a retry is sent', async () => { + const peerAddress = '0xpeerbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb'; + const streamId = `${peerAddress}/Pombo-DM-1`; + channelManager.channels.set(streamId, { + messageStreamId: streamId, type: 'dm', peerAddress, messages: [] + }); + streamrController.getDMPublicKey.mockResolvedValueOnce(null); + dmCrypto.peerPublicKeys.clear(); + + await expect(dmManager.sendMessage(streamId, 'Should fail')).rejects.toThrow(); + + const [keptStream, kept] = channelManager.messageFlow.keepForRetry.mock.calls.at(-1); + expect(keptStream).toBe(streamId); + expect(kept).toMatchObject({ text: 'Should fail', failed: true }); + + await dmManager.resendMessage(streamId, kept.id); + expect(secureStorage.removeFailedOutbox).toHaveBeenCalledWith(streamId, kept.id); + }); + it('should throw when wallet private key is missing', async () => { const peerAddress = '0xpeeraaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa'; const streamId = `${peerAddress}/Pombo-DM-1`; diff --git a/tests/unit/failedOutbox.test.js b/tests/unit/failedOutbox.test.js new file mode 100644 index 0000000..a67d839 --- /dev/null +++ b/tests/unit/failedOutbox.test.js @@ -0,0 +1,178 @@ +/** + * A text whose send failed is kept per conversation, so after a restart the + * bubble comes back "Not sent" and its Retry republishes the same message. + * The restart is real: the encrypted cache is written, dropped and read back. + */ + +import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'; + +// channels.js first: it builds a MessageFlow at load, and the import cycle +// leaves the class undefined if MessageFlow.js is evaluated before it. +await import('../../src/js/channels.js'); +const { MessageFlow } = await import('../../src/js/channels/MessageFlow.js'); +const { DeliveryConfirm } = await import('../../src/js/channels/DeliveryConfirm.js'); +const { secureStorage } = await import('../../src/js/secureStorage.js'); +const { identityManager } = await import('../../src/js/identity.js'); +const { authManager } = await import('../../src/js/auth.js'); + +const ME = '0x' + 'aa'.repeat(20); +const ROOM = `${ME}/room-1`; +const REPLY = { id: 'parent-1', sender: '0x' + 'bb'.repeat(20), senderName: 'Peer', text: 'the question' }; + +const signed = (id) => ({ + id, type: 'text', text: 'hello', sender: ME, senderName: 'me', timestamp: 1_000_000, + channelId: ROOM, replyTo: REPLY, signature: '0xsig' +}); + +async function unlockRealStorage() { + secureStorage.isGuestMode = false; + secureStorage.stateDB = null; + secureStorage.address = ME; + secureStorage.storageKey = await crypto.subtle.generateKey({ name: 'AES-GCM', length: 256 }, false, ['encrypt', 'decrypt']); + secureStorage.cache = secureStorage.createEmptyCache(); + secureStorage.isUnlocked = true; +} + +/** Drop the decrypted cache and read it back from the encrypted copy. */ +async function restart() { + secureStorage.cache = null; + await secureStorage.loadFromStorage(); +} + +function launch() { + const channel = { messageStreamId: ROOM, type: 'public', messages: [] }; + const manager = { + channels: new Map([[ROOM, channel]]), + notifyHandlers: vi.fn(), + sortMessagesByTimestamp: (ch) => ch.messages.sort((a, b) => a.timestamp - b.timestamp), + publishWithRetry: vi.fn().mockRejectedValue(new Error('No network')), + deliveryConfirm: { track: vi.fn() }, + rotationRetry: { settle: vi.fn() }, + sendWakeSignals: vi.fn().mockResolvedValue(undefined) + }; + const flow = new MessageFlow(manager); + manager.messageFlow = flow; + vi.spyOn(flow, '_assertMayPublish').mockResolvedValue(undefined); + return { channel, manager, flow }; +} + +const entries = () => secureStorage.getFailedOutbox(ROOM); + +describe('the failed-send outbox', () => { + let saved; + + beforeEach(async () => { + saved = { + cache: secureStorage.cache, isUnlocked: secureStorage.isUnlocked, isGuestMode: secureStorage.isGuestMode, + address: secureStorage.address, stateDB: secureStorage.stateDB, storageKey: secureStorage.storageKey + }; + localStorage.clear(); + await unlockRealStorage(); + vi.spyOn(authManager, 'getAddress').mockReturnValue(ME); + vi.spyOn(authManager, 'isConnected').mockReturnValue(true); + vi.spyOn(identityManager, 'createSignedMessage').mockImplementation(async () => signed('m1')); + vi.spyOn(identityManager, 'getTrustLevel').mockResolvedValue(0); + vi.spyOn(identityManager, 'resolveENS').mockResolvedValue(null); + }); + + afterEach(() => { + vi.restoreAllMocks(); + Object.assign(secureStorage, saved); + localStorage.clear(); + }); + + async function failSend() { + const { flow } = launch(); + await expect(flow.sendMessage(ROOM, 'hello', REPLY)).rejects.toThrow('No network'); + } + + it('keeps a failed send with everything its retry republishes, and no local state', async () => { + await failSend(); + + const [entry] = entries(); + expect(entry).toMatchObject({ id: 'm1', text: 'hello', sender: ME, timestamp: 1_000_000, signature: '0xsig', failError: 'No network' }); + expect(entry.replyTo).toEqual(REPLY); + for (const local of ['pending', 'failed', 'verified', 'delivered']) expect(entry).not.toHaveProperty(local); + }); + + it('brings the bubble back after a restart, and its retry republishes the same message', async () => { + await failSend(); + await restart(); + + const { channel, manager, flow } = launch(); + await flow.restoreFailedOutbox(channel); + const [bubble] = channel.messages; + expect(bubble).toMatchObject({ id: 'm1', failed: true, pending: false, failError: 'No network', replyTo: REPLY }); + + manager.publishWithRetry.mockResolvedValue({ timestamp: 1 }); + await flow.resendMessage(ROOM, 'm1'); + + const published = manager.publishWithRetry.mock.calls[0][1]; + expect(published).toMatchObject({ id: 'm1', timestamp: 1_000_000, replyTo: REPLY, signature: '0xsig' }); + expect(channel.messages[0].failed).toBe(false); + await restart(); + expect(entries()).toEqual([]); + }); + + it('keeps the entry, with the new reason, when the retry fails again', async () => { + await failSend(); + await restart(); + const { channel, manager, flow } = launch(); + await flow.restoreFailedOutbox(channel); + + manager.publishWithRetry.mockRejectedValue(new Error('Still offline')); + await expect(flow.resendMessage(ROOM, 'm1')).rejects.toThrow('Still offline'); + + await restart(); + expect(entries()).toMatchObject([{ id: 'm1', failError: 'Still offline' }]); + }); + + it('clears the bubble and the entry when our own copy arrives', async () => { + await failSend(); + const { channel, flow } = launch(); + await flow.restoreFailedOutbox(channel); + + await flow.handleTextMessage(ROOM, { ...signed('m1'), replyTo: null }); + + expect(channel.messages).toHaveLength(1); + expect(channel.messages[0].failed).toBe(false); + expect(entries()).toEqual([]); + }); + + it('does not put back what the timeline already holds', async () => { + await failSend(); + const { channel, flow } = launch(); + channel.messages.push({ ...signed('m1'), failed: true }); + + await flow.restoreFailedOutbox(channel); + + expect(channel.messages).toHaveLength(1); + }); + + it('keeps a sent text that storage never recorded as undelivered', async () => { + const { manager } = launch(); + const message = { ...signed('m2'), pending: false }; + const confirm = new DeliveryConfirm(manager); + + confirm._settle(ROOM, message, 'undelivered'); + await vi.waitFor(() => expect(entries()).toHaveLength(1)); + + expect(entries()[0]).toMatchObject({ id: 'm2', undelivered: true }); + }); + + it('never leaves the device: neither the sync nor the backup carries it', async () => { + await failSend(); + + expect(JSON.stringify(secureStorage.exportForSync())).not.toContain('failedOutbox'); + expect(JSON.stringify(secureStorage.exportForBackup())).not.toContain('failedOutbox'); + }); + + it('keeps the newest twenty per conversation', async () => { + for (let i = 1; i <= 21; i++) { + await secureStorage.putFailedOutbox(ROOM, { id: `m${i}`, timestamp: i }); + } + + expect(entries()).toHaveLength(20); + expect(entries()[0].id).toBe('m2'); + }); +});