diff --git a/src/js/app.js b/src/js/app.js
index 5b848b2..dd9c59f 100644
--- a/src/js/app.js
+++ b/src/js/app.js
@@ -848,6 +848,8 @@ class App {
if (data.streamId === currentStreamId) {
Logger.debug(`History loaded: ${data.loaded} messages, hasMore: ${data.hasMore}`);
}
+ } else if (event === 'history_page_due') {
+ if (data.streamId === currentStreamId) chatAreaUI.handleMessagesScroll();
} else if (event === 'history_batch_loaded') {
if (data.streamId === currentStreamId) {
const channel = channelManager.getCurrentChannel();
diff --git a/src/js/channels.js b/src/js/channels.js
index e201cfa..22dc1b1 100644
--- a/src/js/channels.js
+++ b/src/js/channels.js
@@ -1918,6 +1918,8 @@ class ChannelManager {
channel._openReads = {};
channel.historyRetrying = false;
channel.historyReadFailed = false;
+ channel.overridesOwed = false;
+ this.messageFlow.resetPagingBackoff(channel);
}
// Skip network subscription for write-only channels (no subscribe permission)
@@ -2124,6 +2126,7 @@ class ChannelManager {
// refusal closed.
channel.historyError = stats?.readError || null;
channel.hasMoreHistory = !channel.historyError;
+ channel.overridesOwed = !!stats?.overridesFailed;
if (stats?.failed || reactionsRead?.failed) this._retryOpenReads(messageStreamId);
channel.initialLoadInProgress = false;
@@ -2421,8 +2424,16 @@ class ChannelManager {
if (!channel?.gate?.address) return;
clearTimeout(channel._historyRefreshTimer);
channel._historyRefreshTimer = setTimeout(() => {
- this._runHistoryRefresh(messageStreamId).catch(e =>
- Logger.warn('History refresh failed:', e.message));
+ this._runHistoryRefresh(messageStreamId)
+ .then(() => {
+ // Its edits and deletions did not come back and nothing
+ // is reading them again yet.
+ if (channel.overridesOwed && !channel.historyRetrying && !channel.historyReadFailed
+ && this.currentChannel === messageStreamId) {
+ this._retryOpenReads(messageStreamId);
+ }
+ })
+ .catch(e => Logger.warn('History refresh failed:', e.message));
}, 1500);
}
@@ -2513,6 +2524,7 @@ class ChannelManager {
await this.awaitAllFlushes(messageStreamId);
this.applyPendingOverrides(channel);
this.sortMessagesByTimestamp(channel);
+ if (controlRead) channel.overridesOwed = !!controlRead.failed;
// A clean read reopens only what a refusal closed: an exhausted
// history stays exhausted.
diff --git a/src/js/channels/MessageFlow.js b/src/js/channels/MessageFlow.js
index 8dc3b39..b634684 100644
--- a/src/js/channels/MessageFlow.js
+++ b/src/js/channels/MessageFlow.js
@@ -18,6 +18,8 @@ import { messageTime } from '../utils/messageTime.js';
import { storageFetch } from '../storageFetch.js';
import { NO_NETWORK, isOffline } from '../utils/network.js';
+const PAGING_RETRY_DELAYS_MS = [5_000, 15_000, 30_000, 60_000];
+
export class MessageFlow {
/**
* Self-calls go through `manager` on purpose: while the manager is still
@@ -757,6 +759,30 @@ export class MessageFlow {
channel.messages.sort((a, b) => messageTime(a) - messageTime(b));
}
+ /** Forgets an incomplete page's backoff: on a reopen, or once a page came back whole. */
+ resetPagingBackoff(channel) {
+ clearTimeout(channel._pagingNudgeTimer);
+ channel._pagingNudgeTimer = null;
+ channel._pagingFailures = 0;
+ channel._pagingRetryAt = 0;
+ }
+
+ /** Asks for the page again after a backoff; past the last one, only the user's scroll does. */
+ _backOffPaging(channel, messageStreamId, generation) {
+ const failures = channel._pagingFailures || 0;
+ const wait = PAGING_RETRY_DELAYS_MS[Math.min(failures, PAGING_RETRY_DELAYS_MS.length - 1)];
+ channel._pagingFailures = failures + 1;
+ channel._pagingRetryAt = Date.now() + wait;
+ clearTimeout(channel._pagingNudgeTimer);
+ channel._pagingNudgeTimer = null;
+ if (channel._pagingFailures > PAGING_RETRY_DELAYS_MS.length) return;
+ channel._pagingNudgeTimer = setTimeout(() => {
+ channel._pagingNudgeTimer = null;
+ if (this.manager.switchGeneration !== generation) return;
+ this.manager.notifyHandlers('history_page_due', { streamId: messageStreamId });
+ }, wait);
+ }
+
/**
* Load more (older) history from MESSAGE stream - for lazy loading / infinite scroll
* In dual-stream architecture, only messageStream has storage
@@ -816,6 +842,11 @@ export class MessageFlow {
Logger.debug('No more history to load');
return { loaded: 0, hasMore: false };
}
+ // The open's reads are the retry's until they are back: its success
+ // resets hasMoreHistory, and would undo what a page read meanwhile found.
+ if (channel.historyRetrying || Date.now() < (channel._pagingRetryAt || 0)) {
+ return { loaded: 0, hasMore: channel.hasMoreHistory };
+ }
// Use oldest message timestamp, or Date.now() if history was only reactions
const beforeTimestamp = channel.oldestTimestamp || Date.now();
@@ -885,6 +916,18 @@ export class MessageFlow {
fetchContent(), fetchOverrides(), fetchReactions()
]);
+ // Taken whole or not at all: content without its overrides paints
+ // edits and deletions undone, and the cursor would move past
+ // overrides that nothing reads again.
+ if (contentResult.failed || overrideResult.failed) {
+ Logger.warn('loadMoreHistory: page came back incomplete, asking again later');
+ channel.loadingHistory = false;
+ this.manager.notifyHandlers('history_loading', { streamId: messageStreamId, loading: false });
+ this._backOffPaging(channel, messageStreamId, generationAtStart);
+ return { loaded: 0, hasMore: channel.hasMoreHistory };
+ }
+ this.resetPagingBackoff(channel);
+
// Older reactions never gate "is there more history": that is the
// -1's question, and a reaction page that runs dry says nothing
// about the conversation behind it.
@@ -928,7 +971,7 @@ export class MessageFlow {
await new Promise(r => setTimeout(r, backoffMs[attempt]));
if (signal?.aborted || this.manager.switchGeneration !== generationAtStart) break;
const [retryContent, retryOverride] = await Promise.all([fetchContent(), fetchOverrides()]);
- if (!isEmpty(retryContent, retryOverride)) {
+ if (!isEmpty(retryContent, retryOverride) && !retryContent.failed && !retryOverride.failed) {
Logger.info(
`loadMoreHistory: retry #${attempt + 1} recovered messages after empty first response`
);
diff --git a/src/js/streamr.js b/src/js/streamr.js
index 67fd60d..b9c5bab 100644
--- a/src/js/streamr.js
+++ b/src/js/streamr.js
@@ -4768,6 +4768,7 @@ class StreamrController {
controlRequested: historyStats.control.requested,
readError: historyStats.content.readError || historyStats.control.readError,
failed: historyStats.content.failed || historyStats.control.failed,
+ overridesFailed: historyStats.control.failed,
});
} catch (e) { Logger.warn('onHistoryComplete error:', e); }
}
diff --git a/src/js/streamr/History.js b/src/js/streamr/History.js
index 7a9c418..e4d07ba 100644
--- a/src/js/streamr/History.js
+++ b/src/js/streamr/History.js
@@ -217,7 +217,9 @@ export class History {
// response, no error). Without confirmation, one truncated response
// falsely latches `hasMore: false` and kills scroll-up pagination
// for the rest of the session.
+ let lastPassBroke = false;
const collectRange = async () => {
+ lastPassBroke = false;
// Gated: raw resend — same reason as fetchHistoryAsync (the SDK
// validator re-checks stored envelopes against the present gate
// state and erases ex-members' history; authorship comes from the
@@ -267,6 +269,7 @@ export class History {
continue;
}
Logger.warn('fetchOlderHistory iteration error:', iterError.message);
+ lastPassBroke = true;
continue;
}
@@ -385,6 +388,7 @@ export class History {
Logger.debug(`Fetching ${count} older messages before ${new Date(beforeTimestamp).toISOString()}`);
let allMessages = await collectRange();
+ let broke = lastPassBroke;
// Exhaustion claimed (range returned ≤ count messages)? Confirm on a
// fresh connection before trusting it — one truncated response would
@@ -406,6 +410,7 @@ export class History {
if (!byKey.has(k)) byKey.set(k, m);
}
allMessages = Array.from(byKey.values());
+ broke = broke && lastPassBroke;
} catch (confirmError) {
Logger.debug('fetchOlderHistory: confirmation pass failed (keeping first result):', confirmError.message);
}
@@ -432,13 +437,18 @@ export class History {
// P0/P1 breakdown — keep this one at debug level to avoid duplicate noise.
Logger.debug(`fetchOlderHistory partition ${partition}: ${resultMessages.length} messages (hasMore: ${hasMore})`);
+ // Same rule as the open's read: a 4xx is the node's answer, any
+ // other break left the page unread.
+ const readError = storageFetch.lastReadError(messageStreamId, partition);
return {
messages: resultMessages,
- hasMore: hasMore
+ hasMore: hasMore,
+ failed: broke && !(readError?.status >= 400 && readError?.status < 500)
};
} catch (error) {
Logger.warn('Older history fetch error:', error.message);
- return { messages: [], hasMore: false };
+ const noStorage = error?.code === 'NO_STORAGE_NODES' || String(error?.message).includes('NO_STORAGE_NODES');
+ return { messages: [], hasMore: false, failed: !noStorage };
}
}
diff --git a/src/js/ui/ChatAreaUI.js b/src/js/ui/ChatAreaUI.js
index e3c9c50..3f63ff0 100644
--- a/src/js/ui/ChatAreaUI.js
+++ b/src/js/ui/ChatAreaUI.js
@@ -731,9 +731,10 @@ class ChatAreaUI {
const currentAddress = authManager?.getAddress();
let historyStartIndicator = '';
- // A refused or failed read also clears hasMoreHistory, and there the
- // start is unknown, not reached.
- if (!hasMoreHistory && !effectiveChannel?.historyError && !effectiveChannel?.historyReadFailed) {
+ // A refused, failed or unfinished read also leaves hasMoreHistory
+ // false, and there the start is unknown, not reached.
+ if (!hasMoreHistory && !effectiveChannel?.historyError && !effectiveChannel?.historyReadFailed
+ && effectiveChannel?.historyRetrying !== true) {
historyStartIndicator = `
@@ -757,6 +758,16 @@ class ChatAreaUI {
`;
}
+ let overridesBanner = '';
+ if (effectiveChannel?.overridesOwed) {
+ overridesBanner = `
+
+ ${effectiveChannel.historyReadFailed
+ ? 'Edits and deletions could not be loaded. Reopen the channel'
+ : 'Loading edits and deletions…'}
+
+ `;
+ }
// Analyze message groups for Stack Effect
const groupPositions = analyzeMessageGroups(messagesForRender);
@@ -832,7 +843,7 @@ class ChatAreaUI {
messagesHtml += messageRenderer.buildMessageGroupCloseHTML();
}
- this.messagesArea.innerHTML = historyErrorBanner + historyStartIndicator + messagesHtml;
+ this.messagesArea.innerHTML = historyErrorBanner + overridesBanner + historyStartIndicator + messagesHtml;
if (!this.isLoadingMore) {
this.messagesArea.scrollTop = this.messagesArea.scrollHeight;
diff --git a/src/js/ui/PreviewModeUI.js b/src/js/ui/PreviewModeUI.js
index 75bc974..46dc0b5 100644
--- a/src/js/ui/PreviewModeUI.js
+++ b/src/js/ui/PreviewModeUI.js
@@ -1109,6 +1109,10 @@ class PreviewModeUI {
);
let [contentResult, overrideResult] = await Promise.all([fetchContent(), fetchOverrides()]);
+ // Taken whole or not at all, as in channels.loadMoreHistory.
+ if (contentResult.failed || overrideResult.failed) {
+ return { loaded: 0, hasMore: channel.hasMoreHistory !== false };
+ }
// Storage race mitigation (symmetric with `channels.loadMoreHistory`):
// when both partitions return zero across a non-trivial range, the
@@ -1122,7 +1126,7 @@ class PreviewModeUI {
await new Promise(r => setTimeout(r, backoffMs[attempt]));
if (this.previewChannel !== channel) break;
const [retryContent, retryOverride] = await Promise.all([fetchContent(), fetchOverrides()]);
- if (!isEmpty(retryContent, retryOverride)) {
+ if (!isEmpty(retryContent, retryOverride) && !retryContent.failed && !retryOverride.failed) {
Logger.info?.(
`loadMorePreviewHistory: retry #${attempt + 1} recovered messages after empty first response`
);
diff --git a/tests/unit/ChatAreaUI.historyStart.test.js b/tests/unit/ChatAreaUI.historyStart.test.js
index d1f0adc..989a380 100644
--- a/tests/unit/ChatAreaUI.historyStart.test.js
+++ b/tests/unit/ChatAreaUI.historyStart.test.js
@@ -61,14 +61,16 @@ 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, historyReadFailed = false }) {
+function render({ hasMoreHistory, historyError, historyReadFailed = false, historyRetrying = false, overridesOwed = false }) {
const channel = {
streamId: '0xowner/chan-1',
name: 'Chan',
messages: MESSAGES,
hasMoreHistory,
historyError,
- historyReadFailed
+ historyReadFailed,
+ historyRetrying,
+ overridesOwed
};
chatAreaUI.setDependencies({
getActiveChannel: () => channel,
@@ -119,6 +121,10 @@ describe('the start-of-history line', () => {
expect(claimsTheStart(render({ hasMoreHistory: false, historyError: null, historyReadFailed: true }))).toBe(false);
});
+ it('stays away while the reads of the open are being read again', () => {
+ expect(claimsTheStart(render({ hasMoreHistory: false, historyError: null, historyRetrying: true }))).toBe(false);
+ });
+
it('stays away on a lapsed gate, where the subscription strip explains instead', () => {
subscriptionBannerUI.stateOf.mockReturnValue('expired');
const html = render({
@@ -130,6 +136,33 @@ describe('the start-of-history line', () => {
});
});
+describe('the line for edits and deletions that did not come back', () => {
+ beforeEach(() => {
+ vi.clearAllMocks();
+ subscriptionBannerUI.stateOf.mockReturnValue('active');
+ document.body.innerHTML = '
';
+ chatAreaUI.init({
+ messagesArea: document.getElementById('messages-area'),
+ messageInput: document.getElementById('message-input')
+ });
+ });
+
+ it('says they are still loading while they are read again', () => {
+ const html = render({ hasMoreHistory: true, historyError: null, overridesOwed: true, historyRetrying: true });
+ expect(html).toContain('Loading edits and deletions…');
+ });
+
+ it('says they could not be loaded once the reads gave up', () => {
+ const html = render({ hasMoreHistory: false, historyError: null, overridesOwed: true, historyReadFailed: true });
+ expect(html).toContain('Edits and deletions could not be loaded. Reopen the channel');
+ expect(html).not.toContain('Loading edits and deletions');
+ });
+
+ it('is absent when they came back', () => {
+ expect(render({ hasMoreHistory: true, historyError: null })).not.toContain('overrides-owed-banner');
+ });
+});
+
describe('the refusal banner over the messages', () => {
const refused = { hasMoreHistory: false, historyError: { status: 403, signed: true, at: Date.now() } };
diff --git a/tests/unit/PreviewModeUI.test.js b/tests/unit/PreviewModeUI.test.js
index 5d5f50b..4a0b6c5 100644
--- a/tests/unit/PreviewModeUI.test.js
+++ b/tests/unit/PreviewModeUI.test.js
@@ -87,4 +87,28 @@ describe('PreviewModeUI', () => {
expect(previewModeUI.previewChannel.messages).toEqual([]);
expect(result).toEqual({ loaded: 1, hasMore: false });
});
+
+ it('does not take a page whose overrides did not come back', async () => {
+ previewModeUI.previewChannel = {
+ streamId: 'preview-stream',
+ password: null,
+ messages: [],
+ _pendingOverrides: new Map(),
+ loadingHistory: false,
+ oldestTimestamp: 1000,
+ hasMoreHistory: true
+ };
+ deps.streamrController.fetchOlderHistory.mockImplementation(async (_streamId, partition) => (
+ partition === 0
+ ? { messages: [{ id: 'msg-1', text: 'hello', sender: '0x1', timestamp: 900 }], hasMore: false }
+ : { messages: [], hasMore: false, failed: true }
+ ));
+
+ const result = await previewModeUI.loadMorePreviewHistory();
+
+ expect(previewModeUI.previewChannel.messages).toEqual([]);
+ expect(previewModeUI.previewChannel.oldestTimestamp).toBe(1000);
+ expect(previewModeUI.previewChannel.loadingHistory).toBe(false);
+ expect(result).toEqual({ loaded: 0, hasMore: true });
+ });
});
\ No newline at end of file
diff --git a/tests/unit/channels.extended.test.js b/tests/unit/channels.extended.test.js
index 9ba5f64..112d251 100644
--- a/tests/unit/channels.extended.test.js
+++ b/tests/unit/channels.extended.test.js
@@ -1621,6 +1621,18 @@ describe('ChannelManager Extended', () => {
retry.mockRestore();
});
+ it('marks the edits and deletions owed when it was their read that failed', async () => {
+ const retry = vi.spyOn(channelManager, '_retryOpenReads').mockImplementation(() => {});
+
+ const owed = await loadWith({ readError: null, failed: true, overridesFailed: true });
+ expect(owed.overridesOwed).toBe(true);
+
+ channelManager.channels.delete(streamId);
+ const contentOnly = await loadWith({ readError: null, failed: true, overridesFailed: false });
+ expect(contentOnly.overridesOwed).toBe(false);
+ 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
index 788b3ab..95ca6fa 100644
--- a/tests/unit/channels.openReadRetry.test.js
+++ b/tests/unit/channels.openReadRetry.test.js
@@ -33,18 +33,25 @@ function open() {
describe('reading the open again after a failed read', () => {
let failing;
+ let overridesFailing;
let calls;
beforeEach(() => {
vi.useFakeTimers();
failing = true;
+ overridesFailing = false;
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 });
+ if (id !== ID) return;
+ if (partition === STREAM_CONFIG.MESSAGE_STREAM.MESSAGES) {
+ calls.push('content');
+ await done?.({ loaded: 0, requested: count, readError: null, failed: failing });
+ } else if (partition === STREAM_CONFIG.MESSAGE_STREAM.CONTROL) {
+ calls.push('overrides');
+ await done?.({ loaded: 0, requested: count, readError: null, failed: overridesFailing });
+ }
});
vi.spyOn(channelManager, 'refreshAdminState').mockImplementation(async () => { calls.push('admin'); });
for (const method of ['flushBatchVerification', 'awaitAllFlushes']) {
@@ -127,4 +134,48 @@ describe('reading the open again after a failed read', () => {
await vi.advanceTimersByTimeAsync(BACKOFF_MS[0]);
expect(calls).toEqual([]);
});
+
+ it('keeps the edits and deletions owed while their read fails, and clears it when it comes back', async () => {
+ const channel = open();
+ channel._controlPartitionSupported = true;
+ failing = false;
+ overridesFailing = true;
+
+ channelManager._retryOpenReads(ID);
+ await vi.advanceTimersByTimeAsync(BACKOFF_MS[0]);
+ expect(channel.overridesOwed).toBe(true);
+ expect(channel.historyRetrying).toBe(true);
+
+ overridesFailing = false;
+ await vi.advanceTimersByTimeAsync(BACKOFF_MS[1]);
+ expect(channel.overridesOwed).toBe(false);
+ expect(channel.historyRetrying).toBe(false);
+ });
+
+ it('gives up with the edits and deletions still owed', async () => {
+ const channel = open();
+ channel._controlPartitionSupported = true;
+ failing = false;
+ overridesFailing = true;
+
+ channelManager._retryOpenReads(ID);
+ await vi.advanceTimersByTimeAsync(BACKOFF_MS.reduce((a, b) => a + b, 0));
+
+ expect(channel.historyReadFailed).toBe(true);
+ expect(channel.overridesOwed).toBe(true);
+ });
+
+ it('reads again when a refresh left the edits and deletions owed and nothing reads them', async () => {
+ const channel = open();
+ channel._controlPartitionSupported = true;
+ channel.gate = { address: '0x' + '44'.repeat(20) };
+ failing = false;
+ overridesFailing = true;
+
+ channelManager.refreshHistory(ID);
+ await vi.advanceTimersByTimeAsync(1500);
+
+ expect(channel.overridesOwed).toBe(true);
+ expect(channel.historyRetrying).toBe(true);
+ });
});
diff --git a/tests/unit/channels.test.js b/tests/unit/channels.test.js
index 396d8fd..79beb71 100644
--- a/tests/unit/channels.test.js
+++ b/tests/unit/channels.test.js
@@ -1635,6 +1635,65 @@ describe('ChannelManager', () => {
expect(channel.messages.some(m => m.id === 'msg-verify')).toBe(false);
});
+ describe('a page taken whole or not at all', () => {
+ const older = { id: 'msg-old', timestamp: 500, text: 'older', sender: '0x2' };
+ const serveHalves = ({ contentFailed = false, overridesFailed = false }) => {
+ streamrController.fetchOlderHistory.mockImplementation(async (id, partition) => (
+ partition === 1
+ ? { messages: [], hasMore: false, failed: overridesFailed }
+ : { messages: [older], hasMore: true, failed: contentFailed }
+ ));
+ };
+ const notTaken = (result) => {
+ const channel = channelManager.channels.get(streamId);
+ expect(result.loaded).toBe(0);
+ expect(channel.messages.some(m => m.id === 'msg-old')).toBe(false);
+ expect(channel.hasMoreHistory).toBe(true);
+ expect(channel.oldestTimestamp).toBe(1000);
+ expect(channel.loadingHistory).toBe(false);
+ };
+
+ it('does not take a page whose overrides did not come back', async () => {
+ serveHalves({ overridesFailed: true });
+ notTaken(await channelManager.loadMoreHistory(streamId));
+ });
+
+ it('does not take a page whose content did not come back', async () => {
+ serveHalves({ contentFailed: true });
+ notTaken(await channelManager.loadMoreHistory(streamId));
+ });
+
+ it('waits while the open is being read again', async () => {
+ channelManager.channels.get(streamId).historyRetrying = true;
+ serveHalves({});
+
+ notTaken(await channelManager.loadMoreHistory(streamId));
+ expect(streamrController.fetchOlderHistory).not.toHaveBeenCalled();
+ });
+
+ it('asks again only after a backoff, and then says the page is due', async () => {
+ vi.useFakeTimers();
+ try {
+ const notify = vi.spyOn(channelManager, 'notifyHandlers');
+ serveHalves({ overridesFailed: true });
+ await channelManager.loadMoreHistory(streamId);
+ const reads = streamrController.fetchOlderHistory.mock.calls.length;
+
+ await channelManager.loadMoreHistory(streamId);
+ expect(streamrController.fetchOlderHistory.mock.calls.length).toBe(reads);
+
+ await vi.advanceTimersByTimeAsync(5_000);
+ expect(notify).toHaveBeenCalledWith('history_page_due', { streamId });
+
+ serveHalves({});
+ const result = await channelManager.loadMoreHistory(streamId);
+ expect(result.loaded).toBe(1);
+ } finally {
+ vi.useRealTimers();
+ }
+ });
+ });
+
it('should pass abort signal to fetchOlderHistory', async () => {
streamrController.fetchOlderHistory.mockResolvedValue({ messages: [], hasMore: false });
diff --git a/tests/unit/streamr.core.test.js b/tests/unit/streamr.core.test.js
index e928ecb..d0834e7 100644
--- a/tests/unit/streamr.core.test.js
+++ b/tests/unit/streamr.core.test.js
@@ -1290,10 +1290,10 @@ describe('StreamrController Core', () => {
expect(cryptoManager.decryptJSON).toHaveBeenCalled();
});
- it('should return empty on fetch error', async () => {
+ it('should return an empty failed page on fetch error', async () => {
mockClient.resend.mockRejectedValue(new Error('fail'));
const result = await streamrController.fetchOlderHistory('stream-1', 0, 1000);
- expect(result).toEqual({ messages: [], hasMore: false });
+ expect(result).toEqual({ messages: [], hasMore: false, failed: true });
});
it('should confirm exhaustion with a second pass and merge truncated results', async () => {
diff --git a/tests/unit/streamr.dualStream.test.js b/tests/unit/streamr.dualStream.test.js
index f804d4a..0189fcd 100644
--- a/tests/unit/streamr.dualStream.test.js
+++ b/tests/unit/streamr.dualStream.test.js
@@ -112,21 +112,35 @@ describe('subscribeToDualStream', () => {
expect(done).toHaveBeenCalledTimes(1);
expect(done).toHaveBeenCalledWith({
contentLoaded: 12, contentRequested: 30, controlLoaded: 3, controlRequested: 30, readError: null,
- failed: false,
+ failed: false, overridesFailed: false,
});
});
- it('says the history read failed when either stored partition failed', async () => {
+ it('says the history read failed when either stored partition failed, and whether it was the overrides', async () => {
const calls = stubSubscribe();
const done = vi.fn();
await streamrController.subscribeToDualStream(
MESSAGE, EPHEMERAL, { onMessage: () => {}, onOverride: () => {} }, null, 30, done
);
+ // calls[0] is the control partition, calls[1] the content one.
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 }));
+ expect(done).toHaveBeenCalledWith(expect.objectContaining({ failed: true, overridesFailed: true }));
+ });
+
+ it('does not blame the overrides for a content read that failed', async () => {
+ const calls = stubSubscribe();
+ const done = vi.fn();
+
+ await streamrController.subscribeToDualStream(
+ MESSAGE, EPHEMERAL, { onMessage: () => {}, onOverride: () => {} }, null, 30, done
+ );
+ await calls[0].onDone({ loaded: 3, requested: 30 });
+ await calls[1].onDone({ loaded: 0, requested: 30, failed: true });
+
+ expect(done).toHaveBeenCalledWith(expect.objectContaining({ failed: true, overridesFailed: false }));
});
it('signals completion once even if a partition reports twice', async () => {
diff --git a/tests/unit/streamr.olderHistory.test.js b/tests/unit/streamr.olderHistory.test.js
index e51e8c9..cd0e063 100644
--- a/tests/unit/streamr.olderHistory.test.js
+++ b/tests/unit/streamr.olderHistory.test.js
@@ -23,6 +23,7 @@ vi.mock('../../src/js/envelopeSigner.js', async (importOriginal) => ({
}));
const { streamrController } = await import('../../src/js/streamr.js');
+const { storageFetch } = await import('../../src/js/storageFetch.js');
const AUTHOR = '0x' + '11'.repeat(20);
const STREAM = '0xaaa/older-history-1';
@@ -35,10 +36,16 @@ const text = (ts, over = {}) => ({
publisherId: AUTHOR,
});
+/** An Error in the list is thrown there, the way the SDK's iterator surfaces a node that failed. */
async function* streamOf(...messages) {
- for (const m of messages) yield m;
+ for (const m of messages) {
+ if (m instanceof Error) throw m;
+ yield m;
+ }
}
+const storageNodeError = () => Object.assign(new Error('Failed to fetch'), { code: 'STORAGE_NODE_ERROR' });
+
/** Returns the resend calls it recorded, and serves `pages` one call at a time. */
function serve(...pages) {
const calls = [];
@@ -228,12 +235,55 @@ describe('fetchOlderHistory', () => {
expect(out.messages.map((m) => m.id)).toEqual(['id-100', 'id-200']);
});
- it('answers with an empty page when the resend itself fails', async () => {
+ it('answers with a failed page when the resend itself fails, not with the end of history', async () => {
streamrController.client = { resend: vi.fn(async () => { throw new Error('storage down'); }) };
const out = await streamrController.fetchOlderHistory(STREAM, P_MESSAGES, 1000, 10);
- expect(out).toEqual({ messages: [], hasMore: false });
+ expect(out).toEqual({ messages: [], hasMore: false, failed: true });
+ });
+
+ it('takes a stream with no storage as answered', async () => {
+ const noStorage = Object.assign(new Error(`no storage assigned: ${STREAM}`), { code: 'NO_STORAGE_NODES' });
+ streamrController.client = { resend: vi.fn(async () => { throw noStorage; }) };
+
+ const out = await streamrController.fetchOlderHistory(STREAM, P_MESSAGES, 1000, 10);
+
+ expect(out.failed).toBe(false);
+ });
+
+ it('says a page the storage node broke off failed', async () => {
+ serve([text(100), storageNodeError()]);
+
+ const out = await streamrController.fetchOlderHistory(STREAM, P_MESSAGES, 1000, 10);
+
+ expect(out.failed).toBe(true);
+ });
+
+ it('lets a confirmation pass that read the range whole stand for a broken first pass', async () => {
+ serve([text(100), storageNodeError()], [text(90), text(100)]);
+
+ const out = await streamrController.fetchOlderHistory(STREAM, P_MESSAGES, 1000, 10);
+
+ expect(out.failed).toBe(false);
+ expect(out.messages.map((m) => m.id)).toEqual(['id-90', 'id-100']);
+ });
+
+ it('takes a 4xx from the storage node as its answer', async () => {
+ vi.spyOn(storageFetch, 'lastReadError').mockReturnValue({ status: 403, signed: true });
+ serve([storageNodeError()]);
+
+ const out = await streamrController.fetchOlderHistory(STREAM, P_MESSAGES, 1000, 10);
+
+ expect(out.failed).toBe(false);
+ });
+
+ it('does not fail a page over a row it cannot decrypt', async () => {
+ serve([text(100), Object.assign(new Error('no encryption key'), { code: 'DECRYPT_ERROR' })]);
+
+ const out = await streamrController.fetchOlderHistory(STREAM, P_MESSAGES, 1000, 10);
+
+ expect(out.failed).toBe(false);
});
it('stamps the transport timestamp and publisher onto the row', async () => {