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
807 changes: 807 additions & 0 deletions docs/ADMIN-chunk-vectors.json

Large diffs are not rendered by default.

129 changes: 105 additions & 24 deletions src/js/channels/AdminState.js
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,12 @@
*/

import { Logger } from '../logger.js';
import { CONFIG } from '../config.js';
import { cryptoManager } from '../crypto.js';
import { streamrController, STREAM_CONFIG, deriveEphemeralId, deriveAdminId } from '../streamr.js';
import { authManager } from '../auth.js';
import { adminStatePoller } from '../adminStatePoller.js';
import { splitFramed, ADMIN_FRAME, SYNC_CHUNK_CHARS } from '../syncChunks.js';

export class AdminState {
/**
Expand All @@ -24,6 +27,7 @@ export class AdminState {
// adminRev and publish colliding revs — latest-wins would then
// silently drop one of the operations.
this._adminPublishChain = new Map(); // messageStreamId -> Promise
this._signalReads = new Map(); // messageStreamId -> pending -3 read
}

/**
Expand Down Expand Up @@ -277,6 +281,17 @@ export class AdminState {
if (!adminStreamId) return false;

try {
// The newest row alone answers most polls: the owner's snapshot or
// manifest at a rev already held means nothing changed. Anything
// else reads the window as before.
const owner = (channel.createdBy || messageStreamId.split('/')[0] || '').toLowerCase();
const newest = await streamrController.probeAdminState(adminStreamId, {
password: channel.password || null
});
if (newest && owner && String(newest.publisherId).toLowerCase() === owner
&& newest.rev <= (channel.adminRev || 0)) {
return false;
}
const latest = await streamrController.resendAdminState(adminStreamId, {
// Smaller window for cheap polling — only the most recent
// snapshot wins regardless of how many entries we fetch.
Expand Down Expand Up @@ -382,7 +397,16 @@ export class AdminState {
state: next
};

const published = await streamrController.publishAdminState(adminStreamId, adminMsg, channel.password || null);
const password = channel.password || null;
const rows = await this._frameForWire(adminStreamId, adminMsg, password);
const messages = [];
for (const row of rows) {
messages.push(await streamrController.publishAdminState(adminStreamId, row, password));
}
const published = messages.at(-1);
if (rows.length > 1) {
Logger.info(`ADMIN_STATE rev ${newRev} split in ${rows.length - 1} chunks`);
}

// Optimistically apply locally so UI reflects the change immediately.
this.manager.applyAdminState(channel, adminMsg);
Expand All @@ -409,32 +433,89 @@ export class AdminState {
// Fire-and-forget invalidation signal on the ephemeral -2/P0 control
// partition so other clients with the channel active update immediately
// (instead of waiting for the next 30s poller tick). The signal embeds
// the full ADMIN_STATE snapshot so receivers can apply it inline
// without a -3/P0 resend round-trip. The canonical resend path on
// -3/P0 (bootstrap-on-open + periodic poll + on-demand fallback)
// remains as a convergence safety net for clients that miss the
// ephemeral signal. Best-effort: failure here is non-fatal.
try {
const ephemeralStreamId = channel.ephemeralStreamId || deriveEphemeralId(messageStreamId);
if (ephemeralStreamId) {
const signal = {
type: 'admin_invalidate',
rev: newRev,
ts: adminMsg.ts,
snapshot: adminMsg
};
streamrController.publishControl(
ephemeralStreamId,
signal,
channel.password || null
).catch(e => Logger.debug('admin_invalidate publish failed (non-fatal):', e.message));
}
} catch (e) {
Logger.debug('admin_invalidate prepare failed (non-fatal):', e.message);
// the full ADMIN_STATE snapshot when it fits one message, so receivers
// apply it inline; otherwise it carries only the rev and receivers
// read the -3. The canonical resend path on -3/P0 (bootstrap-on-open
// + periodic poll + on-demand fallback) remains as a convergence
// safety net for clients that miss the ephemeral signal. Best-effort:
// failure here is non-fatal.
const ephemeralStreamId = channel.ephemeralStreamId || deriveEphemeralId(messageStreamId);
if (ephemeralStreamId) {
this._signal(ephemeralStreamId, adminMsg, rows.length === 1, password)
.catch(e => Logger.debug('admin_invalidate publish failed (non-fatal):', e.message));
}

Logger.info('Published ADMIN_STATE rev', newRev, 'for', messageStreamId.slice(-20));
return { rev: newRev, state: next, published };
return { rev: newRev, state: next, published, messages };
}

/**
* The rows this snapshot goes out as: itself when it fits the wire once
* encrypted, as it always did, else a run of chunks closed by a manifest,
* each cut to fit. Past `adminStateMaxChunks` nothing goes out, and the
* owner is told instead of the network dropping it unseen.
* @private
*/
async _frameForWire(adminStreamId, adminMsg, password) {
const budget = CONFIG.media.imagePayloadMaxBytes - CONFIG.media.imagePayloadSafetyMarginBytes;
const measure = (row) => streamrController.adminWireBytes(adminStreamId, row, password);
if (await measure(adminMsg) <= budget) return [adminMsg];

const runId = cryptoManager.generateRandomHex(8);
for (let limit = SYNC_CHUNK_CHARS; limit >= 1024; limit = Math.floor(limit * 0.85)) {
const run = splitFramed(adminMsg, runId, ADMIN_FRAME, limit);
if (run.length - 1 > CONFIG.subscriptions.adminStateMaxChunks) break;
let fits = true;
for (const row of run) {
if (await measure(row) > budget) { fits = false; break; }
}
if (fits) return run;
}
const error = new Error("This channel's moderation state is too large to publish. Unpin some messages and try again.");
error.code = 'ADMIN_STATE_TOO_LARGE';
throw error;
}

/** @private */
async _signal(ephemeralStreamId, adminMsg, whole, password) {
const signal = { type: 'admin_invalidate', rev: adminMsg.rev, ts: adminMsg.ts };
if (whole) {
const full = { ...signal, snapshot: adminMsg };
const budget = CONFIG.media.imagePayloadMaxBytes - CONFIG.media.imagePayloadSafetyMarginBytes;
const bytes = await streamrController.channelWireBytes(ephemeralStreamId, full, password)
.catch(() => Infinity);
if (bytes <= budget) return streamrController.publishControl(ephemeralStreamId, full, password);
}
return streamrController.publishControl(ephemeralStreamId, signal, password);
}

/**
* An admin_invalidate announced a snapshot too big to ride along: poll
* the -3 once storage has had time to hold it, and once more later while
* the announced rev has still not landed. Signals arriving while reads
* are pending share them.
* @param {string} messageStreamId - Channel key (-1)
* @param {number} rev - The rev the signal announced
*/
readAfterSignal(messageStreamId, rev) {
if (this._signalReads.has(messageStreamId)) return;
const landed = () => {
const channel = this.manager.channels.get(messageStreamId)
|| (this.manager.previewChannel?.messageStreamId === messageStreamId
? this.manager.previewChannel : null);
return (channel?.adminRev || 0) >= rev;
};
const [first, ...later] = CONFIG.subscriptions.adminSignalReadDelaysMs;
const read = (wait, rest) => this._signalReads.set(messageStreamId, setTimeout(() => {
if (landed() || adminStatePoller.getStreamId() !== messageStreamId) {
this._signalReads.delete(messageStreamId);
return;
}
adminStatePoller.pollNow();
if (rest.length) read(rest[0], rest.slice(1));
else this._signalReads.delete(messageStreamId);
}, wait));
read(first, later);
}

// High-level convenience helpers built on top of publishAdminState ----------
Expand Down
6 changes: 6 additions & 0 deletions src/js/channels/AdminStateConfirm.js
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,12 @@ export class AdminStateConfirm {
state: channel.adminSnapshot || channel.adminState || {}
});
} catch (e) {
if (e?.code === 'ADMIN_STATE_TOO_LARGE') {
this._setPending(messageStreamId, null);
Logger.warn(`${label} cannot be republished: ${e.message}`);
this.manager.notifyHandlers('admin_state_too_large', { streamId: messageStreamId, rev: pending.rev });
return;
}
Logger.warn(`${label} republish failed, kept pending:`, e?.message);
return;
}
Expand Down
4 changes: 4 additions & 0 deletions src/js/channels/MessageFlow.js
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,10 @@ export class MessageFlow {
}
const incomingRev = typeof data.rev === 'number' ? data.rev : 0;
if (incomingRev <= (channel.adminRev || 0)) return;
if (data.snapshot === undefined) {
this.manager.adminState?.readAfterSignal(streamId, incomingRev);
return;
}
if (!this.manager._isValidAdminState(data.snapshot)) return;

