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
70 changes: 63 additions & 7 deletions src/js/channels.js
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -2078,14 +2082,15 @@ 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,
STREAM_CONFIG.INTERACTIONS_STREAM.REACTIONS,
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));
Expand Down Expand Up @@ -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;

Expand Down Expand Up @@ -2433,26 +2439,37 @@ 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<boolean>} 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
// flush verifications, apply pending overrides (which also prunes
// deleted messages), and render ONCE.
channel.initialLoadInProgress = true;
let contentRead = null;
let controlRead = null;
let reactionsRead = null;
try {
await streamrController.fetchHistoryAsync(
messageStreamId,
Expand All @@ -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 }
);
Expand All @@ -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));
Expand Down Expand Up @@ -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 ====================
Expand Down
7 changes: 5 additions & 2 deletions src/js/streamr.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand All @@ -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); }
}
Expand All @@ -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}`);
Expand Down
11 changes: 10 additions & 1 deletion src/js/streamr/History.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)' : ''}...`);

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

Expand Down Expand Up @@ -815,19 +818,25 @@ 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
// exists in storage). Used by `channels.onHistoryComplete` to flip
// `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); }
}
}
Expand Down
14 changes: 11 additions & 3 deletions src/js/ui/ChatAreaUI.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 &&
Expand Down Expand Up @@ -690,6 +691,13 @@ class ChatAreaUI {
`;
this.messagesArea.querySelector('#empty-state-renew-btn')
?.addEventListener('click', () => subscriptionBannerUI.renewCurrent());
} else if (effectiveChannel?.historyReadFailed) {
this.messagesArea.innerHTML = `
<div class="flex flex-col items-center justify-center h-full text-white/40 gap-3">
<span class="text-sm">Channel history could not be loaded</span>
<span class="text-xs text-white/25">Check your connection and reopen the channel</span>
</div>
`;
} else {
this.messagesArea.innerHTML = waitingForKeys
? `
Expand Down Expand Up @@ -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 = `
<div class="flex justify-center py-4">
<div class="text-white/[0.12] text-xs">
Expand Down
9 changes: 7 additions & 2 deletions tests/unit/ChatAreaUI.historyStart.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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({
Expand Down
12 changes: 12 additions & 0 deletions tests/unit/channels.extended.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading
Loading