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
8 changes: 6 additions & 2 deletions index.html
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,10 @@
<!-- Left side: Logo and brand -->
<div class="flex items-center gap-3">
<!-- Logo -->
<img src="favicon/favicon.svg" alt="Pombo" class="w-12 h-12 header-logo">
<span class="header-logo-wrap">
<img src="favicon/favicon.svg" alt="Pombo" class="w-12 h-12 header-logo">
<span class="network-dot" role="status" aria-label="Connecting to the network" title="Connecting…"></span>
</span>
<!-- Pombo name -->
<div class="flex flex-col">
<span class="text-[17px] font-extrabold text-white tracking-tight leading-tight">Pombo</span>
Expand Down Expand Up @@ -453,7 +456,8 @@ <h2 id="current-channel-name" class="text-lg font-semibold truncate">Explore Cha
<div class="relative inline-flex items-center">
<button id="online-header" class="hidden items-center gap-1 hover:text-white/70 transition">
<span id="online-users-count">0</span>
<span>Online</span>
<span class="online-word">Online</span>
<span class="online-connecting">Connecting…</span>
</button>
<div id="online-users-list" class="absolute top-full left-0 mt-1 w-56 bg-[#1e1e1e] border border-white/10 rounded-xl p-3 z-20 hidden max-h-48 overflow-y-auto">
<div class="text-white/30 text-sm text-center">No one online</div>
Expand Down
48 changes: 44 additions & 4 deletions src/js/app.js
Original file line number Diff line number Diff line change
Expand Up @@ -94,12 +94,12 @@ class App {
secureStorage.onSentDataChanged = () => syncManager.scheduleAutoPush(5000);

settingsUI.setDependencies({
connectWallet: () => walletFlows.connectWallet(),
onSyncTransportReconnected: () => this.pullSyncedStateAndBlobs().catch((error) => {
Logger.debug('Sync: Auto-pull after reconnect failed (non-critical):', error.message);
})
connectWallet: () => walletFlows.connectWallet()
});

streamrController.onClientReplaced(() => this.resumeAfterClientReplaced());
streamrController.onNodeStateChange((up) => headerUI.setNetworkDown(!up));

// Wire wallet flows with app-level callbacks
walletFlows.init({
disconnectWallet: (opts) => this.disconnectWallet(opts),
Expand Down Expand Up @@ -266,13 +266,48 @@ class App {
syncManager.stopSnapshotWatch();
pushOnHide();
} else if (document.visibilityState === 'visible') {
streamrController.revival.kick();
if (!syncManager.isAutoSyncAllowed('foreground')) return;
syncManager.startSnapshotWatch();
syncManager.checkForNewSnapshot().catch((error) => {
Logger.debug('Sync: Check on foreground failed (non-critical):', error.message);
});
}
});

window.addEventListener('online', () => streamrController.revival.kick());
}

/**
* A replaced Streamr client took every subscription with it: make them
* again, and retry what only the network can carry (owed rotations, the
* inbox push registration, sync).
*/
async resumeAfterClientReplaced() {
if (!authManager.isConnected()) return;
const steps = [
['DM inbox', () => dmManager.resubscribeInbox()],
['active channel', () => subscriptionManager.resubscribeActive()],
['preview', async () => {
const preview = previewModeUI.isInPreviewMode() ? previewModeUI.getPreviewChannel() : null;
if (preview) await previewModeUI.enterPreviewWithoutHistory(preview.streamId, preview.channelInfo);
}],
['owed rotations', async () => {
if (!authManager.isGuestMode()) channelManager.resumeOwedRotations();
}],
['inbox push', () => dmManager.subscribeInboxPush()],
['sync', async () => {
if (syncManager.isAutoSyncAllowed('foreground')) await syncManager.smartSync();
}]
];
for (const [label, step] of steps) {
try {
await step();
} catch (error) {
Logger.warn(`Resume after client replacement (${label}) failed:`, error?.message || error);
}
}
Logger.info('Subscriptions resumed on the new Streamr client');
}