// Sender authenticity is already validated above (account ===
Expand Down
4 changes: 2 additions & 2 deletions src/js/channels/StorageCopy.js
Original file line number Diff line number Diff line change
Expand Up @@ -284,10 +284,10 @@ export class StorageCopy {
}
if (items.has('admin') && (channel.adminRev || 0) > 0) {
await attempt('admin', async () => {
const { published } = await this.manager.publishAdminState(channel.messageStreamId, {
const { messages } = await this.manager.publishAdminState(channel.messageStreamId, {
state: channel.adminSnapshot || channel.adminState
});
return [row('admin', adminStreamId, ADMIN.MODERATION, published)];
return messages.map((m) => row('admin', adminStreamId, ADMIN.MODERATION, m));
});
}
if (items.has('image') && snapshot?.image && !(snapshot.image.encrypted && !password)) {
Expand Down
9 changes: 9 additions & 0 deletions src/js/config.js
Original file line number Diff line number Diff line change
Expand Up @@ -293,6 +293,15 @@ export const CONFIG = {
// times before the owner is told.
adminConfirmDelaysMs: [5000, 10000, 20000, 40000],
adminConfirmRepublishLimit: 3,
// An ADMIN_STATE too big for one message goes out as a run of this
// many chunks at most; readers reassemble runs up to the second
// bound. The first must stay below the smallest read window (5).
adminStateMaxChunks: 4,
adminReadMaxChunks: 64,
// A snapshot-less admin_invalidate is followed by a -3 read after the
// first wait, and by one more after the second while the announced
// rev has still not landed.
adminSignalReadDelaysMs: [5000, 10000],
// A provider just added is asked for the copy at these delays; the copy
// is published again this many times while still missing there.
storageCopyDelaysMs: [5000, 10000, 20000, 40000],
Expand Down
11 changes: 11 additions & 0 deletions src/js/crypto.js
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,17 @@ class CryptoManager {
}
}

/**
* Length of what encrypt() returns for a plaintext of this many UTF-8
* bytes: base64 over salt, iv, ciphertext and GCM tag. Deterministic, so
* sizing a payload does not pay a key derivation.
* @param {number} plaintextBytes
* @returns {number}
*/
encryptedLength(plaintextBytes) {
return 4 * Math.ceil((16 + 12 + plaintextBytes + 16) / 3);
}

/**
* Encrypt JSON object
* @param {Object} obj - Object to encrypt
Expand Down
Loading
Loading