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
60 changes: 0 additions & 60 deletions src/js/channels.js
Original file line number Diff line number Diff line change
Expand Up @@ -1401,66 +1401,6 @@ class ChannelManager {
}
}

/**
* Sync channel info from The Graph (members, type, owner)
* Use this to refresh on-chain data for a channel
* @param {string} messageStreamId - Message Stream ID
* @returns {Promise<Object|null>} - Updated channel or null if not found
*/
async syncChannelFromGraph(messageStreamId) {
const channel = this.channels.get(messageStreamId);
if (!channel) {
Logger.warn('Channel not found for sync:', messageStreamId);
return null;
}

try {
Logger.debug('Syncing channel from The Graph:', messageStreamId);

// OPTIMIZATION: Fetch all Graph data in parallel
const [streamData, type, membersResult] = await Promise.all([
graphAPI.getStream(messageStreamId),
graphAPI.detectStreamType(messageStreamId),
graphAPI.getStreamMembers(messageStreamId)
]);

if (!streamData) {
Logger.warn('Stream not found in The Graph');
return channel;
}

// Update type based on permissions
if (type !== 'unknown') {
channel.type = type;
}

const members = membersResult.ok ? membersResult.data : [];

// Update members
channel.members = members.map(m => m.address);

// Get owner
const owner = members.find(m => m.isOwner);
if (owner) {
channel.createdBy = owner.address;
}

// Save updated info
await this.saveChannels();

Logger.debug('Channel synced:', {
type: channel.type,
members: channel.members.length,
owner: channel.createdBy?.slice(0, 10)
});

return channel;
} catch (error) {
Logger.warn('Failed to sync channel from Graph:', error);
return channel;
}
}

// Lives in channels/Membership.js; the manager keeps the entry points
// its callers already use.

