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
2 changes: 2 additions & 0 deletions src/js/app.js
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
16 changes: 14 additions & 2 deletions src/js/channels.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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);
}

Expand Down Expand Up @@ -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.
Expand Down
45 changes: 44 additions & 1 deletion src/js/channels/MessageFlow.js
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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`
);
Expand Down
1 change: 1 addition & 0 deletions src/js/streamr.js
Original file line number Diff line number Diff line change
Expand Up @@ -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); }
}
Expand Down
14 changes: 12 additions & 2 deletions src/js/streamr/History.js
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -267,6 +269,7 @@ export class History {
continue;
}
Logger.warn('fetchOlderHistory iteration error:', iterError.message);
lastPassBroke = true;
continue;
}

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

Expand Down
19 changes: 15 additions & 4 deletions src/js/ui/ChatAreaUI.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 = `
<div class="flex justify-center py-4">
<div class="text-white/[0.12] text-xs">
Expand All @@ -757,6 +758,16 @@ class ChatAreaUI {
</div>
`;
}
let overridesBanner = '';
if (effectiveChannel?.overridesOwed) {
overridesBanner = `
<div id="overrides-owed-banner" class="flex flex-col items-center gap-1 py-3 px-4 text-center">
<span class="text-sm text-white/40">${effectiveChannel.historyReadFailed
? 'Edits and deletions could not be loaded. Reopen the channel'
: 'Loading edits and deletions…'}</span>
</div>
`;
}

// Analyze message groups for Stack Effect
const groupPositions = analyzeMessageGroups(messagesForRender);
Expand Down Expand Up @@ -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;
Expand Down
6 changes: 5 additions & 1 deletion src/js/ui/PreviewModeUI.js
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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`
);
Expand Down
37 changes: 35 additions & 2 deletions tests/unit/ChatAreaUI.historyStart.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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({
Expand All @@ -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 = '<div id="messages-area"></div><div id="message-input" contenteditable="true"></div>';
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() } };

Expand Down
24 changes: 24 additions & 0 deletions tests/unit/PreviewModeUI.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 });
});
});
12 changes: 12 additions & 0 deletions tests/unit/channels.extended.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading
Loading