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
5 changes: 5 additions & 0 deletions src/js/channels.js
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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);
Expand Down
1 change: 1 addition & 0 deletions src/js/channels/DeliveryConfirm.js
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
52 changes: 52 additions & 0 deletions src/js/channels/MessageFlow.js
Original file line number Diff line number Diff line change
Expand Up @@ -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];

Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand Down
1 change: 1 addition & 0 deletions src/js/channels/MessageOverrides.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 });

Expand Down
10 changes: 9 additions & 1 deletion src/js/dm.js
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -1408,6 +1409,7 @@ class DMManager {
message,
error: error.message
});
await channelManager.messageFlow.keepForRetry(peerInboxStreamId, message);
throw error;
}
}
Expand Down Expand Up @@ -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 });

Expand Down Expand Up @@ -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;
Expand Down
52 changes: 52 additions & 0 deletions src/js/secureStorage.js
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<Object>}
*/
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
Expand Down
4 changes: 4 additions & 0 deletions tests/unit/channels.extended.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
4 changes: 4 additions & 0 deletions tests/unit/channels.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
46 changes: 44 additions & 2 deletions tests/unit/dm.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}));

Expand Down Expand Up @@ -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) }
}
}));

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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`;
Expand Down
Loading
Loading