/**
Expand Down Expand Up @@ -379,6 +414,7 @@ class App {

headerUI.updateWalletInfo(null);
headerUI.updateNetworkStatus('Disconnected', false);
headerUI.setNetworkDown(false);
uiController.renderChannelList();
uiController.resetToDisconnectedState();

Expand All @@ -393,6 +429,7 @@ class App {
Logger.error('Error disconnecting:', error);
headerUI.updateWalletInfo(null);
headerUI.updateNetworkStatus('Disconnected', false);
headerUI.setNetworkDown(false);
uiController.resetToDisconnectedState();
uiController.showNotification('Error disconnecting: ' + error.message, 'error');
}
Expand Down Expand Up @@ -434,6 +471,7 @@ class App {

headerUI.updateWalletInfo(address, isGuest);
headerUI.updateNetworkStatus('Connecting to Streamr...', false);
headerUI.setNetworkDown(true);

// Bell menu (invites + active transfers): shown the moment the
// connected UI paints. It used to init at the END of this flow,
Expand Down Expand Up @@ -496,6 +534,7 @@ class App {
const streamrAddress = await streamrController.getAddress();
Logger.info('Streamr connected with address:', streamrAddress);
headerUI.updateNetworkStatus('Connected to Streamr', true);
headerUI.setNetworkDown(false);

try {
await identityManager.init();
Expand Down Expand Up @@ -659,6 +698,7 @@ class App {
} catch (error) {
Logger.error('Failed to initialize after wallet connection:', error);
headerUI.updateNetworkStatus('Failed to connect to Streamr', false);
headerUI.setNetworkDown(true);
throw error;
}
}
Expand Down
12 changes: 12 additions & 0 deletions src/js/dm.js
Original file line number Diff line number Diff line change
Expand Up @@ -297,6 +297,18 @@ class DMManager {
}
}

/**
* Subscribe the inbox again on a replaced client, whose predecessor took
* these handles with it. The open conversation's ephemeral comes back
* with the active channel.
*/
async resubscribeInbox() {
this.inboxSubscription = null;
this.inboxNotificationSub = null;
this.inboxEphemeralSubscription = null;
if (this.inboxMessageStreamId) await this.subscribeToInbox();
}

/**
* Unsubscribe from inbox (when leaving DM tab or disconnecting)
*/
Expand Down
141 changes: 119 additions & 22 deletions src/js/streamr.js
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import { CONFIG, getRpcEndpoints } from './config.js';
import { executeWithRetry, executeWithRetryAndVerify } from './utils/retry.js';
import { isRpcError, createPermissionResult } from './utils/rpcErrors.js';
import { authManager } from './auth.js';
import { NodeRevival } from './streamr/NodeRevival.js';
import {
recoverPublisherAccount, applyAccount, stripLocalFields, dropLocalState, clearPublisherProofCache
} from './publisherProof.js';
Expand Down Expand Up @@ -121,6 +122,14 @@ class StreamrController {
this.messages = new MessagePipeline(this);
this._writers = new Map(); // streamId -> { public, writers:Set, ts }
this._writerFetches = new Map(); // streamId -> in-flight promise
this._clientReplacedHandlers = [];
this._nodeStateHandlers = [];
this._replacing = Promise.resolve();
this._session = 0;
this.revival = new NodeRevival({
networkUp: () => this._networkReachable(),
rebuild: () => this.replaceClient()
});
}