Expand Down
4 changes: 2 additions & 2 deletions src/js/channels/Membership.js
Original file line number Diff line number Diff line change
Expand Up @@ -224,7 +224,7 @@ export class Membership {
*/
/**
* Candidate membership answered by the gate: the local cache, the
* KEY_REQUEST authors seen on -4, the -4/P1 roster and, on Closed gates,
* KEY_REQUEST authors seen on -4, the -4/P2 roster and, on Closed gates,
* the contract's own enumeration — which makes the candidate set complete
* there instead of limited to what this client happened to see. Empty on
* failure — each caller picks its own fallback.
Expand Down Expand Up @@ -322,7 +322,7 @@ export class Membership {
const gateAddr = channel.gate.address.toLowerCase();
try {
const { gateManager } = await import('../gate.js');
// Roster (-4/P1) is the persistent, device-independent candidate
// Roster (-4/P2) is the persistent, device-independent candidate
// source; seenRequesters stays as the fallback for channels
// created before the roster partition existed.
const roster = await epochKeyManager.getRosterMembers(channel)
Expand Down
6 changes: 3 additions & 3 deletions src/js/epochKeyManager.js
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ const SEEN_WRAPS_MAX = KEYS_HISTORY_COUNT;
const PENDING_REQUEST_MAX_AGE_MS = 30 * 24 * 60 * 60 * 1000;
const PENDING_REQUESTS_MAX = 8;

// Roster (-4/P1) read cache: the members panel refreshes freely; the resend
// Roster (-4/P2) read cache: the members panel refreshes freely; the resend
// behind it should not.
const ROSTER_CACHE_TTL_MS = 60 * 1000;
const ROSTER_HISTORY_COUNT = 500;
Expand Down Expand Up @@ -1886,7 +1886,7 @@ class EpochKeyManager {
Logger.debug('epochKeys: member hello failed:', e.message));
}

// ==================== ROSTER (-4/P1) ====================
// ==================== ROSTER (-4/P2) ====================

/**
* Does this channel's -4 carry the roster partition? Resolved once per
Expand Down Expand Up @@ -2000,7 +2000,7 @@ class EpochKeyManager {
}

/**
* The channel roster: MEMBER_HELLO authors from -4/P1, deduped by account,
* The channel roster: MEMBER_HELLO authors from -4/P2, deduped by account,
* newest hello wins. Persistent and device-independent, unlike
* seenRequesters — the candidate source the members panel unions in.
* Every entry is authenticated: the hello opens with an epoch key valid at
Expand Down
27 changes: 14 additions & 13 deletions src/js/streamConstants.js
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,8 @@
* -5 → Interactions (WITH storage) — reactions, where members participate
*
* MESSAGE STREAM (-1):
* Regular channels use 11 partitions:
* P0 content, P1 control overrides, P2-P10 storage-file chunks.
* Regular channels use 12 partitions:
* P0 content, P1 control overrides, P2 moderator deltas (gated), P3-P11 storage-file chunks.
* DM inboxes use 13 partitions:
* P0 messages, P1 sync, P2 sync_blobs, P3 notifications, P4-P12 storage-file chunks.
*
Expand All @@ -30,21 +30,22 @@
* 3 partitions: control (presence/typing), media signals, media data.
*
* ADMIN STREAM (-3):
* 3 partitions reserved by protocol; only P0 used in initial scope.
* P0: ADMIN_STATE (moderation: bannedMembers, hiddenMessageIds, pins) — IMPLEMENTED
* P1: CHANNEL_IMAGE — RESERVED
* P2: PASSWORD_CHALLENGE — RESERVED
* 3 partitions:
* P0: ADMIN_STATE (moderation: bannedMembers, hiddenMessageIds, pins)
* P1: CHANNEL_IMAGE
* P2: PASSWORD_CHALLENGE
* Permissions: only owner publishes; readers vary by channel type
* (public/password: public subscribe; gated: clone subscribe).
*
* KEYS STREAM (-4) — gated channels only:
* P0 carries the epoch-key protocol (KEY_ANNOUNCE / KEY_REQUEST / KEY_WRAP).
* P0 carries the announces (KEY_ANNOUNCE), P1 the requests and their wraps
* (KEY_REQUEST / KEY_WRAP).
* Content on -1 is encrypted with a channel-wide epoch key versioned by
* `kid`; this stream is how members obtain those keys.
* P1 carries the member roster (MEMBER_HELLO): one hello per member per
* P2 carries the member roster (MEMBER_HELLO): one hello per member per
* epoch, ALWAYS sealed with that epoch's key — the -4 resend is publicly
* readable over HTTP, so a cleartext roster would be the worst membership
* leak in the system. Channels created before P1 existed have a
* leak in the system. Channels created before the roster existed have a
* single-partition -4 (capability = on-chain partition count).
* Permissions: members publish AND subscribe (any member may answer a request
* with a KEY_WRAP — k-of-n distribution). KEY_ANNOUNCE authority is app-layer:
Expand Down Expand Up @@ -79,7 +80,7 @@ export const MESSAGE_STREAM = Object.freeze({
* Persistent File Sharing over storage nodes (message stream -1, chunk partitions).
*
* Chunks are round-robined over 9 partitions starting right after the last
* "classic" partition of the stream flavor: P2 on regular channels, P4 on DM
* "classic" partition of the stream flavor: P3 on regular channels, P4 on DM
* inboxes. The announcement is a normal signed chat message on P0
* (type 'storage_file_announce') and carries firstChunkPartition/chunkPartitions,
* so readers follow the announce, not these local constants.
Expand All @@ -93,7 +94,7 @@ export const STORAGE_FILE = Object.freeze({
/**
* Partition for storage-file chunk i.
* @param {number} i - Chunk index
* @param {number} firstPartition - First chunk partition (2 regular / 4 DM, or from announce)
* @param {number} firstPartition - First chunk partition (3 regular / 4 DM, or from announce)
* @param {number} [count] - Number of chunk partitions (default 9, or from announce)
* @returns {number}
*/
Expand All @@ -113,7 +114,7 @@ export const EPHEMERAL_STREAM = Object.freeze({

export const ADMIN_STREAM = Object.freeze({
SUFFIX: STREAM_SUFFIX.ADMIN,
PARTITIONS: 3, // P0 implemented; P1 (channel image) and P2 (password challenge) reserved
PARTITIONS: 3, // P0 moderation, P1 channel image, P2 password challenge

// Partition indexes
MODERATION: 0, // ADMIN_STATE: bannedMembers, hiddenMessageIds, pins
Expand Down Expand Up @@ -178,7 +179,7 @@ export const INTERACTIONS_STREAM = Object.freeze({
* account key, so retained requests answer asynchronously.
* Receivers verify sha256(unwrapped) === announced keyHash
* before adopting, both formats.
* MEMBER_HELLO roster entry on P1 (never P0), sealed with the epoch key:
* MEMBER_HELLO roster entry on P2 (never P0), sealed with the epoch key:
* { account, spk, ts } — published on first adoption of each
* CURRENT epoch's key; readers require the envelope signer to
* equal `account` (no planting hellos for someone else)
Expand Down
4 changes: 2 additions & 2 deletions src/js/streamr.js
Original file line number Diff line number Diff line change
Expand Up @@ -616,7 +616,7 @@ class StreamrController {
Logger.info('Creating message stream...');
const startTime = Date.now();

const messageStream = await createStreamWithRetry(messageStreamId, metadata, 'message', STREAM_CONFIG.MESSAGE_STREAM.PARTITIONS); // 11 partitions for channels (content + control + 9 storage-file chunks)
const messageStream = await createStreamWithRetry(messageStreamId, metadata, 'message', STREAM_CONFIG.MESSAGE_STREAM.PARTITIONS); // 12 partitions for channels (content + control + moderation + 9 storage-file chunks)
try { onProgress(); } catch (_) { /* progress callback errors must not break creation */ }

// Step 2: Create EPHEMERAL STREAM (3 partitions: control + media signals + media data)
Expand Down Expand Up @@ -1106,7 +1106,7 @@ class StreamrController {
/**
* Create the DM inbox for the current user (dual-stream: message + ephemeral)
* Idempotent — if streams already exist, returns their IDs without recreating.
* Permissions: public SUBSCRIBE + PUBLISH (Streamr is a blind pipe; E2E encryption at app layer)
* Permissions: many-to-one, public PUBLISH and owner-only SUBSCRIBE (Streamr is a blind pipe; E2E encryption at app layer)
* @param {string} publicKey - Owner's compressed public key (hex, for ECDH)
* @param {Object} options - Storage options
* @param {string} options.storageProvider - 'streamr' or 'custom' (default: 'streamr')
Expand Down
95 changes: 0 additions & 95 deletions tests/unit/channels.extended.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -1023,101 +1023,6 @@ describe('ChannelManager Extended', () => {
});
});

// ==================== syncChannelFromGraph ====================
describe('syncChannelFromGraph', () => {
const streamId = 'stream-sync-1';

beforeEach(() => {
channelManager.channels.set(streamId, {
type: 'public',
members: [],
createdBy: null
});
});

it('returns null if channel not in map', async () => {
const result = await channelManager.syncChannelFromGraph('nonexistent');
expect(result).toBeNull();
});

it('fetches stream data, type, and members in parallel', async () => {
graphAPI.getStream.mockResolvedValue({ id: streamId });
graphAPI.detectStreamType.mockResolvedValue('password');
graphAPI.getStreamMembers.mockResolvedValue({
ok: true,
data: [{ address: '0xowner', isOwner: true }]
});

await channelManager.syncChannelFromGraph(streamId);

expect(graphAPI.getStream).toHaveBeenCalledWith(streamId);
expect(graphAPI.detectStreamType).toHaveBeenCalledWith(streamId);
expect(graphAPI.getStreamMembers).toHaveBeenCalledWith(streamId);
});

it('updates channel type from Graph', async () => {
graphAPI.getStream.mockResolvedValue({ id: streamId });
graphAPI.detectStreamType.mockResolvedValue('password');
graphAPI.getStreamMembers.mockResolvedValue({ ok: true, data: [] });

await channelManager.syncChannelFromGraph(streamId);

const channel = channelManager.channels.get(streamId);
expect(channel.type).toBe('password');
});

it('updates members and owner from Graph', async () => {
graphAPI.getStream.mockResolvedValue({ id: streamId });
graphAPI.detectStreamType.mockResolvedValue('password');
graphAPI.getStreamMembers.mockResolvedValue({
ok: true,
data: [
{ address: '0xOwner', isOwner: true },
{ address: '0xMember1', isOwner: false }
]
});

await channelManager.syncChannelFromGraph(streamId);

const channel = channelManager.channels.get(streamId);
expect(channel.members).toContain('0xOwner');
expect(channel.members).toContain('0xMember1');
expect(channel.createdBy).toBe('0xOwner');
});

it('saves channels after sync', async () => {
const saveSpy = vi.spyOn(channelManager, 'saveChannels').mockResolvedValue(undefined);
graphAPI.getStream.mockResolvedValue({ id: streamId });
graphAPI.detectStreamType.mockResolvedValue('public');
graphAPI.getStreamMembers.mockResolvedValue({ ok: true, data: [] });

await channelManager.syncChannelFromGraph(streamId);

expect(saveSpy).toHaveBeenCalled();
});

it('returns channel unchanged on error', async () => {
graphAPI.getStream.mockRejectedValue(new Error('network'));

const result = await channelManager.syncChannelFromGraph(streamId);

expect(result.type).toBe('public'); // unchanged
});

it('does not update type if detectStreamType returns unknown', async () => {
graphAPI.getStream.mockResolvedValue({ id: streamId });
graphAPI.detectStreamType.mockResolvedValue('unknown');
graphAPI.getStreamMembers.mockResolvedValue({ ok: true, data: [] });

const channelBefore = channelManager.channels.get(streamId);
channelBefore.type = 'password';

await channelManager.syncChannelFromGraph(streamId);

expect(channelManager.channels.get(streamId).type).toBe('password');
});
});

// ==================== publishPresence ====================
describe('publishPresence', () => {
it('does nothing if channel not found', async () => {
Expand Down
Loading