diff --git a/src/js/channels.js b/src/js/channels.js index c8484b5..e201cfa 100644 --- a/src/js/channels.js +++ b/src/js/channels.js @@ -43,6 +43,7 @@ import { ModDeltas, MOD_ACTION_TYPE } from './channels/ModDeltas.js'; import { stampChangedFields } from './syncMerge.js'; const GATE_REPAIR_WAIT_MS = 10_000; +const OPEN_READ_RETRY_DELAYS_MS = [5_000, 15_000, 30_000, 60_000]; class ChannelManager { constructor() { @@ -1914,6 +1915,9 @@ class ChannelManager { channel.initialLoadInProgress = true; // The override read that follows re-establishes every delete. channel._deletedIds = new Set(); + channel._openReads = {}; + channel.historyRetrying = false; + channel.historyReadFailed = false; } // Skip network subscription for write-only channels (no subscribe permission) @@ -2078,6 +2082,7 @@ class ChannelManager { // Reactions live on the -5 — their history is a separate read, // awaited here so the first render already has them instead of // popping in after. + let reactionsRead = null; if (channel.interactionsStreamId) { await streamrController.fetchHistoryAsync( channel.interactionsStreamId, @@ -2085,7 +2090,7 @@ class ChannelManager { STREAM_CONFIG.INITIAL_MESSAGES, (data) => this.handleControlMessage(messageStreamId, data), pwd, - null, + (s) => { reactionsRead = s; }, false, { quiet: true } ).catch(e => Logger.warn('Interactions history failed:', e.message)); @@ -2119,6 +2124,7 @@ class ChannelManager { // refusal closed. channel.historyError = stats?.readError || null; channel.hasMoreHistory = !channel.historyError; + if (stats?.failed || reactionsRead?.failed) this._retryOpenReads(messageStreamId); channel.initialLoadInProgress = false; @@ -2433,19 +2439,28 @@ class ChannelManager { }, CONFIG.subscriptions.renewalHistoryRetryMs); } - async _runHistoryRefresh(messageStreamId) { + /** + * @param {Object} [opts] + * @param {boolean} [opts.adminFirst] - read the -3 before the -1 rather than after + * @returns {Promise} whether every read of the group came back + */ + async _runHistoryRefresh(messageStreamId, { adminFirst = false } = {}) { const channel = this.channels.get(messageStreamId); - if (!channel) return; + if (!channel) return false; // Never overlap the initial load or another refresh: concurrent P0/P1 // fetches let an override land before its target and park forever in // _pendingOverrides. Reschedule through the debounce instead. if (channel.initialLoadInProgress || channel._historyRefreshRunning) { this.refreshHistory(messageStreamId); - return; + return false; } channel._historyRefreshRunning = true; Logger.info('Refreshing history:', messageStreamId.slice(-20)); + if (adminFirst) { + await this.refreshAdminState(messageStreamId) + .catch(e => Logger.warn('Admin state refresh failed:', e?.message)); + } // Same discipline as the initial-load pipeline: gate per-message // renders (no flash of pre-override originals), P0 before P1, then @@ -2453,6 +2468,8 @@ class ChannelManager { // deleted messages), and render ONCE. channel.initialLoadInProgress = true; let contentRead = null; + let controlRead = null; + let reactionsRead = null; try { await streamrController.fetchHistoryAsync( messageStreamId, @@ -2471,7 +2488,7 @@ class ChannelManager { STREAM_CONFIG.INITIAL_MESSAGES, (data) => this.handleOverrideMessage(messageStreamId, data, true), channel.password || null, - null, + (stats) => { controlRead = stats; }, false, { quiet: true } ); @@ -2486,7 +2503,7 @@ class ChannelManager { STREAM_CONFIG.INITIAL_MESSAGES, (data) => this.handleControlMessage(messageStreamId, data), channel.password || null, - null, + (stats) => { reactionsRead = stats; }, false, { quiet: true } ).catch(e => Logger.warn('Interactions history failed:', e.message)); @@ -2514,13 +2531,52 @@ class ChannelManager { // until this key arrived — pins/moderation and a hidden channel's // image need their own re-pull (the refresh above only covers -1; // the admin poller would take up to a full tick to converge). - await this.refreshAdminState(messageStreamId); + if (!adminFirst) await this.refreshAdminState(messageStreamId); const adminId = channel.adminStreamId || deriveAdminId(messageStreamId); if (adminId) { channelImageManager.get(adminId, { password: channel.password || null, force: true }) .catch(() => {}); } this.notifyHandlers('initial_history_complete', { streamId: messageStreamId }); + return !!contentRead && !contentRead.failed && !controlRead?.failed && !reactionsRead?.failed; + } + + /** + * The open's reads failed, which is not an empty channel: read them again, + * the -3 first so nothing paints without its moderation. A cure replaces + * the client and re-reads the active channel itself, so that ends this; + * so does leaving or reopening the channel. + */ + _retryOpenReads(messageStreamId) { + const channel = this.channels.get(messageStreamId); + if (!channel) return; + const openReads = channel._openReads; + const client = streamrController.client; + const stillOpen = () => this.channels.get(messageStreamId) === channel + && channel._openReads === openReads + && this.currentChannel === messageStreamId + && streamrController.client === client; + channel.historyRetrying = true; + (async () => { + for (const [attempt, wait] of OPEN_READ_RETRY_DELAYS_MS.entries()) { + await new Promise(resolve => setTimeout(resolve, wait)); + if (!stillOpen()) return; + Logger.info(`History ${messageStreamId.slice(-20)}: reading the open's group again (attempt ${attempt + 1})`); + const readAll = await this._runHistoryRefresh(messageStreamId, { adminFirst: true }); + if (!stillOpen()) return; + if (readAll) { + channel.historyRetrying = false; + // The open's rule: what a paginate over the failing reads concluded is no verdict. + channel.hasMoreHistory = !channel.historyError; + this.notifyHandlers('initial_history_complete', { streamId: messageStreamId }); + return; + } + } + channel.historyRetrying = false; + channel.historyReadFailed = true; + channel.hasMoreHistory = false; + this.notifyHandlers('initial_history_complete', { streamId: messageStreamId }); + })().catch(e => Logger.warn('Open history retry failed:', e?.message)); } // ==================== Message Overrides ==================== diff --git a/src/js/streamr.js b/src/js/streamr.js index 235e813..67fd60d 100644 --- a/src/js/streamr.js +++ b/src/js/streamr.js @@ -4752,8 +4752,8 @@ class StreamrController { // even when the iterator never signals `done` (e.g. legacy single- // partition channels). const historyStats = { - content: { loaded: 0, requested: 0, readError: null }, - control: { loaded: 0, requested: 0, readError: null }, + content: { loaded: 0, requested: 0, readError: null, failed: false }, + control: { loaded: 0, requested: 0, readError: null, failed: false }, }; const maybeSignalHistoryComplete = async () => { @@ -4767,6 +4767,7 @@ class StreamrController { controlLoaded: historyStats.control.loaded, controlRequested: historyStats.control.requested, readError: historyStats.content.readError || historyStats.control.readError, + failed: historyStats.content.failed || historyStats.control.failed, }); } catch (e) { Logger.warn('onHistoryComplete error:', e); } } @@ -4782,12 +4783,14 @@ class StreamrController { loaded: stats.loaded ?? 0, requested: stats.requested ?? 0, readError: stats.readError ?? null, + failed: !!stats.failed, }; } else if (stats && partition === STREAM_CONFIG.MESSAGE_STREAM.CONTROL) { historyStats.control = { loaded: stats.loaded ?? 0, requested: stats.requested ?? 0, readError: stats.readError ?? null, + failed: !!stats.failed, }; } Logger.debug(`History complete for ${partitionLabel}. Pending: ${pendingHistoryCompletions}`); diff --git a/src/js/streamr/History.js b/src/js/streamr/History.js index 4326ac5..7a9c418 100644 --- a/src/js/streamr/History.js +++ b/src/js/streamr/History.js @@ -564,6 +564,8 @@ export class History { // an out-of-scope variable via typeof and always reported 0) let rawCount = 0; let readError = null; + let failed = false; + let iterationBroke = false; try { Logger.debug(`Fetching ${count} historical messages for partition ${partition}${password ? ' (encrypted)' : ''}...`); @@ -662,6 +664,7 @@ export class History { } // Other iterator errors - log and try to continue Logger.warn('History iteration error:', iterError.message); + iterationBroke = true; continue; } @@ -815,11 +818,17 @@ export class History { } catch (error) { // CORS errors and other network issues are caught here Logger.warn(`History fetch failed for partition ${partition} (may be CORS on localhost):`, error.message); + // A stream the client has no storage for is an answer: the SDK + // keeps that verdict for the client's life, so reading again cannot help. + failed = error?.code !== 'NO_STORAGE_NODES' && !String(error?.message).includes('NO_STORAGE_NODES'); } finally { // A refusal by the storage node surfaces as an iterator error the // loop above skips, so the verdict comes from the fetch layer, which // clears it on the next successful read of this partition. readError = storageFetch.lastReadError(streamId, partition) || null; + // A 4xx is the node's answer about this read; any other break + // (network, 5xx, a page refused for its storedAt) left it unread. + if (iterationBroke && !(readError?.status >= 400 && readError?.status < 500)) failed = true; // Signal that initial history fetch is complete (success or failure). // Pass `loaded`/`requested` so callers can detect exhaustion (when // fewer raw messages came back than requested → no more history @@ -827,7 +836,7 @@ export class History { // `hasMoreHistory=false` deterministically. if (onHistoryComplete) { try { - await onHistoryComplete({ loaded: rawCount, requested: count, readError }); + await onHistoryComplete({ loaded: rawCount, requested: count, readError, failed }); } catch (e) { Logger.warn('onHistoryComplete error:', e); } } } diff --git a/src/js/ui/ChatAreaUI.js b/src/js/ui/ChatAreaUI.js index 118bcb0..e3c9c50 100644 --- a/src/js/ui/ChatAreaUI.js +++ b/src/js/ui/ChatAreaUI.js @@ -638,6 +638,7 @@ class ChatAreaUI { // fallback that raced slow resend iterators. const isTerminalEmpty = effectiveChannel?.initialLoadInProgress !== true && + effectiveChannel?.historyRetrying !== true && hasMoreHistory === false && this._loadOp == null && !effectiveChannel?.loadingHistory && @@ -690,6 +691,13 @@ class ChatAreaUI { `; this.messagesArea.querySelector('#empty-state-renew-btn') ?.addEventListener('click', () => subscriptionBannerUI.renewCurrent()); + } else if (effectiveChannel?.historyReadFailed) { + this.messagesArea.innerHTML = ` +
+ Channel history could not be loaded + Check your connection and reopen the channel +
+ `; } else { this.messagesArea.innerHTML = waitingForKeys ? ` @@ -723,9 +731,9 @@ class ChatAreaUI { const currentAddress = authManager?.getAddress(); let historyStartIndicator = ''; - // A refused read also clears hasMoreHistory, and there the start is - // unknown, not reached. - if (!hasMoreHistory && !effectiveChannel?.historyError) { + // A refused or failed read also clears hasMoreHistory, and there the + // start is unknown, not reached. + if (!hasMoreHistory && !effectiveChannel?.historyError && !effectiveChannel?.historyReadFailed) { historyStartIndicator = `
diff --git a/tests/unit/ChatAreaUI.historyStart.test.js b/tests/unit/ChatAreaUI.historyStart.test.js index 06888b8..d1f0adc 100644 --- a/tests/unit/ChatAreaUI.historyStart.test.js +++ b/tests/unit/ChatAreaUI.historyStart.test.js @@ -61,13 +61,14 @@ import { subscriptionBannerUI } from '../../src/js/ui/SubscriptionBannerUI.js'; const MESSAGES = [{ id: 'm1', text: 'hello', sender: '0xabc', timestamp: 1_789_000_000_000 }]; -function render({ hasMoreHistory, historyError }) { +function render({ hasMoreHistory, historyError, historyReadFailed = false }) { const channel = { streamId: '0xowner/chan-1', name: 'Chan', messages: MESSAGES, hasMoreHistory, - historyError + historyError, + historyReadFailed }; chatAreaUI.setDependencies({ getActiveChannel: () => channel, @@ -114,6 +115,10 @@ describe('the start-of-history line', () => { expect(html).toContain('history-error-banner'); }); + it('stays away when the reads of the open gave up over what the cache holds', () => { + expect(claimsTheStart(render({ hasMoreHistory: false, historyError: null, historyReadFailed: true }))).toBe(false); + }); + it('stays away on a lapsed gate, where the subscription strip explains instead', () => { subscriptionBannerUI.stateOf.mockReturnValue('expired'); const html = render({ diff --git a/tests/unit/channels.extended.test.js b/tests/unit/channels.extended.test.js index b4b8178..9ba5f64 100644 --- a/tests/unit/channels.extended.test.js +++ b/tests/unit/channels.extended.test.js @@ -1609,6 +1609,18 @@ describe('ChannelManager Extended', () => { expect(channel.hasMoreHistory).toBe(false); }); + it('reads the open again after a read that failed, and not after a clean one', async () => { + const retry = vi.spyOn(channelManager, '_retryOpenReads').mockImplementation(() => {}); + + await loadWith({ readError: null, failed: false }); + expect(retry).not.toHaveBeenCalled(); + + channelManager.channels.delete(streamId); + await loadWith({ readError: null, failed: true }); + expect(retry).toHaveBeenCalledWith(streamId); + retry.mockRestore(); + }); + it('reopens it on the next clean read, with no page reload', async () => { const channel = await loadWith({ readError: { status: 403, signed: true, at: Date.now() } }); expect(channel.hasMoreHistory).toBe(false); diff --git a/tests/unit/channels.openReadRetry.test.js b/tests/unit/channels.openReadRetry.test.js new file mode 100644 index 0000000..788b3ab --- /dev/null +++ b/tests/unit/channels.openReadRetry.test.js @@ -0,0 +1,130 @@ +/** + * A read of the open that failed is not an empty channel: it is read again, + * the -3 first, on a short backoff, and gives up only after the last attempt. + * A cure replaces the client and re-reads the active channel itself, so a + * retry that sees a new client stops instead of reading a second time. + */ +import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'; +import { ethers } from 'ethers'; + +globalThis.ethers = ethers; + +const { channelManager } = await import('../../src/js/channels.js'); +const { streamrController, STREAM_CONFIG } = await import('../../src/js/streamr.js'); +const { channelImageManager } = await import('../../src/js/channelImageManager.js'); + +const ID = '0xowner/room-1'; +const BACKOFF_MS = [5_000, 15_000, 30_000, 60_000]; + +function open() { + const channel = { + messageStreamId: ID, + streamId: ID, + adminStreamId: '0xowner/room-3', + messages: [], + hasMoreHistory: true, + _controlPartitionSupported: false, + _openReads: {} + }; + channelManager.channels.set(ID, channel); + channelManager.currentChannel = ID; + return channel; +} + +describe('reading the open again after a failed read', () => { + let failing; + let calls; + + beforeEach(() => { + vi.useFakeTimers(); + failing = true; + calls = []; + streamrController.client = { id: 'first' }; + vi.spyOn(streamrController, 'fetchHistoryAsync').mockImplementation( + async (id, partition, count, handler, password, done) => { + if (id !== ID || partition !== STREAM_CONFIG.MESSAGE_STREAM.MESSAGES) return; + calls.push('content'); + await done?.({ loaded: 0, requested: count, readError: null, failed: failing }); + }); + vi.spyOn(channelManager, 'refreshAdminState').mockImplementation(async () => { calls.push('admin'); }); + for (const method of ['flushBatchVerification', 'awaitAllFlushes']) { + vi.spyOn(channelManager, method).mockResolvedValue(undefined); + } + vi.spyOn(channelManager, 'applyPendingOverrides').mockImplementation(() => {}); + vi.spyOn(channelManager, 'sortMessagesByTimestamp').mockImplementation(() => {}); + vi.spyOn(channelManager, 'notifyHandlers').mockImplementation(() => {}); + vi.spyOn(channelImageManager, 'get').mockResolvedValue(null); + }); + + afterEach(() => { + channelManager.channels.delete(ID); + channelManager.currentChannel = null; + streamrController.client = null; + vi.useRealTimers(); + vi.restoreAllMocks(); + }); + + it('reads the -3 first, then the history, and stops once it came back', async () => { + const channel = open(); + failing = false; + + channelManager._retryOpenReads(ID); + expect(channel.historyRetrying).toBe(true); + await vi.advanceTimersByTimeAsync(BACKOFF_MS[0]); + + expect(calls).toEqual(['admin', 'content']); + expect(channel.historyRetrying).toBe(false); + expect(channel.historyReadFailed).toBeFalsy(); + + await vi.advanceTimersByTimeAsync(BACKOFF_MS[1] + BACKOFF_MS[2] + BACKOFF_MS[3]); + expect(calls).toEqual(['admin', 'content']); + }); + + it('reopens the paging a paginate over the failing reads closed', async () => { + const channel = open(); + channel.hasMoreHistory = false; + failing = false; + + channelManager._retryOpenReads(ID); + await vi.advanceTimersByTimeAsync(BACKOFF_MS[0]); + + expect(channel.hasMoreHistory).toBe(true); + }); + + it('gives up after the last attempt and closes paging', async () => { + const channel = open(); + + channelManager._retryOpenReads(ID); + await vi.advanceTimersByTimeAsync(BACKOFF_MS.reduce((a, b) => a + b, 0)); + + expect(calls.filter(c => c === 'content')).toHaveLength(BACKOFF_MS.length); + expect(channel.historyRetrying).toBe(false); + expect(channel.historyReadFailed).toBe(true); + expect(channel.hasMoreHistory).toBe(false); + }); + + it('stops when a cure replaced the client, which re-reads on its own', async () => { + const channel = open(); + + channelManager._retryOpenReads(ID); + streamrController.client = { id: 'rebuilt' }; + await vi.advanceTimersByTimeAsync(BACKOFF_MS.reduce((a, b) => a + b, 0)); + + expect(calls).toEqual([]); + expect(channel.historyReadFailed).toBeFalsy(); + }); + + it('stops when the user left the channel, or opened it again', async () => { + open(); + channelManager._retryOpenReads(ID); + channelManager.currentChannel = '0xowner/other-1'; + await vi.advanceTimersByTimeAsync(BACKOFF_MS[0]); + expect(calls).toEqual([]); + + const reopened = open(); + channelManager._retryOpenReads(ID); + reopened._openReads = {}; + await vi.advanceTimersByTimeAsync(BACKOFF_MS[0]); + expect(calls).toEqual([]); + }); +}); diff --git a/tests/unit/streamr.dualStream.test.js b/tests/unit/streamr.dualStream.test.js index 01ae944..f804d4a 100644 --- a/tests/unit/streamr.dualStream.test.js +++ b/tests/unit/streamr.dualStream.test.js @@ -112,9 +112,23 @@ describe('subscribeToDualStream', () => { expect(done).toHaveBeenCalledTimes(1); expect(done).toHaveBeenCalledWith({ contentLoaded: 12, contentRequested: 30, controlLoaded: 3, controlRequested: 30, readError: null, + failed: false, }); }); + it('says the history read failed when either stored partition failed', async () => { + const calls = stubSubscribe(); + const done = vi.fn(); + + await streamrController.subscribeToDualStream( + MESSAGE, EPHEMERAL, { onMessage: () => {}, onOverride: () => {} }, null, 30, done + ); + await calls[0].onDone({ loaded: 0, requested: 30, failed: true }); + await calls[1].onDone({ loaded: 12, requested: 30 }); + + expect(done).toHaveBeenCalledWith(expect.objectContaining({ failed: true })); + }); + it('signals completion once even if a partition reports twice', async () => { const calls = stubSubscribe(); const done = vi.fn(); diff --git a/tests/unit/streamr.initialHistory.test.js b/tests/unit/streamr.initialHistory.test.js index 25af397..00d9601 100644 --- a/tests/unit/streamr.initialHistory.test.js +++ b/tests/unit/streamr.initialHistory.test.js @@ -24,6 +24,7 @@ vi.mock('../../src/js/envelopeSigner.js', async (importOriginal) => ({ const { streamrController } = await import('../../src/js/streamr.js'); const { cryptoManager } = await import('../../src/js/crypto.js'); +const { storageFetch } = await import('../../src/js/storageFetch.js'); const OWNER = '0x' + '77'.repeat(20); const AUTHOR = '0x' + '88'.repeat(20); @@ -45,6 +46,26 @@ function serve(...messages) { return calls; } +/** + * The SDK's shape for a storage node that fails the read: `resend()` resolves + * and the error comes out of the iterator. Each step is a row or an error. + */ +function serveSteps(...steps) { + let i = 0; + const iterator = { + next: async () => { + if (i >= steps.length) return { done: true, value: undefined }; + const step = steps[i++]; + if (step instanceof Error) throw step; + return { done: false, value: step }; + }, + [Symbol.asyncIterator]() { return this; }, + }; + streamrController.client = { resend: vi.fn(async () => iterator) }; +} + +const storageNodeError = () => Object.assign(new Error('Failed to fetch'), { code: 'STORAGE_NODE_ERROR' }); + const row = (ts, content) => ({ content, timestamp: ts, @@ -86,7 +107,7 @@ describe('fetchHistoryAsync', () => { await streamrController.fetchHistoryAsync(MESSAGE, 0, 40, handler, null, done); expect(handler).toHaveBeenCalledTimes(1); - expect(done).toHaveBeenCalledWith({ loaded: 3, requested: 40, readError: null }); + expect(done).toHaveBeenCalledWith({ loaded: 3, requested: 40, readError: null, failed: false }); }); it('still reports when storage is unreachable, so the channel does not hang', async () => { @@ -95,7 +116,62 @@ describe('fetchHistoryAsync', () => { await streamrController.fetchHistoryAsync(MESSAGE, 0, 40, () => {}, null, done); - expect(done).toHaveBeenCalledWith({ loaded: 0, requested: 40, readError: null }); + expect(done).toHaveBeenCalledWith({ loaded: 0, requested: 40, readError: null, failed: true }); + }); + + it('reports a stream with no storage as answered, not as a failed read', async () => { + const noStorage = Object.assign(new Error(`no storage assigned: ${MESSAGE}`), { code: 'NO_STORAGE_NODES' }); + streamrController.client = { resend: vi.fn(async () => { throw noStorage; }) }; + const done = vi.fn(); + + await streamrController.fetchHistoryAsync(MESSAGE, 0, 40, () => {}, null, done); + + expect(done).toHaveBeenCalledWith({ loaded: 0, requested: 40, readError: null, failed: false }); + }); + + it('says the read failed when the storage node breaks off the iteration', async () => { + serveSteps(row(100, { type: 'text', id: 'a', text: 'x', sender: AUTHOR, timestamp: 100 }), storageNodeError()); + const handler = vi.fn(); + const done = vi.fn(); + + await streamrController.fetchHistoryAsync(MESSAGE, 0, 40, handler, null, done); + + expect(handler).toHaveBeenCalledTimes(1); + expect(done).toHaveBeenCalledWith({ loaded: 1, requested: 40, readError: null, failed: true }); + }); + + it('counts a page refused for the storedAt its node owes as a failed read', async () => { + const refused = { status: 503, signed: false, reason: 'storedAt' }; + vi.spyOn(storageFetch, 'lastReadError').mockReturnValue(refused); + serveSteps(storageNodeError()); + const done = vi.fn(); + + await streamrController.fetchHistoryAsync(MESSAGE, 0, 40, () => {}, null, done); + + expect(done).toHaveBeenCalledWith({ loaded: 0, requested: 40, readError: refused, failed: true }); + }); + + it('takes a 4xx from the storage node as its answer, not as a failed read', async () => { + const refused = { status: 403, signed: true }; + vi.spyOn(storageFetch, 'lastReadError').mockReturnValue(refused); + serveSteps(storageNodeError()); + const done = vi.fn(); + + await streamrController.fetchHistoryAsync(MESSAGE, 0, 40, () => {}, null, done); + + expect(done).toHaveBeenCalledWith({ loaded: 0, requested: 40, readError: refused, failed: false }); + }); + + it('does not fail the read over a row it cannot decrypt', async () => { + const decrypt = Object.assign(new Error('no encryption key'), { code: 'DECRYPT_ERROR' }); + serveSteps(decrypt, row(100, { type: 'text', id: 'a', text: 'x', sender: AUTHOR, timestamp: 100 })); + const handler = vi.fn(); + const done = vi.fn(); + + await streamrController.fetchHistoryAsync(MESSAGE, 0, 40, handler, null, done); + + expect(handler).toHaveBeenCalledTimes(1); + expect(done).toHaveBeenCalledWith({ loaded: 1, requested: 40, readError: null, failed: false }); }); it('survives a completion callback that throws', async () => {