/**
Expand Down Expand Up @@ -233,8 +242,10 @@ class StreamrController {
/**
* Initialize Streamr client with signer
* @param {Object} signer - Ethers signer from wallet (must have privateKey)
* @param {Object} [options]
* @param {boolean} [options.resubscribe] - run the onClientReplaced handlers once the node is up
*/
async init(signer) {
async init(signer, { resubscribe = false } = {}) {
try {
// Get StreamrClient from the global window object (exposed by the
// self-hosted vendor bundle — see src/streamr-bundle.js)
Expand All @@ -246,7 +257,7 @@ class StreamrController {
throw new Error('Signer must have a privateKey');
}

this.client = new StreamrClient({
const client = new StreamrClient({
auth: {
privateKey: signer.privateKey
},
Expand Down Expand Up @@ -283,9 +294,11 @@ class StreamrController {
rpcQuorum: 1
}
});

this.address = await this.client.getAddress();
this.client = client;

this.address = await client.getAddress();
Logger.info('Streamr client initialized with address:', this.address);
this._watchNode(client, resubscribe);

// Account identity for publishAs-based ACCOUNT publishes (keys
// stream). The -4 stream must publish as the account — its grant is
Expand All @@ -308,6 +321,98 @@ class StreamrController {
}
}

/**
* The SDK keeps a failed node start for the client's whole life, so the
* start is watched and a failure hands the client to the revival.
*/
_watchNode(client, resubscribe = false) {
client.getNodeId().then(() => {
if (this.client !== client) return;
this.revival.onAlive();
this._nodeStateHandlers.forEach((handler) => handler(true));
if (resubscribe) this._resubscribe(client);
}, (error) => {
if (this.client !== client) return;
Logger.warn('Streamr node failed to start:', error?.message || error);
this.revival.onDead();
this._nodeStateHandlers.forEach((handler) => handler(false));
});
}

/** Called with true when a client's node comes up, false when its start fails. */
onNodeStateChange(handler) {
this._nodeStateHandlers.push(handler);
}

/** One JSON-RPC round trip to any of the chosen endpoints. */
async _networkReachable() {
const endpoints = getRpcEndpoints().slice(0, 3);
const probe = ({ url }) => fetch(url, {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ jsonrpc: '2.0', id: 1, method: 'eth_chainId', params: [] }),
signal: AbortSignal.timeout(5000)
}).then((response) => {
if (!response.ok) throw new Error(`HTTP ${response.status}`);
});
try {
await Promise.any(endpoints.map(probe));
return true;
} catch {
return false;
}
}

/** Called after every client replacement, to subscribe again. */
onClientReplaced(handler) {
this._clientReplacedHandlers.push(handler);
}

/**
* Before the node is up a subscribe only waits on its start, and after a
* failed start there is nothing to subscribe on: so this runs once it is up.
*/
async _resubscribe(client) {
for (const handler of this._clientReplacedHandlers) {
if (this.client !== client) return;
try {
await handler();
} catch (e) {
Logger.warn('Resubscribe after client replacement failed:', e?.message || e);
}
}
}

/**
* A new SDK client for this same session: the pseudonyms and caches of the
* session stay, the subscriptions go with the old client and every
* onClientReplaced handler makes its own again once the new node is up.
*/
replaceClient() {
const run = this._replacing.catch(() => {}).then(() => this._replaceClientNow());
this._replacing = run;
return run;
}

async _replaceClientNow() {
const session = this._session;
const signer = authManager.getSigner();
if (!signer) throw new Error('No signer to rebuild the Streamr client with');
const old = this.client;
this.client = null;
this.subscriptions.clear();
if (old) {
// A stop that hangs must not hold the new client back.
await Promise.race([
old.destroy().catch((e) => Logger.warn('Streamr client destroy error (ignored):', e.message)),
new Promise((resolve) => setTimeout(resolve, 5000))
]);
}
// Logged out or switched account meanwhile: the next login makes its own client.
if (session !== this._session) return;
await this.init(signer, { resubscribe: true });
}

/**
* Pre-warm the network node for faster first subscription.
* Subscribes briefly to a stream to force node startup.
Expand Down Expand Up @@ -3939,7 +4044,14 @@ class StreamrController {
* Disconnect Streamr client
*/
async disconnect() {
if (this.client) {
this._session++;
this.revival.stop();
const client = this.client;
if (client) {
// Out before anything is awaited: a node start still pending settles
// during the teardown, and its watch must not see this client as current.
this.client = null;

// Unsubscribe from all streams
for (const streamId of this.subscriptions.keys()) {
try {
Expand All @@ -3950,11 +4062,10 @@ class StreamrController {
}

try {
await this.client.destroy();
await client.destroy();
} catch (e) {
Logger.warn('Streamr client destroy error (ignored):', e.message);
}
this.client = null;
Logger.info('Streamr client disconnected');
}

Expand All @@ -3973,26 +4084,12 @@ class StreamrController {

/**
* Reconnect Streamr client with new RPC endpoints
* Gets signer from authManager (avoids storing private key)
* Note: Active subscriptions are cleared - user needs to rejoin channels
* @returns {Promise<boolean>} - Success status
*/
async reconnect() {
const signer = authManager.getSigner();
if (!signer) {
Logger.warn('Cannot reconnect: no signer available from authManager');
return false;
}

try {
Logger.info('Reconnecting Streamr client with new RPC endpoints...');

// Disconnect current client (clears subscriptions)
await this.disconnect();

// Re-initialize with fresh signer (will use new RPC endpoints)
await this.init(signer);

await this.replaceClient();
Logger.info('Streamr client reconnected successfully');
return true;
} catch (error) {
Expand Down
Loading
Loading